本文目录导读:

我来详细说明Python处理分片同步失败任务的几种常见方法和最佳实践:
基础的重试机制
import time
from functools import wraps
def retry_on_failure(max_retries=3, delay=2, backoff=2):
"""
重试装饰器
:param max_retries: 最大重试次数
:param delay: 初始延迟
:param backoff: 退避因子
"""
def decorator(func):
@wraps(func)
def wrapper(*args, **kwargs):
retries = 0
current_delay = delay
while retries < max_retries:
try:
return func(*args, **kwargs)
except Exception as e:
retries += 1
if retries == max_retries:
raise Exception(f"重试{max_retries}次后仍然失败: {str(e)}")
print(f"第{retries}次重试,等待{current_delay}秒...")
time.sleep(current_delay)
current_delay *= backoff # 指数退避
return None
return wrapper
return decorator
# 使用示例
@retry_on_failure(max_retries=3, delay=1)
def sync_chunk(chunk_data):
# 同步分片数据的逻辑
if not upload_chunk(chunk_data):
raise Exception("分片同步失败")
return True
分片任务管理器
import json
import threading
from datetime import datetime
from typing import List, Dict, Any
class ShardSyncManager:
def __init__(self, max_retries=3, retry_delay=5):
self.max_retries = max_retries
self.retry_delay = retry_delay
self.failed_tasks = [] # 失败任务队列
self.sync_status = {} # 同步状态跟踪
self.lock = threading.Lock()
def add_task(self, task_id: str, shard_data: Dict[str, Any]):
"""添加分片任务"""
task = {
'task_id': task_id,
'shard_data': shard_data,
'retry_count': 0,
'status': 'pending',
'created_at': datetime.now(),
'last_retry': None
}
with self.lock:
self.sync_status[task_id] = task
def sync_shard(self, task_id: str) -> bool:
"""执行单个分片同步"""
task = self.sync_status.get(task_id)
if not task:
return False
while task['retry_count'] < self.max_retries:
try:
# 执行分片同步逻辑
result = self._execute_sync(task['shard_data'])
if result:
task['status'] = 'success'
task['completed_at'] = datetime.now()
return True
else:
raise Exception("同步返回失败")
except Exception as e:
task['retry_count'] += 1
task['last_retry'] = datetime.now()
if task['retry_count'] >= self.max_retries:
task['status'] = 'failed'
task['error'] = str(e)
self._handle_failure(task)
return False
print(f"任务{task_id}第{task['retry_count']}次重试")
time.sleep(self.retry_delay * task['retry_count'])
return False
def _execute_sync(self, shard_data: Dict) -> bool:
"""执行实际的分片同步(需要实现具体逻辑)"""
# 这里应该实现您的分片同步逻辑
pass
def _handle_failure(self, task: Dict):
"""处理失败任务"""
with self.lock:
self.failed_tasks.append(task)
# 可以在这里添加告警通知
print(f"任务{task['task_id']}同步失败,已加入失败队列")
def retry_failed_tasks(self):
"""重试所有失败任务"""
failed_tasks_copy = self.failed_tasks.copy()
self.failed_tasks.clear()
for task in failed_tasks_copy:
task['retry_count'] = 0
task['status'] = 'pending'
self.sync_shard(task['task_id'])
def get_failed_tasks(self) -> List[Dict]:
"""获取失败任务列表"""
return self.failed_tasks
def get_sync_stats(self) -> Dict:
"""获取同步统计信息"""
total = len(self.sync_status)
success = sum(1 for t in self.sync_status.values() if t['status'] == 'success')
failed = len(self.failed_tasks)
pending = total - success - failed
return {
'total': total,
'success': success,
'failed': failed,
'pending': pending
}
使用消息队列处理失败任务
import queue
import threading
class ShardSyncQueue:
def __init__(self, num_workers=3):
self.task_queue = queue.Queue()
self.retry_queue = queue.Queue()
self.failed_queue = queue.Queue()
self.workers = []
self.num_workers = num_workers
self.running = False
def start(self):
"""启动工作线程"""
self.running = True
for i in range(self.num_workers):
worker = threading.Thread(target=self._worker, args=(i,))
worker.daemon = True
worker.start()
self.workers.append(worker)
# 启动重试线程
retry_thread = threading.Thread(target=self._retry_worker)
retry_thread.daemon = True
retry_thread.start()
def add_task(self, shard_task):
"""添加任务到队列"""
self.task_queue.put(shard_task)
def _worker(self, worker_id):
"""工作线程函数"""
while self.running:
try:
task = self.task_queue.get(timeout=1)
print(f"Worker {worker_id} 处理任务: {task['id']}")
if self._process_task(task):
print(f"任务 {task['id']} 处理成功")
else:
# 处理失败,加入重试队列
task['retry_count'] = task.get('retry_count', 0) + 1
self.retry_queue.put(task)
except queue.Empty:
continue
except Exception as e:
print(f"Worker异常: {str(e)}")
def _retry_worker(self):
"""重试线程函数"""
while self.running:
try:
task = self.retry_queue.get(timeout=5)
if task['retry_count'] > 3:
# 重试次数过多,加入失败队列
self.failed_queue.put(task)
print(f"任务 {task['id']} 超过最大重试次数")
else:
print(f"重试任务 {task['id']}, 第{task['retry_count']}次重试")
time.sleep(5 * task['retry_count'])
self.task_queue.put(task)
except queue.Empty:
continue
def _process_task(self, task):
"""处理单个任务(需要实现具体逻辑)"""
# 这里实现您的分片同步逻辑
pass
def stop(self):
"""停止工作"""
self.running = False
使用数据库持久化失败任务
import sqlite3
from datetime import datetime
class PersistentShardSync:
def __init__(self, db_path='sync_tasks.db'):
self.db_path = db_path
self._init_database()
def _init_database(self):
"""初始化数据库表"""
with sqlite3.connect(self.db_path) as conn:
conn.execute('''
CREATE TABLE IF NOT EXISTS sync_tasks (
id INTEGER PRIMARY KEY AUTOINCREMENT,
shard_id TEXT NOT NULL,
shard_data TEXT,
status TEXT DEFAULT 'pending',
retry_count INTEGER DEFAULT 0,
max_retries INTEGER DEFAULT 3,
error_message TEXT,
created_at TIMESTAMP,
last_retry_at TIMESTAMP,
completed_at TIMESTAMP
)
''')
conn.execute('''
CREATE TABLE IF NOT EXISTS sync_logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
task_id INTEGER,
action TEXT,
message TEXT,
timestamp TIMESTAMP,
FOREIGN KEY (task_id) REFERENCES sync_tasks(id)
)
''')
conn.commit()
def add_task(self, shard_id: str, shard_data: dict):
"""添加同步任务"""
with sqlite3.connect(self.db_path) as conn:
conn.execute(
'INSERT INTO sync_tasks (shard_id, shard_data, status, created_at) VALUES (?, ?, ?, ?)',
(shard_id, json.dumps(shard_data), 'pending', datetime.now())
)
conn.commit()
def get_pending_tasks(self, limit=10):
"""获取待处理任务"""
with sqlite3.connect(self.db_path) as conn:
cursor = conn.execute(
'SELECT * FROM sync_tasks WHERE status = "pending" AND retry_count < max_retries LIMIT ?',
(limit,)
)
return cursor.fetchall()
def update_task_status(self, task_id: int, status: str, error_message=None):
"""更新任务状态"""
with sqlite3.connect(self.db_path) as conn:
if status == 'success':
conn.execute(
'UPDATE sync_tasks SET status = ?, completed_at = ? WHERE id = ?',
(status, datetime.now(), task_id)
)
elif status == 'failed':
conn.execute(
'UPDATE sync_tasks SET status = ?, retry_count = retry_count + 1, '
'error_message = ?, last_retry_at = ? WHERE id = ?',
(status, error_message, datetime.now(), task_id)
)
conn.commit()
def retry_failed_tasks(self):
"""重试失败任务"""
with sqlite3.connect(self.db_path) as conn:
# 将失败任务重置为待处理
conn.execute('''
UPDATE sync_tables
SET status = 'pending'
WHERE status = 'failed' AND retry_count < max_retries
''')
conn.commit()
实践建议
错误处理最佳实践
class ShardSyncError(Exception):
"""分片同步自定义异常"""
pass
class NetworkError(ShardSyncError):
"""网络错误"""
pass
class DataCorruptionError(ShardSyncError):
"""数据损坏错误"""
pass
def handle_sync_error(task_id, shard_data, error):
"""智能错误处理"""
if isinstance(error, NetworkError):
# 网络错误可以重试
return 'retry'
elif isinstance(error, DataCorruptionError):
# 数据损坏需要重新获取数据
return 'reacquire'
elif isinstance(error, PermissionError):
# 权限错误需要检查配置
return 'alert'
else:
# 未知错误需要记录并人工处理
return 'manual_review'
监控和告警
class SyncMonitor:
def __init__(self):
self.metrics = {
'total_tasks': 0,
'success_tasks': 0,
'failed_tasks': 0,
'avg_duration': 0,
'error_types': {}
}
def record_task(self, task_id, duration, success, error_type=None):
"""记录任务统计"""
self.metrics['total_tasks'] += 1
if success:
self.metrics['success_tasks'] += 1
else:
self.metrics['failed_tasks'] += 1
if error_type:
self.metrics['error_types'][error_type] = \
self.metrics['error_types'].get(error_type, 0) + 1
# 更新平均耗时
self.metrics['avg_duration'] = (
self.metrics['avg_duration'] * (self.metrics['total_tasks'] - 1) + duration
) / self.metrics['total_tasks']
def check_health(self):
"""健康检查"""
if self.metrics['failed_tasks'] > 10:
self.send_alert("失败任务过多,需要人工介入")
if self.metrics.get('error_types', {}).get('NetworkError', 0) > 5:
self.send_alert("网络错误频繁,请检查网络状况")
def send_alert(self, message):
"""发送告警"""
print(f"[ALERT] {message}")
# 实际可以集成邮件、短信、钉钉等通知方式
完整使用示例
def main():
# 初始化同步管理器
sync_manager = ShardSyncManager(max_retries=3)
# 模拟分片数据
shards = [
{'id': 'shard_1', 'data': '...'},
{'id': 'shard_2', 'data': '...'},
# ...
]
# 处理所有分片
for shard in shards:
sync_manager.add_task(shard['id'], shard)
# 同步所有分片
for task_id in sync_manager.sync_status:
if not sync_manager.sync_shard(task_id):
print(f"分片 {task_id} 同步失败")
# 检查失败任务并重试
failed_tasks = sync_manager.get_failed_tasks()
if failed_tasks:
print(f"有 {len(failed_tasks)} 个失败任务,准备重试...")
sync_manager.retry_failed_tasks()
# 输出统计信息
stats = sync_manager.get_sync_stats()
print(f"同步完成: {stats['success']}/{stats['total']}")
if __name__ == "__main__":
main()
这些方法可以组合使用,根据您的具体需求选择合适的方案,关键是要实现:
- 自动重试机制
- 失败任务持久化
- 任务状态跟踪
- 异常处理和告警