From 5cf3b008d3c2ce557b3992838197cb33861088a3 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Mon, 8 Dec 2025 15:16:17 +0800 Subject: [PATCH] =?UTF-8?q?celery=20worker=E5=9C=A8windows=E4=B8=8A?= =?UTF-8?q?=E5=90=AF=E5=8A=A8=E5=A4=B1=E8=B4=A5=EF=BC=8C=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E5=8F=82=E6=95=B0--pool=3Dsolo?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/celery_app.py | 2 +- app/celery_tasks.py | 2 +- main.py | 97 ++------------------------------------------- 3 files changed, 6 insertions(+), 95 deletions(-) diff --git a/app/celery_app.py b/app/celery_app.py index cd98ddd..ebacaf8 100644 --- a/app/celery_app.py +++ b/app/celery_app.py @@ -31,7 +31,7 @@ celery_app.conf.update( beat_schedule={ 'push-data-to-dtc-every-minute': { 'task': 'push_data_to_dtc', - 'schedule': 120.0, # 每60秒执行一次(1分钟) + 'schedule': 60.0, # 每60秒执行一次(1分钟) }, }, ) diff --git a/app/celery_tasks.py b/app/celery_tasks.py index d087c01..b45d14c 100644 --- a/app/celery_tasks.py +++ b/app/celery_tasks.py @@ -33,7 +33,7 @@ def push_data_to_dtc_task(self): """ task_name = 'push_data_to_dtc' logger.info(f"🌿 开始推送数据给DTC任务") - + # 获取分布式锁,使用任务名称作为锁标识 lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时 diff --git a/main.py b/main.py index 61b3a9d..32b94bb 100644 --- a/main.py +++ b/main.py @@ -32,9 +32,8 @@ def start_celery_worker(): "-A", "app.celery_app", # 指定celery应用模块 "worker", '--loglevel=info', - '--concurrency=4', - '--prefetch-multiplier=1', - '--max-tasks-per-child=1000', + '--pool=solo', + '--concurrency=1', '--time-limit=300', # 5分钟任务超时 '--soft-time-limit=240', # 4分钟软超时 ], check=True) @@ -91,79 +90,6 @@ 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): # 启动时初始化 @@ -291,12 +217,10 @@ 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() # 如果没有传递任何参数,输出完整的提示信息 - if not args.mode and not args.stop: + if not args.mode: print("🚀 AI Talk Callback API 管理指南") print("=" * 50) print("\n📋 启动模式:") @@ -304,20 +228,11 @@ if __name__ == "__main__": print(" worker - 启动 Celery Worker 任务处理器") print(" beat - 启动 Celery Beat 定时任务调度器") print(" flower - 启动 Flower 监控服务") - 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") @@ -328,11 +243,7 @@ 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":