diff --git a/README.md b/README.md index a9852ab..bee9002 100644 --- a/README.md +++ b/README.md @@ -1,3 +1,106 @@ -# task-executer +# Celery 定时任务项目 + Flower监控 -定时调用外呼任务api \ No newline at end of file +## 项目结构 +- `celery_app.py`: Celery应用配置 +- `tasks.py`: 定义任务函数 +- `scheduler.py`: 定时任务调度配置 +- `run_worker.py`: 启动Worker的脚本 +- `run_beat.py`: 启动Beat调度器的脚本 +- `run_flower.py`: 启动Flower监控的脚本 +- `flower_config.py`: Flower监控配置 +- `start_all.py`: 一键启动所有服务 +- `test_task.py`: 测试脚本 + +## 安装依赖 +```bash +pip install -r requirements.txt +``` + +## 启动Redis +确保Redis服务已启动,默认监听localhost:6379 + +## 运行项目 + +### 方式一:一键启动(推荐) +```bash +python start_all.py +``` + +### 方式二:分别启动 +1. 启动Celery Worker +```bash +python run_worker.py +``` + +2. 启动Celery Beat调度器 +```bash +python run_beat.py +``` + +3. 启动Flower监控 +```bash +python run_flower.py +``` + +## Flower监控 +访问 http://localhost:5555 查看监控面板,包含: +- 📊 实时任务状态 +- 📈 Worker状态和性能 +- ⏰ 定时任务执行历史 +- 🔍 任务执行详情 +- 📋 队列监控 + +## 定时任务配置 +任务会在每天上午9点准时执行,执行时会打印日志信息。 + +## 测试任务 + +### 1. 测试API接口 +```bash +python test_api.py +``` + +### 2. 测试重试机制 +```bash +# 正常重试测试 +python test_retry.py + +# 模拟失败场景测试 +python simulate_failure.py +``` + +### 3. 测试HTTP方法 +```bash +# 测试不同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 +``` + +### 4. 测试Body配置 +```bash +# 测试不同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) +``` + +### 3. 手动触发定时任务 +```python +from tasks import daily_morning_task +result = daily_morning_task.delay() +print(result.get()) +``` + +## 监控功能 +- 实时查看任务执行状态 +- 监控Worker性能指标 +- 查看任务执行历史和日志 +- 管理定时任务 +- 监控队列状态 \ No newline at end of file diff --git a/api_config.py b/api_config.py new file mode 100644 index 0000000..ec1d20f --- /dev/null +++ b/api_config.py @@ -0,0 +1,85 @@ +"""API接口配置文件""" + +# API配置 +API_CONFIG = { + # 主要API接口 + 'main_api': { + 'url': 'https://jsonplaceholder.typicode.com/posts', + 'method': 'POST', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + 'title': '每日定时任务', + 'body': '这是通过配置文件设置的请求体数据', + 'userId': 1 + } + }, + + # 备用API接口 + 'backup_api': { + 'url': 'https://jsonplaceholder.typicode.com/comments', + 'method': 'POST', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + 'name': '备用任务', + 'email': 'backup@example.com', + 'body': '这是备用API的请求体' + } + }, + + # 健康检查API + 'health_check': { + 'url': 'https://jsonplaceholder.typicode.com/posts/1', + 'method': 'GET', + 'headers': { + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 10, + 'body': {} # GET请求通常不需要body,但可以设置查询参数 + }, + + # GET请求示例API + 'get_api': { + 'url': 'https://jsonplaceholder.typicode.com/posts', + 'method': 'GET', + 'headers': { + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + 'userId': 1, # 这将作为查询参数 + 'limit': 10 + } + }, + + # PUT请求示例API + 'put_api': { + 'url': 'https://jsonplaceholder.typicode.com/posts/1', + 'method': 'PUT', + 'headers': { + 'Content-Type': 'application/json', + 'User-Agent': 'Celery-Task/1.0' + }, + 'timeout': 30, + 'body': { + 'id': 1, + 'title': '更新后的标题', + 'body': '这是通过PUT请求更新的内容', + 'userId': 1 + } + } +} + +# 重试配置 +RETRY_CONFIG = { + 'max_retries': 3, + 'retry_delay': 60, # 重试间隔(秒) + 'retry_on_status': [500, 502, 503, 504, 429] +} \ No newline at end of file diff --git a/celery_app.py b/celery_app.py new file mode 100644 index 0000000..174e7a1 --- /dev/null +++ b/celery_app.py @@ -0,0 +1,31 @@ +from celery import Celery +from datetime import datetime +import logging + +# 配置日志 +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + +# 创建Celery应用 +app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0') + +# 配置Celery +app.conf.update( + timezone='Asia/Shanghai', + enable_utc=True, + broker_url='redis://localhost:6379/0', + result_backend='redis://localhost:6379/0', + task_serializer='json', + accept_content=['json'], + result_serializer='json', + beat_schedule={ + 'daily-morning-task': { + 'task': 'tasks.daily_morning_task', + 'schedule': 60.0 * 2.0, # 每2分钟执行一次 + 'options': {'queue': 'default'} + }, + }, + # Flower监控配置 + worker_send_task_events=True, + task_send_sent_event=True, +) \ No newline at end of file diff --git a/check_configs.py b/check_configs.py new file mode 100644 index 0000000..c1f722a --- /dev/null +++ b/check_configs.py @@ -0,0 +1,23 @@ +#!/usr/bin/env python3 +"""简化的API列表测试""" + +import sys +import os +sys.path.append(os.path.dirname(os.path.abspath(__file__))) + +from api_config import API_CONFIG + +def check_api_configs(): + """检查配置文件中的API配置""" + print("配置文件中的API配置列表:") + print("=" * 30) + + for api_name, api_config in API_CONFIG.items(): + print(f"🔧 {api_name}:") + print(f" URL: {api_config['url']}") + print(f" Method: {api_config['method']}") + print(f" Timeout: {api_config['timeout']}") + print() + +if __name__ == "__main__": + check_api_configs() \ No newline at end of file diff --git a/flower_config.py b/flower_config.py new file mode 100644 index 0000000..a5d5a44 --- /dev/null +++ b/flower_config.py @@ -0,0 +1,20 @@ +"""Flower监控配置""" + +# Flower监控配置 +FLOWER_CONFIG = { + 'broker_url': 'redis://localhost:6379/0', + 'result_backend': 'redis://localhost:6379/0', + 'port': 5555, + 'address': '0.0.0.0', # 允许外部访问 + 'basic_auth': None, # 基础认证,格式: 'username:password' + 'oauth_redirect_url': None, + 'db': None, # 使用SQLite存储监控数据 + 'inspect_timeout': 1000, + 'purge_offline_workers': 60, + 'max_tasks': 10000, + 'max_workers': 5000, + 'format_task': 'json', + 'enable_events': True, + 'tasks_columns': ['uuid', 'name', 'args', 'kwargs', 'state', 'runtime', 'worker', 'timestamp'], + 'workers_columns': ['hostname', 'pid', 'sw_ver', 'loadavg', 'active', 'processed', 'failed', 'status', 'last_heartbeat'], +} \ No newline at end of file diff --git a/generic_api_task.py b/generic_api_task.py new file mode 100644 index 0000000..0a26767 --- /dev/null +++ b/generic_api_task.py @@ -0,0 +1,120 @@ +"""呼叫任务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, max_retries=None) +def generic_api_call(self, api_key='main_api', custom_payload=None): + """调用呼叫任务API + + Args: + api_key: 呼叫任务API配置的键名(如 'call_api', 'voice_api') + custom_payload: 自定义呼叫参数,如果为None则使用默认参数 + """ + 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: + # 直接使用配置文件中的body作为呼叫参数 + 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, + 'retry_count': getattr(self.request, 'retries', 0), + 'call_api_key': api_key + } + 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}) - 状态码: {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失败,已达到最大重试次数 - 状态码: {response.status_code}") + return { + 'status': 'failed', + 'message': f'呼叫任务API失败,已达到最大重试次数: {max_retries}', + 'response': response.text, + 'timestamp': current_time, + 'retry_count': current_retry + 1, + 'call_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, + 'retry_count': getattr(self.request, 'retries', 0), + 'call_api_key': api_key + } + + 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}) - 异常: {str(e)}") + raise self.retry(countdown=retry_delay, exc=e) + else: + logger.error(f"呼叫任务API异常,已达到最大重试次数: {max_retries}") + print(f"[{current_time}] ❌ 呼叫任务API异常,已达到最大重试次数 - 异常: {str(e)}") + return { + 'status': 'error', + 'message': f'呼叫任务API异常,已达到最大重试次数: {max_retries}', + 'error': str(e), + 'timestamp': current_time, + 'retry_count': current_retry + 1, + 'call_api_key': api_key + } \ No newline at end of file diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..9618557 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,4 @@ +redis==5.0.1 +celery==5.3.4 +flower==2.0.1 +requests==2.31.0 \ No newline at end of file diff --git a/run_beat.py b/run_beat.py new file mode 100644 index 0000000..d478c03 --- /dev/null +++ b/run_beat.py @@ -0,0 +1,6 @@ +#!/usr/bin/env python3 +from celery_app import app + +if __name__ == '__main__': + print("启动Celery Beat调度器...") + app.worker_main(['beat', '--loglevel=info']) \ No newline at end of file diff --git a/run_flower.py b/run_flower.py new file mode 100644 index 0000000..4830805 --- /dev/null +++ b/run_flower.py @@ -0,0 +1,25 @@ +#!/usr/bin/env python3 +"""启动Flower监控服务""" +from celery import Celery +from flower import command as flower_command + +def start_flower(): + """启动Flower监控""" + print("启动Flower监控服务...") + print("监控地址: http://localhost:5555") + print("按 Ctrl+C 停止服务") + + # Flower配置 + flower_options = { + 'broker': 'redis://localhost:6379/0', + 'port': 5555, + 'basic_auth': None, # 可以设置用户名:密码,如 'admin:password' + 'inspect_timeout': 1000, + 'purge_offline_workers': 60, + } + + # 启动Flower + flower_command.flower_command(flower_options) + +if __name__ == '__main__': + start_flower() \ No newline at end of file diff --git a/run_worker.py b/run_worker.py new file mode 100644 index 0000000..440e1e0 --- /dev/null +++ b/run_worker.py @@ -0,0 +1,6 @@ +#!/usr/bin/env python3 +from celery_app import app + +if __name__ == '__main__': + print("启动Celery Worker...") + app.worker_main(['worker', '--loglevel=info','--pool=solo','--concurrency=1']) \ No newline at end of file diff --git a/scheduler.py b/scheduler.py new file mode 100644 index 0000000..0d8f73d --- /dev/null +++ b/scheduler.py @@ -0,0 +1,11 @@ +from celery.schedules import crontab +from celery_app import app + +# 使用crontab设置精确的定时任务 +app.conf.beat_schedule = { + 'daily-morning-9am': { + 'task': 'tasks.daily_morning_task', + 'schedule': crontab(hour=9, minute=0), # 每天上午9点执行 + 'options': {'queue': 'default'} + }, +} \ No newline at end of file diff --git a/simulate_failure.py b/simulate_failure.py new file mode 100644 index 0000000..cf466f1 --- /dev/null +++ b/simulate_failure.py @@ -0,0 +1,58 @@ +#!/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/start_all.py b/start_all.py new file mode 100644 index 0000000..43fce57 --- /dev/null +++ b/start_all.py @@ -0,0 +1,67 @@ +#!/usr/bin/env python3 +"""一键启动所有服务:Celery Worker + Beat + Flower""" +import subprocess +import sys +import time +import signal +import os + +def start_services(): + """启动所有服务""" + processes = [] + + try: + print("正在启动Celery服务...") + + # 启动Worker + print("1. 启动Celery Worker...") + worker_proc = subprocess.Popen([ + sys.executable, 'run_worker.py' + ], stdout=subprocess.PIPE, stderr=subprocess.STDOUT) + processes.append(('Worker', worker_proc)) + + # 等待Worker启动 + time.sleep(2) + + # 启动Beat + print("2. 启动Celery Beat调度器...") + beat_proc = subprocess.Popen([ + sys.executable, 'run_beat.py' + ], stdout=subprocess.PIPE, stderr=subprocess.STDOUT) + processes.append(('Beat', beat_proc)) + + # 等待Beat启动 + time.sleep(2) + + # 启动Flower + print("3. 启动Flower监控服务...") + flower_proc = subprocess.Popen([ + sys.executable, 'run_flower.py' + ], stdout=subprocess.PIPE, stderr=subprocess.STDOUT) + processes.append(('Flower', flower_proc)) + + print("\n✅ 所有服务已启动!") + print("🌸 Flower监控地址: http://localhost:5555") + print("📋 查看定时任务状态和执行历史") + print("\n按 Ctrl+C 停止所有服务") + + # 等待中断信号 + while True: + time.sleep(1) + + except KeyboardInterrupt: + print("\n正在停止所有服务...") + for name, proc in processes: + print(f"停止 {name}...") + proc.terminate() + proc.wait() + print("所有服务已停止") + + except Exception as e: + print(f"启动服务时出错: {e}") + for name, proc in processes: + proc.terminate() + sys.exit(1) + +if __name__ == '__main__': + start_services() \ No newline at end of file diff --git a/tasks.py b/tasks.py new file mode 100644 index 0000000..3ed9fb5 --- /dev/null +++ b/tasks.py @@ -0,0 +1,136 @@ +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 new file mode 100644 index 0000000..5c86815 --- /dev/null +++ b/test_api.py @@ -0,0 +1,43 @@ +#!/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 new file mode 100644 index 0000000..5e4ebdd --- /dev/null +++ b/test_api_list.py @@ -0,0 +1,37 @@ +#!/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_beat_schedule.py b/test_beat_schedule.py new file mode 100644 index 0000000..afa204d --- /dev/null +++ b/test_beat_schedule.py @@ -0,0 +1,53 @@ +#!/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 new file mode 100644 index 0000000..27d6d58 --- /dev/null +++ b/test_body_config.py @@ -0,0 +1,84 @@ +#!/usr/bin/env python3 +"""测试body配置功能""" +from tasks import daily_morning_task +from generic_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 new file mode 100644 index 0000000..5a98405 --- /dev/null +++ b/test_call_api.py @@ -0,0 +1,47 @@ +#!/usr/bin/env python3 +"""测试呼叫任务API功能""" + +import sys +import os +sys.path.append(os.path.dirname(os.path.abspath(__file__))) + +from generic_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_merged_status.py b/test_merged_status.py new file mode 100644 index 0000000..71a816b --- /dev/null +++ b/test_merged_status.py @@ -0,0 +1,77 @@ +#!/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 new file mode 100644 index 0000000..fa85cff --- /dev/null +++ b/test_methods.py @@ -0,0 +1,53 @@ +#!/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_retry.py b/test_retry.py new file mode 100644 index 0000000..4c425c1 --- /dev/null +++ b/test_retry.py @@ -0,0 +1,38 @@ +#!/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 new file mode 100644 index 0000000..93bee4b --- /dev/null +++ b/test_schedule.py @@ -0,0 +1,39 @@ +#!/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_task.py b/test_task.py new file mode 100644 index 0000000..89f3e75 --- /dev/null +++ b/test_task.py @@ -0,0 +1,10 @@ +#!/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