Python脚本如何处理分片同步失败任务

wen python案例 29

本文目录导读:

Python脚本如何处理分片同步失败任务

  1. 基础的重试机制
  2. 分片任务管理器
  3. 使用消息队列处理失败任务
  4. 使用数据库持久化失败任务
  5. 实践建议
  6. 完整使用示例

我来详细说明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()

这些方法可以组合使用,根据您的具体需求选择合适的方案,关键是要实现:

  1. 自动重试机制
  2. 失败任务持久化
  3. 任务状态跟踪
  4. 异常处理和告警

抱歉,评论功能暂时关闭!