celery worker在windows上启动失败,增加参数--pool=solo
This commit is contained in:
@@ -31,7 +31,7 @@ celery_app.conf.update(
|
|||||||
beat_schedule={
|
beat_schedule={
|
||||||
'push-data-to-dtc-every-minute': {
|
'push-data-to-dtc-every-minute': {
|
||||||
'task': 'push_data_to_dtc',
|
'task': 'push_data_to_dtc',
|
||||||
'schedule': 120.0, # 每60秒执行一次(1分钟)
|
'schedule': 60.0, # 每60秒执行一次(1分钟)
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -33,7 +33,7 @@ def push_data_to_dtc_task(self):
|
|||||||
"""
|
"""
|
||||||
task_name = 'push_data_to_dtc'
|
task_name = 'push_data_to_dtc'
|
||||||
logger.info(f"🌿 开始推送数据给DTC任务")
|
logger.info(f"🌿 开始推送数据给DTC任务")
|
||||||
|
|
||||||
# 获取分布式锁,使用任务名称作为锁标识
|
# 获取分布式锁,使用任务名称作为锁标识
|
||||||
lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时
|
lock = redis_manager.create_lock(f"celery_task:{task_name}", timeout=300) # 5分钟超时
|
||||||
|
|
||||||
|
|||||||
97
main.py
97
main.py
@@ -32,9 +32,8 @@ def start_celery_worker():
|
|||||||
"-A", "app.celery_app", # 指定celery应用模块
|
"-A", "app.celery_app", # 指定celery应用模块
|
||||||
"worker",
|
"worker",
|
||||||
'--loglevel=info',
|
'--loglevel=info',
|
||||||
'--concurrency=4',
|
'--pool=solo',
|
||||||
'--prefetch-multiplier=1',
|
'--concurrency=1',
|
||||||
'--max-tasks-per-child=1000',
|
|
||||||
'--time-limit=300', # 5分钟任务超时
|
'--time-limit=300', # 5分钟任务超时
|
||||||
'--soft-time-limit=240', # 4分钟软超时
|
'--soft-time-limit=240', # 4分钟软超时
|
||||||
], check=True)
|
], check=True)
|
||||||
@@ -91,79 +90,6 @@ 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):
|
||||||
# 启动时初始化
|
# 启动时初始化
|
||||||
@@ -291,12 +217,10 @@ 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()
|
||||||
|
|
||||||
# 如果没有传递任何参数,输出完整的提示信息
|
# 如果没有传递任何参数,输出完整的提示信息
|
||||||
if not args.mode and not args.stop:
|
if not args.mode:
|
||||||
print("🚀 AI Talk Callback API 管理指南")
|
print("🚀 AI Talk Callback API 管理指南")
|
||||||
print("=" * 50)
|
print("=" * 50)
|
||||||
print("\n📋 启动模式:")
|
print("\n📋 启动模式:")
|
||||||
@@ -304,20 +228,11 @@ if __name__ == "__main__":
|
|||||||
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(" worker - 停止 Celery Worker 服务")
|
|
||||||
print(" beat - 停止 Celery Beat 调度器")
|
|
||||||
print(" flower - 停止 Flower 监控服务")
|
|
||||||
print(" all - 停止所有 Celery 相关服务")
|
|
||||||
print("\n🔧 使用示例:")
|
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")
|
||||||
@@ -328,11 +243,7 @@ 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":
|
||||||
|
|||||||
Reference in New Issue
Block a user