PCNUSMSE/transcript_service
0
1"""错误处理和容错机制模块2 3提供统一的错误处理、重试逻辑和异常恢复功能。4"""5 6import asyncio7import functools8import time9from typing import Any, Callable, Dict, Optional, Type, Union10from enum import Enum11 12from ..core.config import get_config13from ..utils.logger import get_task_logger14 15 16class ErrorCode(Enum):17 """错误代码"""18 # 文件相关错误19 FILE_NOT_FOUND = "FILE_001"20 FILE_TOO_LARGE = "FILE_002"21 FILE_FORMAT_UNSUPPORTED = "FILE_003"22 FILE_CORRUPTED = "FILE_004"23 24 # 网络相关错误25 NETWORK_TIMEOUT = "NET_001"26 NETWORK_CONNECTION_ERROR = "NET_002"27 NETWORK_DNS_ERROR = "NET_003"28 29 # API相关错误30 API_KEY_INVALID = "API_001"31 API_QUOTA_EXCEEDED = "API_002"32 API_SERVICE_UNAVAILABLE = "API_003"33 API_RATE_LIMITED = "API_004"34 35 # OSS相关错误36 OSS_ACCESS_DENIED = "OSS_001"37 OSS_BUCKET_NOT_FOUND = "OSS_002"38 OSS_UPLOAD_FAILED = "OSS_003"39 40 # 系统相关错误41 SYSTEM_OUT_OF_MEMORY = "SYS_001"42 SYSTEM_DISK_FULL = "SYS_002"43 SYSTEM_PERMISSION_DENIED = "SYS_003"44 45 # 通用错误46 UNKNOWN_ERROR = "GEN_001"47 TIMEOUT_ERROR = "GEN_002"48 VALIDATION_ERROR = "GEN_003"49 50 51class TranscriptServiceError(Exception):52 """服务自定义异常基类"""53 54 def __init__(self, message: str, error_code: ErrorCode = ErrorCode.UNKNOWN_ERROR, details: Dict = None):55 """初始化异常56 57 Args:58 message: 错误消息59 error_code: 错误代码60 details: 额外详情61 """62 super().__init__(message)63 self.message = message64 self.error_code = error_code65 self.details = details or {}66 self.timestamp = time.time()67 68 def to_dict(self) -> Dict[str, Any]:69 """转换为字典格式"""70 return {71 'error_code': self.error_code.value,72 'message': self.message,73 'details': self.details,74 'timestamp': self.timestamp75 }76 77 78class FileValidationError(TranscriptServiceError):79 """文件验证错误"""80 pass81 82 83class NetworkError(TranscriptServiceError):84 """网络相关错误"""85 pass86 87 88class APIError(TranscriptServiceError):89 """API调用错误"""90 pass91 92 93class OSSError(TranscriptServiceError):94 """OSS操作错误"""95 pass96 97 98class SystemError(TranscriptServiceError):99 """系统错误"""100 pass101 102 103class RetryStrategy:104 """重试策略"""105 106 def __init__(107 self,108 max_attempts: int = 3,109 base_delay: float = 1.0,110 max_delay: float = 60.0,111 exponential_base: float = 2.0,112 jitter: bool = True113 ):114 """初始化重试策略115 116 Args:117 max_attempts: 最大重试次数118 base_delay: 基础延迟时间(秒)119 max_delay: 最大延迟时间(秒)120 exponential_base: 指数退避基数121 jitter: 是否添加随机抖动122 """123 self.max_attempts = max_attempts124 self.base_delay = base_delay125 self.max_delay = max_delay126 self.exponential_base = exponential_base127 self.jitter = jitter128 129 def calculate_delay(self, attempt: int) -> float:130 """计算延迟时间131 132 Args:133 attempt: 当前尝试次数(从1开始)134 135 Returns:136 延迟时间(秒)137 """138 delay = self.base_delay * (self.exponential_base ** (attempt - 1))139 delay = min(delay, self.max_delay)140 141 if self.jitter:142 import random143 delay *= (0.5 + random.random() * 0.5) # 添加±50%的随机抖动144 145 return delay146 147 148class ErrorHandler:149 """错误处理器"""150 151 def __init__(self):152 """初始化错误处理器"""153 self.config = get_config()154 self.logger = get_task_logger(logger_name="transcript_service.error")155 156 # 错误分类映射157 self.error_mapping = {158 # 文件错误159 FileNotFoundError: (FileValidationError, ErrorCode.FILE_NOT_FOUND),160 PermissionError: (SystemError, ErrorCode.SYSTEM_PERMISSION_DENIED),161 162 # 网络错误163 asyncio.TimeoutError: (NetworkError, ErrorCode.NETWORK_TIMEOUT),164 ConnectionError: (NetworkError, ErrorCode.NETWORK_CONNECTION_ERROR),165 166 # 通用错误167 ValueError: (TranscriptServiceError, ErrorCode.VALIDATION_ERROR),168 RuntimeError: (TranscriptServiceError, ErrorCode.UNKNOWN_ERROR),169 }170 171 # 可重试的错误类型172 self.retryable_errors = {173 ErrorCode.NETWORK_TIMEOUT,174 ErrorCode.NETWORK_CONNECTION_ERROR,175 ErrorCode.API_RATE_LIMITED,176 ErrorCode.OSS_UPLOAD_FAILED,177 ErrorCode.API_SERVICE_UNAVAILABLE178 }179 180 def classify_error(self, error: Exception) -> TranscriptServiceError:181 """分类和包装错误182 183 Args:184 error: 原始异常185 186 Returns:187 分类后的服务异常188 """189 if isinstance(error, TranscriptServiceError):190 return error191 192 error_type = type(error)193 if error_type in self.error_mapping:194 exception_class, error_code = self.error_mapping[error_type]195 return exception_class(str(error), error_code)196 197 # 根据错误消息内容进行分类198 error_msg = str(error).lower()199 200 if "timeout" in error_msg:201 return NetworkError(str(error), ErrorCode.NETWORK_TIMEOUT)202 elif "permission denied" in error_msg:203 return SystemError(str(error), ErrorCode.SYSTEM_PERMISSION_DENIED)204 elif "api key" in error_msg:205 return APIError(str(error), ErrorCode.API_KEY_INVALID)206 elif "quota" in error_msg or "limit" in error_msg:207 return APIError(str(error), ErrorCode.API_QUOTA_EXCEEDED)208 else:209 return TranscriptServiceError(str(error), ErrorCode.UNKNOWN_ERROR)210 211 def is_retryable(self, error: TranscriptServiceError) -> bool:212 """判断错误是否可重试213 214 Args:215 error: 服务异常216 217 Returns:218 是否可重试219 """220 return error.error_code in self.retryable_errors221 222 def handle_error(self, error: Exception, context: str = "") -> TranscriptServiceError:223 """处理错误224 225 Args:226 error: 原始异常227 context: 错误上下文228 229 Returns:230 处理后的服务异常231 """232 classified_error = self.classify_error(error)233 234 # 记录错误日志235 log_msg = f"错误处理 - {context}: {classified_error.message}"236 if classified_error.error_code in [ErrorCode.UNKNOWN_ERROR, ErrorCode.SYSTEM_OUT_OF_MEMORY]:237 self.logger.exception(log_msg)238 else:239 self.logger.error(log_msg)240 241 return classified_error242 243 244# 全局错误处理器实例245error_handler = ErrorHandler()246 247 248def retry_async(249 strategy: Optional[RetryStrategy] = None,250 exceptions: tuple = (Exception,),251 context: str = ""252):253 """异步函数重试装饰器254 255 Args:256 strategy: 重试策略257 exceptions: 需要重试的异常类型258 context: 上下文信息259 """260 if strategy is None:261 strategy = RetryStrategy()262 263 def decorator(func: Callable):264 @functools.wraps(func)265 async def wrapper(*args, **kwargs):266 logger = get_task_logger(logger_name="transcript_service.retry")267 268 for attempt in range(1, strategy.max_attempts + 1):269 try:270 return await func(*args, **kwargs)271 except exceptions as e:272 classified_error = error_handler.classify_error(e)273 274 # 检查是否可重试275 if attempt == strategy.max_attempts or not error_handler.is_retryable(classified_error):276 logger.error(f"{context} 最终失败 (尝试 {attempt}/{strategy.max_attempts}): {str(e)}")277 raise classified_error278 279 # 计算延迟时间280 delay = strategy.calculate_delay(attempt)281 logger.warning(f"{context} 第 {attempt} 次尝试失败,{delay:.1f}秒后重试: {str(e)}")282 283 await asyncio.sleep(delay)284 285 # 理论上不会执行到这里286 raise TranscriptServiceError("重试逻辑异常", ErrorCode.UNKNOWN_ERROR)287 288 return wrapper289 return decorator290 291 292def retry_sync(293 strategy: Optional[RetryStrategy] = None,294 exceptions: tuple = (Exception,),295 context: str = ""296):297 """同步函数重试装饰器298 299 Args:300 strategy: 重试策略301 exceptions: 需要重试的异常类型302 context: 上下文信息303 """304 if strategy is None:305 strategy = RetryStrategy()306 307 def decorator(func: Callable):308 @functools.wraps(func)309 def wrapper(*args, **kwargs):310 logger = get_task_logger(logger_name="transcript_service.retry")311 312 for attempt in range(1, strategy.max_attempts + 1):313 try:314 return func(*args, **kwargs)315 except exceptions as e:316 classified_error = error_handler.classify_error(e)317 318 # 检查是否可重试319 if attempt == strategy.max_attempts or not error_handler.is_retryable(classified_error):320 logger.error(f"{context} 最终失败 (尝试 {attempt}/{strategy.max_attempts}): {str(e)}")321 raise classified_error322 323 # 计算延迟时间324 delay = strategy.calculate_delay(attempt)325 logger.warning(f"{context} 第 {attempt} 次尝试失败,{delay:.1f}秒后重试: {str(e)}")326 327 time.sleep(delay)328 329 # 理论上不会执行到这里330 raise TranscriptServiceError("重试逻辑异常", ErrorCode.UNKNOWN_ERROR)331 332 return wrapper333 return decorator334 335 336def safe_execute(func: Callable, *args, **kwargs) -> tuple[bool, Any, Optional[TranscriptServiceError]]:337 """安全执行函数338 339 Args:340 func: 要执行的函数341 *args: 位置参数342 **kwargs: 关键字参数343 344 Returns:345 (是否成功, 结果或None, 错误或None)346 """347 try:348 result = func(*args, **kwargs)349 return True, result, None350 except Exception as e:351 error = error_handler.handle_error(e, f"执行 {func.__name__}")352 return False, None, error353 354 355async def safe_execute_async(func: Callable, *args, **kwargs) -> tuple[bool, Any, Optional[TranscriptServiceError]]:356 """安全执行异步函数357 358 Args:359 func: 要执行的异步函数360 *args: 位置参数361 **kwargs: 关键字参数362 363 Returns:364 (是否成功, 结果或None, 错误或None)365 """366 try:367 result = await func(*args, **kwargs)368 return True, result, None369 except Exception as e:370 error = error_handler.handle_error(e, f"执行 {func.__name__}")371 return False, None, error372 373 374def get_error_handler() -> ErrorHandler:375 """获取错误处理器实例376 377 Returns:378 错误处理器实例379 """380 return error_handler