本文目录导读:

在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库 - 需要精细控制:手动实现重试逻辑
- 批量任务:使用重试管理器
- 异步任务:使用异步版本的重试管理器
记得根据具体场景选择合适的重试策略,避免过度重试导致资源浪费。