From 0af7eebed8da31dd23d38c4755a4b1025c4c0b6b Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Tue, 9 Dec 2025 17:21:32 +0800 Subject: [PATCH] =?UTF-8?q?=E8=B0=83=E8=AF=95=E5=AE=8C=E6=88=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 3 + README.md | 14 +- invoke_call_api_task.py => call_api_task.py | 4 +- celery_app.py | 19 ++- tasks.py | 136 -------------------- test_api.py | 43 ------- test_api_list.py | 37 ------ test_api_list_execution.py | 67 ---------- test_beat_schedule.py | 53 -------- test_body_config.py | 84 ------------ test_call_api.py | 47 ------- test_celery_api_list.py | 73 ----------- test_direct_api.py | 129 ------------------- test_merged_status.py | 77 ----------- test_methods.py | 53 -------- test_quick_api.py | 45 ------- test_quick_updated.py | 32 ----- test_retry.py | 38 ------ test_schedule.py | 39 ------ test_simple_api_list.py | 67 ---------- test_simulate_failure.py | 58 --------- test_task.py | 10 -- test_updated_invoke_call.py | 39 ------ 23 files changed, 25 insertions(+), 1142 deletions(-) rename invoke_call_api_task.py => call_api_task.py (98%) delete mode 100644 tasks.py delete mode 100644 test_api.py delete mode 100644 test_api_list.py delete mode 100644 test_api_list_execution.py delete mode 100644 test_beat_schedule.py delete mode 100644 test_body_config.py delete mode 100644 test_call_api.py delete mode 100644 test_celery_api_list.py delete mode 100644 test_direct_api.py delete mode 100644 test_merged_status.py delete mode 100644 test_methods.py delete mode 100644 test_quick_api.py delete mode 100644 test_quick_updated.py delete mode 100644 test_retry.py delete mode 100644 test_schedule.py delete mode 100644 test_simple_api_list.py delete mode 100644 test_simulate_failure.py delete mode 100644 test_task.py delete mode 100644 test_updated_invoke_call.py diff --git a/.gitignore b/.gitignore index d62945a..d9f9482 100644 --- a/.gitignore +++ b/.gitignore @@ -176,3 +176,6 @@ cython_debug/ .vscode +/celerybeat-schedule.bak +/celerybeat-schedule.dat +/celerybeat-schedule.dir diff --git a/README.md b/README.md index bee9002..3724417 100644 --- a/README.md +++ b/README.md @@ -74,10 +74,9 @@ python simulate_failure.py # 测试不同HTTP方法的API调用 python test_methods.py -# 使用通用API任务 -from generic_api_task import generic_api_call -result = generic_api_call.delay('get_api') # 调用GET API -result = generic_api_call.delay('put_api', {'key': 'value'}) # 调用PUT API +# 使用呼叫API任务 +from call_api_task import execute_call_api_task +result = execute_call_api_task.delay() # 调用所有配置的API ``` ### 4. 测试Body配置 @@ -85,10 +84,9 @@ result = generic_api_call.delay('put_api', {'key': 'value'}) # 调用PUT API # 测试不同API配置的body参数 python test_body_config.py -# 使用自定义body参数 -from generic_api_task import generic_api_call -custom_body = {'title': '自定义标题', 'content': '自定义内容'} -result = generic_api_call.delay('main_api', custom_body) +# 使用呼叫API任务(自动遍历所有配置) +from call_api_task import execute_call_api_task +result = execute_call_api_task.delay() # 自动调用所有配置的API ``` ### 3. 手动触发定时任务 diff --git a/invoke_call_api_task.py b/call_api_task.py similarity index 98% rename from invoke_call_api_task.py rename to call_api_task.py index 422a878..e1d6d1b 100644 --- a/invoke_call_api_task.py +++ b/call_api_task.py @@ -7,8 +7,8 @@ from api_config import API_CONFIG, RETRY_CONFIG logger = logging.getLogger(__name__) -@app.task(bind=True, max_retries=None) -def invoke_call_api_task(self): +@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}") diff --git a/celery_app.py b/celery_app.py index 174e7a1..3eb8295 100644 --- a/celery_app.py +++ b/celery_app.py @@ -7,24 +7,33 @@ logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) # 创建Celery应用 -app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') +app = Celery( + 'tasks', + broker='redis://localhost:6379/0', + backend='redis://localhost:6379/0', + include=['call_api_task'] +) # 配置Celery app.conf.update( timezone='Asia/Shanghai', - enable_utc=True, + enable_utc=False, broker_url='redis://localhost:6379/0', result_backend='redis://localhost:6379/0', task_serializer='json', accept_content=['json'], result_serializer='json', + ask_track_started=True, + + # Beat 调度配置 + # crontab(hour=9, minute=0), # 每天上午9点执行 beat_schedule={ 'daily-morning-task': { - 'task': 'tasks.daily_morning_task', - 'schedule': 60.0 * 2.0, # 每2分钟执行一次 - 'options': {'queue': 'default'} + 'task': 'call_api', + 'schedule': 60.0 * 2.0 }, }, + # Flower监控配置 worker_send_task_events=True, task_send_sent_event=True, diff --git a/tasks.py b/tasks.py deleted file mode 100644 index 3ed9fb5..0000000 --- a/tasks.py +++ /dev/null @@ -1,136 +0,0 @@ -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, max_retries=None) -def daily_morning_task(self): - """每天定时运行配置中的API接口并支持重试机制""" - current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S') - logger.info(f"开始执行,当前时间: {current_time}") - - # 获取重试配置 - max_retries = RETRY_CONFIG['max_retries'] - retry_delay = RETRY_CONFIG['retry_delay'] - retry_on_status = RETRY_CONFIG['retry_on_status'] - - # 遍历配置文件中的所有API接口 - results = [] - for api_name, api_config in API_CONFIG.items(): - logger.info(f"开始处理API配置: {api_name}") - - api_url = api_config['url'] - headers = api_config['headers'] - timeout = api_config['timeout'] - payload = api_config['body'] - method = api_config['method'].upper() - - try: - logger.info(f"正在调用API接口: {method} {api_url}") - - # 根据method选择请求方式 - if method == 'POST': - response = requests.post(api_url, headers=headers, json=payload, timeout=timeout) - elif method == 'GET': - response = requests.get(api_url, headers=headers, params=payload, timeout=timeout) - elif method == 'PUT': - response = requests.put(api_url, headers=headers, json=payload, timeout=timeout) - elif method == 'DELETE': - response = requests.delete(api_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}") - - # 构建成功消息 - if response.status_code == 201: - print(f"[{current_time}] ✅ API调用成功 ({api_name}) - ID: {result_data.get('id')}") - message = 'API调用成功' - response_id = result_data.get('id') - else: - print(f"[{current_time}] ✅ API调用成功 ({api_name}) - {method}") - message = f'API调用成功 ({method})' - response_id = None - - results.append({ - 'api_name': api_name, - 'status': 'success', - 'message': message, - 'response': result_data, - 'response_id': response_id, - '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_name} - 状态码: {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_name} - 状态码: {response.status_code}") - results.append({ - 'api_name': api_name, - '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_name} - 状态码: {response.status_code}") - results.append({ - 'api_name': api_name, - '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_name} - 异常: {str(e)}") - raise self.retry(countdown=retry_delay, exc=e) - else: - logger.error(f"API调用异常,已达到最大重试次数: {max_retries}") - print(f"[{current_time}] ❌ API调用异常,已达到最大重试次数 - {api_name} - 异常: {str(e)}") - results.append({ - 'api_name': api_name, - 'status': 'error', - 'message': f'API调用异常,已达到最大重试次数: {max_retries}', - 'error': str(e), - 'timestamp': current_time, - 'retry_count': current_retry + 1 - }) - - # 汇总结果 - success_count = len([r for r in results if r['status'] == 'success']) - total_count = len(results) - logger.info(f"定时任务执行完成 - 成功: {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 - } \ No newline at end of file diff --git a/test_api.py b/test_api.py deleted file mode 100644 index 5c86815..0000000 --- a/test_api.py +++ /dev/null @@ -1,43 +0,0 @@ -#!/usr/bin/env python3 -"""测试API调用脚本""" -import requests -from datetime import datetime -import json - -def test_api_call(): - """测试API调用""" - current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S') - - api_url = "https://jsonplaceholder.typicode.com/posts" - headers = { - 'Content-Type': 'application/json', - 'User-Agent': 'Celery-Task/1.0' - } - payload = { - 'title': f'API测试 - {current_time}', - 'body': f'这是测试API调用的数据', - 'userId': 1 - } - - try: - print(f"正在测试API调用: {api_url}") - response = requests.post(api_url, headers=headers, json=payload, timeout=30) - - if response.status_code == 201: - result_data = response.json() - print(f"✅ API调用成功") - print(f"响应数据: {json.dumps(result_data, indent=2, ensure_ascii=False)}") - return True - else: - print(f"❌ API调用失败 - 状态码: {response.status_code}") - print(f"响应内容: {response.text}") - return False - - except requests.exceptions.RequestException as e: - print(f"❌ API调用异常: {str(e)}") - return False - -if __name__ == '__main__': - print("=== API调用测试 ===") - success = test_api_call() - print(f"\n测试结果: {'成功' if success else '失败'}") \ No newline at end of file diff --git a/test_api_list.py b/test_api_list.py deleted file mode 100644 index 5e4ebdd..0000000 --- a/test_api_list.py +++ /dev/null @@ -1,37 +0,0 @@ -#!/usr/bin/env python3 -"""测试遍历API配置列表功能""" - -import os -import sys -from datetime import datetime -from tasks import daily_morning_task - -def test_api_list_execution(): - """测试遍历API配置列表的执行""" - try: - print("开始测试遍历API配置列表功能...") - - # 直接调用任务函数进行测试 - result = daily_morning_task() - - print("\n" + "=" * 40) - print("测试结果汇总:") - print("=" * 40) - print(f"总API数量: {result['total_apis']}") - print(f"成功数量: {result['success_count']}") - print(f"失败数量: {result['failed_count']}") - print(f"总体状态: {result['status']}") - - print("\n各API执行状态:") - for api_result in result['results']: - status_icon = "✅" if api_result['status'] == 'success' else "❌" - print(f"{status_icon} {api_result['api_name']}") - - return result - - except Exception as e: - print(f"❌ 测试失败: {str(e)}") - return None - -if __name__ == "__main__": - test_api_list_execution() \ No newline at end of file diff --git a/test_api_list_execution.py b/test_api_list_execution.py deleted file mode 100644 index 36fdcc3..0000000 --- a/test_api_list_execution.py +++ /dev/null @@ -1,67 +0,0 @@ -#!/usr/bin/env python3 -"""测试遍历API配置列表执行调用任务""" - -from invoke_call_api_task import invoke_call_api_task -from api_config import API_CONFIG -import time - -def main(): - """遍历API配置文件中的所有API并执行调用""" - print("🚀 开始遍历API配置列表执行调用任务") - print("=" * 60) - - results = [] - - for api_key, config in API_CONFIG.items(): - print(f"\n📞 正在调用API: {api_key}") - print(f" URL: {config['url']}") - print(f" 方法: {config['method']}") - print(f" 超时: {config['timeout']}秒") - - # 执行API调用任务 - try: - result = invoke_call_api_task.delay(api_key=api_key) - - # 等待任务完成并获取结果 - task_result = result.get(timeout=config['timeout'] + 10) - - print(f" 结果: {task_result['status']}") - print(f" 消息: {task_result['message']}") - if task_result['status'] == 'success': - print(f" ✅ 成功") - else: - print(f" ❌ 失败") - - results.append({ - 'api_key': api_key, - 'status': task_result['status'], - 'result': task_result - }) - - except Exception as e: - print(f" ❌ 异常: {str(e)}") - results.append({ - 'api_key': api_key, - 'status': 'error', - 'error': str(e) - }) - - # 短暂延迟避免请求过于频繁 - time.sleep(1) - - # 汇总结果 - print("\n" + "=" * 60) - print("📊 任务执行汇总:") - - success_count = len([r for r in results if r['status'] == 'success']) - total_count = len(results) - - for result in results: - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" {status_icon} {result['api_key']}: {result['status']}") - - print(f"\n总计: {success_count}/{total_count} 成功") - print("🎯 任务执行完成") - -if __name__ == "__main__": - main() \ No newline at end of file diff --git a/test_beat_schedule.py b/test_beat_schedule.py deleted file mode 100644 index afa204d..0000000 --- a/test_beat_schedule.py +++ /dev/null @@ -1,53 +0,0 @@ -#!/usr/bin/env python3 -"""测试新的调度配置 - 运行2次观察执行情况""" - -import subprocess -import time -import sys -import threading -from datetime import datetime - -def run_beat(): - """运行Celery Beat调度器""" - try: - print("🚀 启动Celery Beat调度器...") - process = subprocess.Popen([ - sys.executable, '-m', 'celery', '-A', 'celery_app', 'beat', '--loglevel=info' - ], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1) - - # 实时输出日志 - for line in iter(process.stdout.readline, ''): - print(f"[BEAT] {line.strip()}") - - except KeyboardInterrupt: - print("\n🛑 停止Beat调度器") - process.terminate() - except Exception as e: - print(f"❌ Beat启动失败: {e}") - -def test_schedule(): - """测试新的调度配置""" - print("=" * 60) - print(f"测试新的调度配置 - 每2分钟执行一次") - print(f"开始时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}") - print("=" * 60) - print() - - print("⚠️ 注意: 这是测试模式,将运行大约4分钟来观察2次任务执行") - print("💡 如果Redis没有运行,请先启动Redis服务") - print() - - # 启动Beat调度器 - beat_thread = threading.Thread(target=run_beat) - beat_thread.daemon = True - beat_thread.start() - - # 等待4分钟观察执行 - try: - time.sleep(240) # 4分钟,应该能看到2次任务执行 - print("\n✅ 测试完成!") - except KeyboardInterrupt: - print("\n🛑 手动停止测试") - -if __name__ == "__main__": - test_schedule() \ No newline at end of file diff --git a/test_body_config.py b/test_body_config.py deleted file mode 100644 index 9079854..0000000 --- a/test_body_config.py +++ /dev/null @@ -1,84 +0,0 @@ -#!/usr/bin/env python3 -"""测试body配置功能""" -from tasks import daily_morning_task -from invoke_call_api_task import generic_api_call -import api_config -import time -import json - -def test_body_configuration(): - """测试不同API配置的body参数""" - print("=== 测试Body配置功能 ===") - - # 测试不同的API配置 - test_apis = ['main_api', 'backup_api', 'get_api', 'put_api'] - - for api_key in test_apis: - print(f"\n--- 测试 {api_key} 配置 ---") - - # 显示配置信息 - config = api_config.API_CONFIG[api_key] - print(f"URL: {config['url']}") - print(f"Method: {config['method']}") - print(f"Body: {json.dumps(config.get('body', {}), indent=2, ensure_ascii=False)}") - - # 执行任务 - result = generic_api_call.delay(api_key) - print(f"任务ID: {result.id}") - - # 等待完成 - for i in range(10): - try: - if result.ready(): - print(f"任务状态: {result.status}") - if result.successful(): - task_result = result.get() - print(f"✅ 任务成功: {task_result.get('message')}") - print(f"发送的Body: {json.dumps(api_config.get('body', {}), indent=2, ensure_ascii=False)}") - else: - print(f"❌ 任务失败: {result.result}") - break - time.sleep(1) - except Exception as e: - print(f"检查任务状态时出错: {e}") - time.sleep(1) - - time.sleep(2) # 等待一下再进行下一个测试 - - print("\n✅ 所有body配置测试完成") - -def test_custom_body(): - """测试自定义body参数""" - print("\n=== 测试自定义Body参数 ===") - - custom_payload = { - 'title': '自定义任务', - 'body': '这是自定义的请求体数据', - 'userId': 999, - 'extra_field': '额外字段' - } - - print(f"自定义Body: {json.dumps(custom_payload, indent=2, ensure_ascii=False)}") - - result = generic_api_call.delay('main_api', custom_payload) - print(f"任务ID: {result.id}") - - # 等待完成 - for i in range(10): - try: - if result.ready(): - print(f"任务状态: {result.status}") - if result.successful(): - task_result = result.get() - print(f"✅ 自定义任务成功: {task_result.get('message')}") - else: - print(f"❌ 自定义任务失败: {result.result}") - break - time.sleep(1) - except Exception as e: - print(f"检查任务状态时出错: {e}") - time.sleep(1) - -if __name__ == '__main__': - test_body_configuration() - test_custom_body() \ No newline at end of file diff --git a/test_call_api.py b/test_call_api.py deleted file mode 100644 index 4ebfcc9..0000000 --- a/test_call_api.py +++ /dev/null @@ -1,47 +0,0 @@ -#!/usr/bin/env python3 -"""测试呼叫任务API功能""" - -import sys -import os -sys.path.append(os.path.dirname(os.path.abspath(__file__))) - -from invoke_call_api_task import generic_api_call -from datetime import datetime - -def test_call_api_task(): - """测试呼叫任务API功能""" - print("=" * 60) - print(f"测试呼叫任务API功能 - {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}") - print("=" * 60) - - try: - # 测试默认呼叫任务API - print("🔧 测试默认呼叫任务API (main_api):") - result1 = generic_api_call() - print(f"✅ 默认呼叫结果: {result1['status']} - {result1['message']}") - - print("\n🔧 测试备用呼叫任务API (backup_api):") - result2 = generic_api_call(api_key='backup_api') - print(f"✅ 备用呼叫结果: {result2['status']} - {result2['message']}") - - print("\n🔧 测试自定义呼叫参数:") - custom_payload = { - 'title': '呼叫任务测试', - 'body': '这是自定义的呼叫任务参数', - 'userId': 888 - } - result3 = generic_api_call(api_key='main_api', custom_payload=custom_payload) - print(f"✅ 自定义呼叫结果: {result3['status']} - {result3['message']}") - - print("\n" + "=" * 60) - print("📊 呼叫任务API测试完成!") - print("=" * 60) - - return [result1, result2, result3] - - except Exception as e: - print(f"❌ 测试失败: {str(e)}") - return None - -if __name__ == "__main__": - test_call_api_task() \ No newline at end of file diff --git a/test_celery_api_list.py b/test_celery_api_list.py deleted file mode 100644 index ef65979..0000000 --- a/test_celery_api_list.py +++ /dev/null @@ -1,73 +0,0 @@ -#!/usr/bin/env python3 -"""使用Celery任务遍历API配置列表执行调用""" - -from invoke_call_api_task import invoke_call_api_task -from api_config import API_CONFIG - -def test_celery_api_list(): - """使用Celery异步任务遍历所有API配置""" - print("🚀 使用Celery任务遍历API配置列表") - print("=" * 50) - - results = [] - task_ids = [] - - # 提交所有任务 - for api_key in API_CONFIG.keys(): - print(f"📤 提交任务: {api_key}") - - # 提交异步任务 - task_result = invoke_call_api_task.delay(api_key=api_key) - task_ids.append((api_key, task_result)) - - results.append({ - 'api_key': api_key, - 'task_id': task_result.id, - 'status': 'submitted' - }) - - print(f"\n⏳ 等待 {len(task_ids)} 个任务完成...") - - # 等待所有任务完成 - for api_key, task_result in task_ids: - try: - # 等待任务结果 - result = task_result.get(timeout=60) - - # 更新结果 - for r in results: - if r['api_key'] == api_key: - r['status'] = result['status'] - r['message'] = result['message'] - r['task_result'] = result - break - - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" {status_icon} {api_key}: {result['status']}") - - except Exception as e: - for r in results: - if r['api_key'] == api_key: - r['status'] = 'error' - r['message'] = str(e) - break - print(f" ❌ {api_key}: 错误 - {str(e)[:30]}") - - # 汇总 - print("\n" + "=" * 50) - print("📊 Celery任务执行汇总:") - - success_count = len([r for r in results if r['status'] == 'success']) - total_count = len(results) - - for result in results: - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" {status_icon} {result['api_key']}: {result['status']}") - - print(f"\n🎯 总计: {success_count}/{total_count} 成功") - print("✨ 所有任务执行完成") - - return results - -if __name__ == "__main__": - test_celery_api_list() \ No newline at end of file diff --git a/test_direct_api.py b/test_direct_api.py deleted file mode 100644 index e1093a3..0000000 --- a/test_direct_api.py +++ /dev/null @@ -1,129 +0,0 @@ -#!/usr/bin/env python3 -"""直接测试API调用(不使用Celery)""" - -import requests -from datetime import datetime -from api_config import API_CONFIG, RETRY_CONFIG -import logging - -logging.basicConfig(level=logging.INFO) -logger = logging.getLogger(__name__) - -def call_api_direct(api_key='main_api', custom_payload=None): - """直接调用API(复制invoke_call_api_task的核心逻辑)""" - 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配置 - call_config = API_CONFIG[api_key] - call_url = call_config['url'] - method = call_config['method'].upper() - headers = call_config['headers'] - timeout = call_config['timeout'] - - # 构建呼叫参数 - if custom_payload is None: - payload = call_config['body'] - else: - payload = custom_payload - - 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成功 - {method} {call_url}") - return { - 'status': 'success', - 'message': f'呼叫任务API调用成功 ({method})', - 'response': result_data, - 'timestamp': current_time, - 'api_key': api_key - } - else: - logger.error(f"呼叫任务API失败,状态码: {response.status_code}, 响应: {response.text}") - print(f"[{current_time}] ❌ 呼叫任务API失败 - 状态码: {response.status_code}") - return { - 'status': 'failed', - 'message': f'呼叫任务API失败,状态码: {response.status_code}', - 'response': response.text, - 'timestamp': current_time, - 'api_key': api_key - } - - except requests.exceptions.RequestException as e: - logger.error(f"呼叫任务API异常: {str(e)}") - print(f"[{current_time}] ❌ 呼叫任务API异常 - {str(e)}") - return { - 'status': 'error', - 'message': f'呼叫任务API异常: {str(e)}', - 'timestamp': current_time, - 'api_key': api_key - } - -def main(): - """遍历所有API配置并执行调用""" - print("🚀 开始遍历API配置列表执行直接调用") - print("=" * 60) - - results = [] - - for api_key, config in API_CONFIG.items(): - print(f"\n📞 正在调用API: {api_key}") - print(f" URL: {config['url']}") - print(f" 方法: {config['method']}") - print(f" 超时: {config['timeout']}秒") - - # 直接调用API函数 - result = call_api_direct(api_key=api_key) - - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" 结果: {status_icon} {result['status']}") - print(f" 消息: {result['message']}") - - results.append({ - 'api_key': api_key, - 'status': result['status'], - 'message': result['message'], - 'result': result - }) - - # 汇总结果 - print("\n" + "=" * 60) - print("📊 任务执行汇总:") - - success_count = len([r for r in results if r['status'] == 'success']) - total_count = len(results) - - for result in results: - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" {status_icon} {result['api_key']}: {result['status']}") - - print(f"\n🎯 总计: {success_count}/{total_count} 成功") - print("✨ 任务执行完成") - -if __name__ == "__main__": - main() \ No newline at end of file diff --git a/test_merged_status.py b/test_merged_status.py deleted file mode 100644 index 71a816b..0000000 --- a/test_merged_status.py +++ /dev/null @@ -1,77 +0,0 @@ -#!/usr/bin/env python3 -"""测试合并状态码功能""" - -import sys -import os -sys.path.append(os.path.dirname(os.path.abspath(__file__))) - -from datetime import datetime -from api_config import API_CONFIG, RETRY_CONFIG -import requests -import logging - -# 配置日志 -logging.basicConfig(level=logging.INFO) -logger = logging.getLogger(__name__) - -def test_status_code_handling(): - """测试不同状态码的处理""" - print("测试合并状态码 200, 201, 204 的处理") - print("=" * 40) - - current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S') - - # 测试几个主要的API - test_apis = ['main_api', 'health_check', 'get_api'] - - for api_name in test_apis: - if api_name in API_CONFIG: - api_config = API_CONFIG[api_name] - api_url = api_config['url'] - headers = api_config['headers'] - timeout = api_config['timeout'] - payload = api_config['body'] - method = api_config['method'].upper() - - try: - print(f"\n🔧 测试 {api_name}: {method} {api_url}") - - # 根据method选择请求方式 - if method == 'POST': - response = requests.post(api_url, headers=headers, json=payload, timeout=timeout) - elif method == 'GET': - response = requests.get(api_url, headers=headers, params=payload, timeout=timeout) - elif method == 'PUT': - response = requests.put(api_url, headers=headers, json=payload, timeout=timeout) - elif method == 'DELETE': - response = requests.delete(api_url, headers=headers, timeout=timeout) - - print(f"📊 状态码: {response.status_code}") - - # 使用新的合并状态码逻辑 - if response.status_code in [200, 201, 204]: # 成功状态码 - try: - result_data = response.json() - except: - result_data = {'response': response.text} - - # 构建成功消息 - if response.status_code == 201: - print(f"✅ 成功 (创建) - ID: {result_data.get('id')}") - elif response.status_code == 200: - print(f"✅ 成功 (OK) - 数据长度: {len(str(result_data))}") - elif response.status_code == 204: - print(f"✅ 成功 (无内容)") - - else: - print(f"❌ 失败 - 状态码: {response.status_code}") - - except requests.exceptions.RequestException as e: - print(f"❌ 网络异常: {str(e)}") - continue # 继续测试下一个API - except Exception as e: - print(f"❌ 其他异常: {str(e)}") - continue # 继续测试下一个API - -if __name__ == "__main__": - test_status_code_handling() \ No newline at end of file diff --git a/test_methods.py b/test_methods.py deleted file mode 100644 index fa85cff..0000000 --- a/test_methods.py +++ /dev/null @@ -1,53 +0,0 @@ -#!/usr/bin/env python3 -"""测试不同HTTP方法的API调用""" -from tasks import daily_morning_task -import api_config -import time - -def test_different_methods(): - """测试不同HTTP方法的API调用""" - print("=== 测试不同HTTP方法 ===") - - # 保存原始配置 - original_config = api_config.API_CONFIG['main_api'].copy() - - test_cases = [ - ('POST', 'https://jsonplaceholder.typicode.com/posts'), - ('GET', 'https://jsonplaceholder.typicode.com/posts/1'), - ('PUT', 'https://jsonplaceholder.typicode.com/posts/1'), - ] - - for method, url in test_cases: - print(f"\n--- 测试 {method} 方法 ---") - - # 临时修改配置 - api_config.API_CONFIG['main_api']['method'] = method - api_config.API_CONFIG['main_api']['url'] = url - - # 执行任务 - result = daily_morning_task.delay() - print(f"任务ID: {result.id}") - - # 等待完成 - for i in range(10): - try: - if result.ready(): - print(f"任务状态: {result.status}") - if result.successful(): - print(f"任务结果: {result.get()}") - else: - print(f"任务失败: {result.result}") - break - time.sleep(1) - except Exception as e: - print(f"检查任务状态时出错: {e}") - time.sleep(1) - - time.sleep(2) # 等待一下再进行下一个测试 - - # 恢复原始配置 - api_config.API_CONFIG['main_api'] = original_config - print("\n✅ 所有测试完成,配置已恢复") - -if __name__ == '__main__': - test_different_methods() \ No newline at end of file diff --git a/test_quick_api.py b/test_quick_api.py deleted file mode 100644 index 941aa34..0000000 --- a/test_quick_api.py +++ /dev/null @@ -1,45 +0,0 @@ -#!/usr/bin/env python3 -"""快速测试API列表调用""" - -import requests -from api_config import API_CONFIG -from datetime import datetime - -def quick_test(): - print("🚀 快速测试API列表调用") - print("=" * 40) - - success = 0 - total = 0 - - for api_key, config in API_CONFIG.items(): - total += 1 - print(f"\n📞 {api_key} ({config['method']})") - - try: - if config['method'] == 'GET': - resp = requests.get(config['url'], headers=config['headers'], - params=config['body'], timeout=config['timeout']) - elif config['method'] == 'POST': - resp = requests.post(config['url'], headers=config['headers'], - json=config['body'], timeout=config['timeout']) - elif config['method'] == 'PUT': - resp = requests.put(config['url'], headers=config['headers'], - json=config['body'], timeout=config['timeout']) - else: - resp = requests.delete(config['url'], headers=config['headers'], - timeout=config['timeout']) - - if resp.status_code in [200, 201, 204]: - print(f" ✅ 成功 - 状态码: {resp.status_code}") - success += 1 - else: - print(f" ❌ 失败 - 状态码: {resp.status_code}") - - except Exception as e: - print(f" ❌ 异常: {str(e)[:50]}") - - print(f"\n📊 结果: {success}/{total} 成功") - -if __name__ == "__main__": - quick_test() \ No newline at end of file diff --git a/test_quick_updated.py b/test_quick_updated.py deleted file mode 100644 index 1fed842..0000000 --- a/test_quick_updated.py +++ /dev/null @@ -1,32 +0,0 @@ -#!/usr/bin/env python3 -"""快速测试更新后的invoke_call_api_task方法""" - -import logging -logging.basicConfig(level=logging.INFO) - -from invoke_call_api_task import invoke_call_api_task -from api_config import API_CONFIG - -def quick_test(): - print("🚀 快速测试更新后的invoke_call_api_task") - print("=" * 40) - - # 创建任务实例 - task_instance = invoke_call_api_task() - - # 执行任务 - result = task_instance.run() - - print(f"\n📊 执行汇总:") - print(f" 状态: {result['status']}") - print(f" 总数: {result['total_apis']}") - print(f" 成功: {result['success_count']}") - print(f" 失败: {result['failed_count']}") - - print(f"\n📋 API调用结果:") - for api_result in result['results']: - icon = "✅" if api_result['status'] == 'success' else "❌" - print(f" {icon} {api_result['call_api_key']} - {api_result['status']}") - -if __name__ == "__main__": - quick_test() \ No newline at end of file diff --git a/test_retry.py b/test_retry.py deleted file mode 100644 index 4c425c1..0000000 --- a/test_retry.py +++ /dev/null @@ -1,38 +0,0 @@ -#!/usr/bin/env python3 -"""测试重试机制脚本""" -from tasks import daily_morning_task -import time - -def test_retry_mechanism(): - """测试重试机制""" - print("=== 测试API重试机制 ===") - print("启动任务测试...") - - # 手动触发任务 - result = daily_morning_task.delay() - - print(f"任务ID: {result.id}") - print("等待任务执行完成...") - - # 监控任务状态 - for i in range(30): # 最多等待30秒 - try: - status = result.status - if status in ['SUCCESS', 'FAILURE']: - print(f"\n任务最终状态: {status}") - if status == 'SUCCESS': - print(f"任务结果: {result.get()}") - else: - print(f"任务失败信息: {result.result}") - break - else: - print(f"任务状态: {status}... (等待中)") - time.sleep(1) - except Exception as e: - print(f"检查任务状态时出错: {e}") - time.sleep(1) - - print("测试完成") - -if __name__ == '__main__': - test_retry_mechanism() \ No newline at end of file diff --git a/test_schedule.py b/test_schedule.py deleted file mode 100644 index 93bee4b..0000000 --- a/test_schedule.py +++ /dev/null @@ -1,39 +0,0 @@ -#!/usr/bin/env python3 -"""测试调度配置""" - -import sys -import os -sys.path.append(os.path.dirname(os.path.abspath(__file__))) - -from celery_app import app -from datetime import datetime - -def check_schedule_config(): - """检查当前的调度配置""" - print("当前Celery调度配置:") - print("=" * 40) - - beat_schedule = app.conf.beat_schedule - - for task_name, task_config in beat_schedule.items(): - schedule_seconds = task_config['schedule'] - task_path = task_config['task'] - - # 转换为可读格式 - if schedule_seconds < 60: - interval = f"{schedule_seconds}秒" - elif schedule_seconds < 3600: - minutes = schedule_seconds / 60 - interval = f"{minutes:.1f}分钟" - else: - hours = schedule_seconds / 3600 - interval = f"{hours:.1f}小时" - - print(f"🕐 任务名称: {task_name}") - print(f"📝 任务路径: {task_path}") - print(f"⏰ 执行间隔: {interval} ({schedule_seconds}秒)") - print(f"🔧 队列: {task_config.get('options', {}).get('queue', 'default')}") - print() - -if __name__ == "__main__": - check_schedule_config() \ No newline at end of file diff --git a/test_simple_api_list.py b/test_simple_api_list.py deleted file mode 100644 index b9cc69f..0000000 --- a/test_simple_api_list.py +++ /dev/null @@ -1,67 +0,0 @@ -#!/usr/bin/env python3 -"""简单测试:遍历API配置列表执行调用(同步方式)""" - -from invoke_call_api_task import invoke_call_api_task -from api_config import API_CONFIG -import time - -def main(): - """直接调用函数遍历所有API配置""" - print("🚀 开始遍历API配置列表执行调用") - print("=" * 50) - - results = [] - - for api_key, config in API_CONFIG.items(): - print(f"\n📞 正在调用API: {api_key}") - print(f" URL: {config['url']}") - print(f" 方法: {config['method']}") - - try: - # 使用异步方式执行Celery任务 - from celery import current_app - - # 提交异步任务 - async_result = invoke_call_api_task.delay(api_key=api_key) - - # 等待任务完成 - result = async_result.get(timeout=60) - - print(f" 状态: {result['status']}") - print(f" 消息: {result['message']}") - - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" 结果: {status_icon} {result['status']}") - - results.append({ - 'api_key': api_key, - 'status': result['status'], - 'message': result['message'] - }) - - except Exception as e: - print(f" ❌ 异常: {str(e)}") - results.append({ - 'api_key': api_key, - 'status': 'error', - 'message': str(e) - }) - - # 短暂延迟 - time.sleep(1) - - # 汇总 - print("\n" + "=" * 50) - print("📊 执行汇总:") - - success_count = len([r for r in results if r['status'] == 'success']) - total_count = len(results) - - for result in results: - status_icon = "✅" if result['status'] == 'success' else "❌" - print(f" {status_icon} {result['api_key']}: {result['message']}") - - print(f"\n🎯 总计: {success_count}/{total_count} 成功") - -if __name__ == "__main__": - main() \ No newline at end of file diff --git a/test_simulate_failure.py b/test_simulate_failure.py deleted file mode 100644 index cf466f1..0000000 --- a/test_simulate_failure.py +++ /dev/null @@ -1,58 +0,0 @@ -#!/usr/bin/env python3 -"""模拟API失败情况,测试重试机制""" - -def modify_config_for_test(): - """临时修改配置以测试重试机制""" - import api_config - - # 保存原始配置 - original_url = api_config.API_CONFIG['main_api']['url'] - - # 修改为无效URL以模拟失败 - api_config.API_CONFIG['main_api']['url'] = 'https://nonexistent-api.example.com/posts' - - print(f"已临时修改API URL为: {api_config.API_CONFIG['main_api']['url']}") - print("现在运行测试任务将触发重试机制") - print("测试完成后会恢复原始配置") - - return original_url - -def restore_config(original_url): - """恢复原始配置""" - import api_config - api_config.API_CONFIG['main_api']['url'] = original_url - print(f"已恢复原始API URL: {original_url}") - -def test_retry_with_failure(): - """测试失败情况下的重试机制""" - original_url = modify_config_for_test() - - try: - from tasks import daily_morning_task - import time - - print("\n=== 启动失败重试测试 ===") - result = daily_morning_task.delay() - print(f"任务ID: {result.id}") - - # 监控重试过程 - for i in range(60): # 等待足够长时间观察重试 - try: - status = result.status - if status in ['SUCCESS', 'FAILURE']: - print(f"\n任务最终状态: {status}") - print(f"任务结果: {result.result}") - break - else: - print(f"任务状态: {status}... (重试中)") - time.sleep(2) - except Exception as e: - print(f"检查任务状态时出错: {e}") - time.sleep(2) - - finally: - restore_config(original_url) - print("\n配置已恢复,测试完成") - -if __name__ == '__main__': - test_retry_with_failure() \ No newline at end of file diff --git a/test_task.py b/test_task.py deleted file mode 100644 index 89f3e75..0000000 --- a/test_task.py +++ /dev/null @@ -1,10 +0,0 @@ -#!/usr/bin/env python3 -"""测试脚本 - 手动触发任务""" -from tasks import daily_morning_task - -if __name__ == '__main__': - print("手动触发定时任务进行测试...") - result = daily_morning_task.delay() - print(f"任务已提交,任务ID: {result.id}") - print("等待任务执行结果...") - print(f"任务结果: {result.get(timeout=10)}") \ No newline at end of file diff --git a/test_updated_invoke_call.py b/test_updated_invoke_call.py deleted file mode 100644 index 1f56093..0000000 --- a/test_updated_invoke_call.py +++ /dev/null @@ -1,39 +0,0 @@ -#!/usr/bin/env python3 -"""测试修改后的invoke_call_api_task方法""" - -import logging -logging.basicConfig(level=logging.INFO) - -from invoke_call_api_task import invoke_call_api_task - -def test_updated_method(): - """测试更新后的invoke_call_api_task方法""" - print("🚀 测试修改后的invoke_call_api_task方法") - print("=" * 50) - - try: - # 直接调用任务实例来测试逻辑 - task_instance = invoke_call_api_task() - result = task_instance.run() - - print(f"\n📊 执行结果:") - print(f" 总体状态: {result['status']}") - print(f" 总API数量: {result['total_apis']}") - print(f" 成功数量: {result['success_count']}") - print(f" 失败数量: {result['failed_count']}") - print(f" 时间戳: {result['timestamp']}") - - print(f"\n📋 详细结果:") - for api_result in result['results']: - status_icon = "✅" if api_result['status'] == 'success' else "❌" - print(f" {status_icon} {api_result['call_api_key']}: {api_result['status']}") - - print(f"\n🎯 任务执行完成") - return result - - except Exception as e: - print(f"❌ 测试失败: {str(e)}") - return None - -if __name__ == "__main__": - test_updated_method() \ No newline at end of file