Python脚本如何重试失败数据同步任务

wen python案例 29

本文目录导读:

Python脚本如何重试失败数据同步任务

  1. 使用 tenacity 库(推荐)
  2. 手动实现重试逻辑
  3. 使用 backoff
  4. 完整的重试管理器
  5. 实践建议
  6. 选择建议

在Python中重试失败的数据同步任务,有几种常见的方法:

使用 tenacity 库(推荐)

最简单且功能强大的重试库:

from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type
import requests
@retry(
    stop=stop_after_attempt(3),  # 最多重试3次
    wait=wait_exponential(multiplier=1, min=4, max=10),  # 指数退避
    retry=retry_if_exception_type((requests.RequestException, ConnectionError))
)
def sync_data():
    # 数据同步逻辑
    response = requests.get('https://api.example.com/sync')
    response.raise_for_status()
    return response.json()
# 使用
try:
    result = sync_data()
except Exception as e:
    print(f"同步失败: {e}")

手动实现重试逻辑

更灵活的自定义实现:

import time
import logging
from functools import wraps
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
def retry_sync(max_retries=3, delay=1, backoff=2, exceptions=(Exception,)):
    def decorator(func):
        @wraps(func)
        def wrapper(*args, **kwargs):
            retries = 0
            current_delay = delay
            while retries < max_retries:
                try:
                    return func(*args, **kwargs)
                except exceptions as e:
                    retries += 1
                    if retries == max_retries:
                        logger.error(f"重试{max_retries}次后仍然失败: {e}")
                        raise
                    logger.warning(f"同步失败 (第{retries}次重试): {e}")
                    time.sleep(current_delay)
                    current_delay *= backoff  # 指数退避
            return None
        return wrapper
    return decorator
@retry_sync(max_retries=5, delay=2, backoff=2)
def sync_data_to_db():
    # 数据库同步逻辑
    pass
@retry_sync(max_retries=3, exceptions=(ConnectionError, TimeoutError))
def sync_data_to_api():
    # API同步逻辑
    pass

使用 backoff

import backoff
import requests
@backoff.on_exception(
    backoff.expo,  # 指数退避
    requests.RequestException,
    max_tries=5,
    max_time=60  # 最大重试时间
)
def sync_data_with_backoff():
    # 数据同步
    response = requests.post('https://api.example.com/sync', json={'data': 'value'})
    response.raise_for_status()
    return response.json()

完整的重试管理器

import time
import logging
from typing import Callable, List, Dict, Any
from dataclasses import dataclass
from enum import Enum
logger = logging.getLogger(__name__)
class SyncStatus(Enum):
    SUCCESS = "success"
    FAILED = "failed"
    RETRYING = "retrying"
@dataclass
class SyncResult:
    status: SyncStatus
    retry_count: int
    error: str = ""
    data: Any = None
class RetryManager:
    def __init__(
        self,
        max_retries: int = 3,
        base_delay: float = 1.0,
        max_delay: float = 60.0,
        exponential_backoff: bool = True
    ):
        self.max_retries = max_retries
        self.base_delay = base_delay
        self.max_delay = max_delay
        self.exponential_backoff = exponential_backoff
    def get_delay(self, retry_count: int) -> float:
        """计算重试延迟时间"""
        if self.exponential_backoff:
            delay = min(
                self.base_delay * (2 ** retry_count),
                self.max_delay
            )
        else:
            delay = self.base_delay
        # 添加随机抖动
        import random
        jitter = random.uniform(0, 0.1 * delay)
        return delay + jitter
    def execute(self, sync_func: Callable, *args, **kwargs) -> SyncResult:
        """执行带重试的同步任务"""
        last_error = ""
        for attempt in range(self.max_retries + 1):
            try:
                logger.info(f"开始同步 (尝试 {attempt + 1}/{self.max_retries + 1})")
                result = sync_func(*args, **kwargs)
                return SyncResult(
                    status=SyncStatus.SUCCESS,
                    retry_count=attempt,
                    data=result
                )
            except Exception as e:
                last_error = str(e)
                if attempt < self.max_retries:
                    delay = self.get_delay(attempt)
                    logger.warning(
                        f"同步失败 (尝试 {attempt + 1}): {e}. "
                        f"将在 {delay:.2f} 秒后重试..."
                    )
                    time.sleep(delay)
                else:
                    logger.error(f"重试 {self.max_retries} 次后仍然失败: {e}")
        return SyncResult(
            status=SyncStatus.FAILED,
            retry_count=self.max_retries,
            error=last_error
        )
# 使用示例
def sync_database():
    # 模拟数据库同步
    import random
    if random.random() < 0.7:  # 70%概率失败
        raise ConnectionError("数据库连接失败")
    return {"synced": True}
def batch_sync(tasks: List[Dict]) -> List[SyncResult]:
    """批量同步多个任务"""
    manager = RetryManager(max_retries=3, exponential_backoff=True)
    results = []
    for task in tasks:
        result = manager.execute(sync_database, **task)
        results.append(result)
        if result.status == SyncStatus.SUCCESS:
            logger.info(f"任务 {task} 同步成功")
        else:
            logger.error(f"任务 {task} 同步失败: {result.error}")
    return results
# 异步版本的重试管理器
import asyncio
class AsyncRetryManager:
    def __init__(self, max_retries=3, base_delay=1.0, max_delay=60.0):
        self.max_retries = max_retries
        self.base_delay = base_delay
        self.max_delay = max_delay
    async def execute(self, sync_func, *args, **kwargs):
        last_error = ""
        for attempt in range(self.max_retries + 1):
            try:
                result = await sync_func(*args, **kwargs)
                return SyncResult(
                    status=SyncStatus.SUCCESS,
                    retry_count=attempt,
                    data=result
                )
            except Exception as e:
                last_error = str(e)
                if attempt < self.max_retries:
                    delay = min(self.base_delay * (2 ** attempt), self.max_delay)
                    await asyncio.sleep(delay)
        return SyncResult(
            status=SyncStatus.FAILED,
            retry_count=self.max_retries,
            error=last_error
        )
# 异步使用示例
async def async_sync_example():
    manager = AsyncRetryManager()
    async def sync_task():
        # 异步同步逻辑
        await asyncio.sleep(1)
        return {"status": "success"}
    result = await manager.execute(sync_task)
    return result

实践建议

# 1. 区分可重试和不可重试的错误
class RetryableError(Exception):
    """可重试的错误"""
    pass
class NonRetryableError(Exception):
    """不可重试的错误"""
    pass
def sync_with_retry_policy():
    manager = RetryManager()
    def sync_logic():
        try:
            # 同步逻辑
            pass
        except NonRetryableError:
            raise  # 不重试
        except RetryableError:
            raise  # 会重试
        except Exception as e:
            # 其他异常根据情况决定
            if "timeout" in str(e).lower():
                raise RetryableError(e)
            raise NonRetryableError(e)
    result = manager.execute(sync_logic)
    return result
# 2. 记录重试日志
import logging
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
# 3. 设置监控和告警
def monitor_sync_result(result: SyncResult):
    if result.status == SyncStatus.FAILED:
        # 发送告警
        send_alert(f"数据同步失败: {result.error}")
    # 记录指标
    record_metric(
        name="sync_retry_count",
        value=result.retry_count,
        tags={"status": result.status.value}
    )

选择建议

  • 简单场景:使用 tenacity
  • 需要精细控制:手动实现重试逻辑
  • 批量任务:使用重试管理器
  • 异步任务:使用异步版本的重试管理器

记得根据具体场景选择合适的重试策略,避免过度重试导致资源浪费。

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