增加停止服务开关

This commit is contained in:
mark.tian
2025-12-08 13:03:02 +08:00
parent 4b0ab3260a
commit a55c321fce
2 changed files with 113 additions and 16 deletions

View File

@@ -1,7 +1,7 @@
""" """
Celery任务定义 Celery任务定义
""" """
from celery import current_task
from sqlalchemy import create_engine from sqlalchemy import create_engine
import json import json
from app.celery_app import celery_app from app.celery_app import celery_app
@@ -26,7 +26,7 @@ engine = create_engine(
) )
@celery_app.task(bind=True, name='push_data_to_dtc') @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任务 推送数据给DTC的Celery任务
自动获取一条未完成的回调请求进行处理 自动获取一条未完成的回调请求进行处理
@@ -48,7 +48,7 @@ def push_data_to_dtc_task():
try: try:
with engine.connect() as conn: with engine.connect() as conn:
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', state='PROGRESS',
meta={'current': 0, 'total': 100, 'status': f'正在获取下一条回调请求日志...'} 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}") logger.info(f"📋 获取到未完成的回调请求: ID={callback_log_id}, site_id={site_id}")
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', state='PROGRESS',
meta={'current': 25, 'total': 100, 'status': f'获取回调请求成功(ID={callback_log_id}, site_id={site_id}),开始分析处理...'} 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}失败,任务停止执行") raise Exception(f"保存callback_data:{callback_log_id}失败,任务停止执行")
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', state='PROGRESS',
meta={'current': 50, 'total': 100, 'status': f'请求体中通话明细处理成功(ID={callback_log_id}, site_id={site_id}),开始过滤需要转发的通话记录...'} 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": "没有有效的通过记录需要处理"} return {"status": "completed", "message": "没有有效的通过记录需要处理"}
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', state='PROGRESS',
meta={'current': 70, 'total': 100, 'status': f'需要推送的通过记录已获取成功(ID={callback_log_id}, site_id={site_id}),准备转发...'} 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 filtered_request_body['data'] = related_records
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', 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( 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) mark_callback_log_completed(conn, callback_log_id)
# 更新任务状态 # 更新任务状态
current_task.update_state( self.update_state(
state='PROGRESS', 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 { return {

109
main.py
View File

@@ -5,6 +5,7 @@ from sqlalchemy import text
import redis.asyncio as redis import redis.asyncio as redis
import subprocess import subprocess
import sys import sys
import os
from app.config import settings from app.config import settings
from app.database import engine from app.database import engine
@@ -46,13 +47,18 @@ def start_celery_beat():
try: try:
logger.info("📅 启动Celery Beat调度器...") logger.info("📅 启动Celery Beat调度器...")
# 使用跨平台的调度文件路径
import tempfile
import os
schedule_file = os.path.join(tempfile.gettempdir(), 'celerybeat-schedule')
# 使用subprocess启动独立的celery beat进程 # 使用subprocess启动独立的celery beat进程
subprocess.run([ subprocess.run([
sys.executable, "-m", "celery", sys.executable, "-m", "celery",
"-A", "app.celery_app", # 指定celery应用模块 "-A", "app.celery_app", # 指定celery应用模块
"beat", "beat",
'--loglevel=info', '--loglevel=info',
'--schedule=/tmp/celerybeat-schedule', f'--schedule={schedule_file}',
], check=True) ], check=True)
except Exception as e: except Exception as e:
logger.error(f"❌ Celery Beat 启动失败: {e}") logger.error(f"❌ Celery Beat 启动失败: {e}")
@@ -85,6 +91,79 @@ def start_flower():
logger.error(f"❌ Flower 监控服务启动失败: {e}") 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 @asynccontextmanager
async def lifespan(app: FastAPI): async def lifespan(app: FastAPI):
# 启动时初始化 # 启动时初始化
@@ -212,22 +291,33 @@ if __name__ == "__main__":
parser = argparse.ArgumentParser(description="AI Talk Callback API") parser = argparse.ArgumentParser(description="AI Talk Callback API")
parser.add_argument("--mode", choices=["api", "worker", "beat", "flower"], parser.add_argument("--mode", choices=["api", "worker", "beat", "flower"],
help="启动模式: api(仅API), worker(仅Celery Worker), beat(仅Celery Beat), flower(仅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() args = parser.parse_args()
# 如果没有传递 mode 参数,输出完整的提示信息 # 如果没有传递任何参数,输出完整的提示信息
if not args.mode: if not args.mode and not args.stop:
print("🚀 AI Talk Callback API 启动指南") print("🚀 AI Talk Callback API 管理指南")
print("=" * 50) print("=" * 50)
print("\n📋 可用的启动模式:") print("\n📋 启动模式:")
print(" api - 启动 FastAPI Web 应用服务 (端口: 8000)") print(" api - 启动 FastAPI Web 应用服务 (端口: 8000)")
print(" worker - 启动 Celery Worker 任务处理器") print(" worker - 启动 Celery Worker 任务处理器")
print(" beat - 启动 Celery Beat 定时任务调度器") print(" beat - 启动 Celery Beat 定时任务调度器")
print(" flower - 启动 Flower 监控服务") 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=api # 启动 Web API 服务")
print(" python main.py --mode=worker # 启动任务处理器") print(" python main.py --mode=worker # 启动任务处理器")
print(" python main.py --mode=beat # 启动定时任务调度器") print(" python main.py --mode=beat # 启动定时任务调度器")
print(f" python main.py --mode=flower # 启动监控服务 (访问: {settings.flower_url})") 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("\n🌐 服务地址:")
print(" API 服务: http://localhost:8000") print(" API 服务: http://localhost:8000")
print(" API 文档: http://localhost:8000/docs") print(" API 文档: http://localhost:8000/docs")
@@ -238,6 +328,13 @@ if __name__ == "__main__":
print(" - 建议在多个终端中分别启动不同服务") print(" - 建议在多个终端中分别启动不同服务")
sys.exit(0) sys.exit(0)
# 处理停止服务请求
if args.stop:
logger.info(f"🛑 正在停止 {args.stop} 服务...")
stop_celery_service(args.stop)
sys.exit(0)
# 处理启动服务请求
if args.mode == "api": if args.mode == "api":
# 仅启动 FastAPI 应用 # 仅启动 FastAPI 应用
logger.info("🚀 启动FastAPI应用...") logger.info("🚀 启动FastAPI应用...")