"""呼叫任务API函数""" from celery_app import app from datetime import datetime import logging import requests from api_config import API_CONFIG, RETRY_CONFIG logger = logging.getLogger(__name__) @app.task(bind=True, name='call_api', max_retries=None) def execute_call_api_task(self): """遍历API配置文件中的所有API信息并调用接口""" current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S') logger.info(f"开始遍历API配置文件执行调用任务 - 当前时间: {current_time}") # 获取重试配置 max_retries = RETRY_CONFIG['max_retries'] retry_delay = RETRY_CONFIG['retry_delay'] retry_on_status = RETRY_CONFIG['retry_on_status'] # 遍历所有API信息 results = [] success_count = 0 for api_key, call_config in API_CONFIG.items(): logger.info(f"开始处理API配置: {api_key}") call_url = call_config['url'] method = call_config['method'].upper() headers = call_config['headers'] timeout = call_config['timeout'] payload = call_config['body'] try: logger.info(f"正在调用任务API: {method} {call_url}") # 根据method选择请求方式 if method == 'POST': response = requests.post(call_url, headers=headers, json=payload, timeout=timeout) elif method == 'GET': response = requests.get(call_url, headers=headers, params=payload, timeout=timeout) elif method == 'PUT': response = requests.put(call_url, headers=headers, json=payload, timeout=timeout) elif method == 'DELETE': response = requests.delete(call_url, headers=headers, timeout=timeout) else: raise ValueError(f"不支持的HTTP方法: {method}") # 检查响应状态 if response.status_code in [200, 201, 204]: # 成功状态码 try: result_data = response.json() except: result_data = {'response': response.text} logger.info(f"任务API调用成功,响应数据: {result_data}") print(f"[{current_time}] ✅ 任务API成功 ({api_key}) - {method} {call_url}") success_count += 1 results.append({ 'call_api_key': api_key, 'status': 'success', 'message': f'任务API调用成功 ({method})', 'response': result_data, 'timestamp': current_time, 'retry_count': getattr(self.request, 'retries', 0) }) elif response.status_code in retry_on_status: # 需要重试的状态码 current_retry = getattr(self.request, 'retries', 0) if current_retry < max_retries: logger.warning(f"任务API需要重试,当前重试次数: {current_retry + 1}/{max_retries}") print(f"[{current_time}] 🔄 任务API失败,正在重试 ({current_retry + 1}/{max_retries}) - {api_key} - 状态码: {response.status_code}") raise self.retry(countdown=retry_delay, exc=Exception(f"任务API失败,状态码: {response.status_code}")) else: logger.error(f"任务API失败,已达到最大重试次数: {max_retries}") print(f"[{current_time}] ❌ 任务API失败,已达到最大重试次数 - {api_key} - 状态码: {response.status_code}") results.append({ 'call_api_key': api_key, 'status': 'failed', 'message': f'任务API失败,已达到最大重试次数: {max_retries}', 'response': response.text, 'timestamp': current_time, 'retry_count': current_retry + 1 }) else: # 其他失败状态码,不重试 logger.error(f"任务API失败,状态码: {response.status_code}, 响应: {response.text}") print(f"[{current_time}] ❌ 任务API失败 - {api_key} - 状态码: {response.status_code}") results.append({ 'call_api_key': api_key, 'status': 'failed', 'message': f'任务API失败,状态码: {response.status_code}', 'response': response.text, 'timestamp': current_time, 'retry_count': getattr(self.request, 'retries', 0) }) except requests.exceptions.RequestException as e: # 网络异常,需要重试 current_retry = getattr(self.request, 'retries', 0) if current_retry < max_retries: logger.warning(f"任务API异常,需要重试,当前重试次数: {current_retry + 1}/{max_retries}") print(f"[{current_time}] 🔄 任务API异常,正在重试 ({current_retry + 1}/{max_retries}) - {api_key} - 异常: {str(e)}") raise self.retry(countdown=retry_delay, exc=e) else: logger.error(f"任务API异常,已达到最大重试次数: {max_retries}") print(f"[{current_time}] ❌ 任务API异常,已达到最大重试次数 - {api_key} - 异常: {str(e)}") results.append({ 'call_api_key': api_key, 'status': 'error', 'message': f'任务API异常,已达到最大重试次数: {max_retries}', 'error': str(e), 'timestamp': current_time, 'retry_count': current_retry + 1 }) # 汇总结果 total_count = len(results) logger.info(f"遍历API配置文件完成 - 成功: {success_count}/{total_count}") print(f"[{current_time}] 📊 任务汇总 - 成功: {success_count}/{total_count}") return { 'status': 'completed', 'total_apis': total_count, 'success_count': success_count, 'failed_count': total_count - success_count, 'results': results, 'timestamp': current_time }