From a55c321fcea0c887ca2bddcfc20543faeedb5889 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Mon, 8 Dec 2025 13:03:02 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E5=81=9C=E6=AD=A2=E6=9C=8D?= =?UTF-8?q?=E5=8A=A1=E5=BC=80=E5=85=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/celery_tasks.py | 20 ++++---- main.py | 109 +++++++++++++++++++++++++++++++++++++++++--- 2 files changed, 113 insertions(+), 16 deletions(-) diff --git a/app/celery_tasks.py b/app/celery_tasks.py index c7b2e4b..1830bfc 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -1,7 +1,7 @@ """ Celery任务定义 """ -from celery import current_task + from sqlalchemy import create_engine import json from app.celery_app import celery_app @@ -26,7 +26,7 @@ engine = create_engine( ) @celery_app.task(bind=True, name='push_data_to_dtc') -def push_data_to_dtc_task(): +def push_data_to_dtc_task(self): """ 推送数据给DTC的Celery任务 自动获取一条未完成的回调请求进行处理 @@ -48,7 +48,7 @@ def push_data_to_dtc_task(): try: with engine.connect() as conn: # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'} ) @@ -70,7 +70,7 @@ def push_data_to_dtc_task(): logger.info(f"📋 获取到未完成的回调请求: ID={callback_log_id}, site_id={site_id}") # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'} ) @@ -88,7 +88,7 @@ def push_data_to_dtc_task(): raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行") # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', meta={'current': 50, 'total': 100, 'status': f'请求体中通话明细处理成功(ID={callback_log_id}, site_id={site_id}),开始过滤需要转发的通话记录...'} ) @@ -140,7 +140,7 @@ def push_data_to_dtc_task(): return {"status": "completed", "message": "没有有效的通过记录需要处理"} # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', meta={'current': 70, 'total': 100, 'status': f'需要推送的通过记录已获取成功(ID={callback_log_id}, site_id={site_id}),准备转发...'} ) @@ -151,9 +151,9 @@ def push_data_to_dtc_task(): filtered_request_body['data'] = related_records # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', - meta={'current': 85, 'total': 100, 'status': '开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'} + meta={'current': 85, 'total': 100, 'status': f'开始推送数据给DTC(ID={callback_log_id}, site_id={site_id})...'} ) success, retry_count = call_external_api_with_retry( @@ -177,9 +177,9 @@ def push_data_to_dtc_task(): mark_callback_log_completed(conn, callback_log_id) # 更新任务状态 - current_task.update_state( + self.update_state( state='PROGRESS', - meta={'current': 100, 'total': 100, 'status': '标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'} + meta={'current': 100, 'total': 100, 'status': f'标记回调请求日志为已完成(ID={callback_log_id}, site_id={site_id})'} ) return { diff --git a/main.py b/main.py index 9ea5ffb..6b0def6 100644 --- a/main.py +++ b/main.py @@ -5,6 +5,7 @@ from sqlalchemy import text import redis.asyncio as redis import subprocess import sys +import os from app.config import settings from app.database import engine @@ -46,13 +47,18 @@ def start_celery_beat(): try: logger.info("📅 启动Celery Beat调度器...") + # 使用跨平台的调度文件路径 + import tempfile + import os + schedule_file = os.path.join(tempfile.gettempdir(), 'celerybeat-schedule') + # 使用subprocess启动独立的celery beat进程 subprocess.run([ sys.executable, "-m", "celery", "-A", "app.celery_app", # 指定celery应用模块 "beat", '--loglevel=info', - '--schedule=/tmp/celerybeat-schedule', + f'--schedule={schedule_file}', ], check=True) except Exception as e: logger.error(f"❌ Celery Beat 启动失败: {e}") @@ -85,6 +91,79 @@ def start_flower(): logger.error(f"❌ Flower 监控服务启动失败: {e}") +def stop_celery_service(service_name): + """停止 Celery 服务""" + try: + logger.info(f"🛑 正在停止 {service_name} 服务...") + + # 查找并停止相关进程 + if service_name == "worker": + # 停止 Celery Worker + subprocess.run([ + sys.executable, "-m", "celery", + "-A", "app.celery_app", + "control", + "shutdown" + ], check=False) + logger.info("✅ Celery Worker 停止命令已发送") + + elif service_name == "beat": + # 查找并停止 beat 进程 + try: + if os.name == 'nt': # Windows + # 使用 PowerShell 查找并停止相关进程 + subprocess.run([ + "powershell", "-Command", + "Get-Process python | Where-Object {$_.ProcessName -eq 'python' -and $_.MainWindowTitle -like '*beat*'} | Stop-Process -Force" + ], check=False) + # 备用方案:查找包含 celery beat 的进程 + subprocess.run([ + "powershell", "-Command", + "Get-WmiObject Win32_Process | Where-Object {$_.Name -eq 'python.exe' -and $_.CommandLine -like '*beat*'} | ForEach-Object {Stop-Process -Id $_.ProcessId -Force}" + ], check=False) + else: # Linux/Mac + subprocess.run([ + "pkill", "-f", "celery.*beat" + ], check=False) + + logger.info(f"✅ Celery Beat 停止命令已发送") + except Exception as e: + logger.warning(f"⚠️ 停止 {service_name} 时出现警告: {e}") + + elif service_name == "flower": + # 查找并停止 flower 进程 + try: + if os.name == 'nt': # Windows + # 使用 PowerShell 查找并停止相关进程 + subprocess.run([ + "powershell", "-Command", + "Get-Process python | Where-Object {$_.ProcessName -eq 'python' -and $_.MainWindowTitle -like '*flower*'} | Stop-Process -Force" + ], check=False) + # 备用方案:查找包含 celery flower 的进程 + subprocess.run([ + "powershell", "-Command", + "Get-WmiObject Win32_Process | Where-Object {$_.Name -eq 'python.exe' -and $_.CommandLine -like '*flower*'} | ForEach-Object {Stop-Process -Id $_.ProcessId -Force}" + ], check=False) + else: # Linux/Mac + subprocess.run([ + "pkill", "-f", "celery.*flower" + ], check=False) + + logger.info(f"✅ Flower 监控服务停止命令已发送") + except Exception as e: + logger.warning(f"⚠️ 停止 {service_name} 时出现警告: {e}") + + elif service_name == "all": + # 停止所有 Celery 相关服务 + logger.info("🛑 正在停止所有 Celery 服务...") + stop_celery_service("worker") + stop_celery_service("beat") + stop_celery_service("flower") + + except Exception as e: + logger.error(f"❌ 停止 {service_name} 服务失败: {e}") + + @asynccontextmanager async def lifespan(app: FastAPI): # 启动时初始化 @@ -212,22 +291,33 @@ if __name__ == "__main__": parser = argparse.ArgumentParser(description="AI Talk Callback API") parser.add_argument("--mode", choices=["api", "worker", "beat", "flower"], help="启动模式: api(仅API), worker(仅Celery Worker), beat(仅Celery Beat), flower(仅Flower监控)") + parser.add_argument("--stop", choices=["worker", "beat", "flower", "all"], + help="停止服务: worker(停止Worker), beat(停止Beat), flower(停止Flower), all(停止所有Celery服务)") args = parser.parse_args() - # 如果没有传递 mode 参数,输出完整的提示信息 - if not args.mode: - print("🚀 AI Talk Callback API 启动指南") + # 如果没有传递任何参数,输出完整的提示信息 + if not args.mode and not args.stop: + print("🚀 AI Talk Callback API 管理指南") print("=" * 50) - print("\n📋 可用的启动模式:") + print("\n📋 启动模式:") print(" api - 启动 FastAPI Web 应用服务 (端口: 8000)") print(" worker - 启动 Celery Worker 任务处理器") print(" beat - 启动 Celery Beat 定时任务调度器") print(" flower - 启动 Flower 监控服务") - print("\n🔧 启动示例:") + print("\n🛑 停止服务:") + print(" worker - 停止 Celery Worker 服务") + print(" beat - 停止 Celery Beat 调度器") + print(" flower - 停止 Flower 监控服务") + print(" all - 停止所有 Celery 相关服务") + print("\n🔧 使用示例:") print(" python main.py --mode=api # 启动 Web API 服务") print(" python main.py --mode=worker # 启动任务处理器") print(" python main.py --mode=beat # 启动定时任务调度器") print(f" python main.py --mode=flower # 启动监控服务 (访问: {settings.flower_url})") + print("\n python main.py --stop=worker # 停止 Worker 服务") + print(" python main.py --stop=beat # 停止 Beat 调度器") + print(" python main.py --stop=flower # 停止 Flower 监控") + print(" python main.py --stop=all # 停止所有 Celery 服务") print("\n🌐 服务地址:") print(" API 服务: http://localhost:8000") print(" API 文档: http://localhost:8000/docs") @@ -238,6 +328,13 @@ if __name__ == "__main__": print(" - 建议在多个终端中分别启动不同服务") sys.exit(0) + # 处理停止服务请求 + if args.stop: + logger.info(f"🛑 正在停止 {args.stop} 服务...") + stop_celery_service(args.stop) + sys.exit(0) + + # 处理启动服务请求 if args.mode == "api": # 仅启动 FastAPI 应用 logger.info("🚀 启动FastAPI应用...")