From 4b0ab3260aaf7d8ba6cb3acd618bdd3ca85ad627 Mon Sep 17 00:00:00 2001 From: "mark.tian" Date: Mon, 8 Dec 2025 12:31:05 +0800 Subject: [PATCH] =?UTF-8?q?=E6=A2=B3=E7=90=86main.py=E4=B8=AD=E7=9A=84?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 7 +++ app/celery_app.py | 7 +++ app/config.py | 14 ++++-- celery_worker.py | 28 ----------- main.py | 117 ++++++++++++++++++++++++++++------------------ requirements.txt | 1 + 6 files changed, 97 insertions(+), 77 deletions(-) delete mode 100644 celery_worker.py diff --git a/.env.example b/.env.example index c1572a3..e72e6d6 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,13 @@ DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/ai_talk_callback_ CELERY_BROKER_URL=redis://localhost:6379/0 CELERY_RESULT_BACKEND=redis://localhost:6379/0 +# Flower监控配置 +FLOWER_ENABLED=true +FLOWER_PORT=5555 +FLOWER_BASIC_AUTH=admin:admin123 +FLOWER_URL_PREFIX= +FLOWER_URL=http://localhost:5555 + # Redis配置 (扩展配置,如果需要覆盖默认值) REDIS_PASSWORD= REDIS_MAX_CONNECTIONS=20 diff --git a/app/celery_app.py b/app/celery_app.py index 13bce2b..b705306 100644 --- a/app/celery_app.py +++ b/app/celery_app.py @@ -27,6 +27,13 @@ celery_app.conf.update( task_soft_time_limit=25 * 60, # 25分钟软超时 worker_prefetch_multiplier=1, worker_max_tasks_per_child=1000, + # Beat 调度配置 + beat_schedule={ + 'push-data-to-dtc-every-minute': { + 'task': 'push_data_to_dtc', + 'schedule': 120.0, # 每60秒执行一次(1分钟) + }, + }, ) logger.info("🌿 Celery应用配置完成") \ No newline at end of file diff --git a/app/config.py b/app/config.py index 81ffb62..d585a18 100644 --- a/app/config.py +++ b/app/config.py @@ -6,17 +6,25 @@ class Settings(BaseSettings): def redis_url(self) -> str: """从Celery配置获取Redis URL""" return self.celery_broker_url + # 数据库配置 database_url: str = "" # Celery配置 celery_broker_url: str = "redis://localhost:6379/0" - celery_result_backend: str = "redis://localhost:6379/0" + celery_result_backend: str = "redis://localhost:6379/1" celery_task_serializer: str = "json" celery_result_serializer: str = "json" celery_accept_content: list = ["json"] - celery_timezone: str = "UTC" - celery_enable_utc: bool = True + celery_timezone: str = "Asia/Shanghai" + celery_enable_utc: bool = False + + # Flower监控配置 + flower_enabled: bool = True # 是否启用Flower监控 + flower_url: str = "http://localhost:5555" # Flower访问URL + flower_port: int = 5555 # Flower服务端口 + flower_url_prefix: str = "" # Flower URL前缀 + flower_basic_auth: str = "admin:admin123" # Flower基础认证,格式:username:password # Redis配置 (从Celery配置获取) redis_password: str = "" diff --git a/celery_worker.py b/celery_worker.py deleted file mode 100644 index afc742a..0000000 --- a/celery_worker.py +++ /dev/null @@ -1,28 +0,0 @@ -#!/usr/bin/env python3 -""" -Celery Worker 启动脚本 -""" -import os -import sys - -# 添加项目根目录到Python路径 -sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) - -from app.celery_app import celery_app -from app.logger import get_logger - -logger = get_logger("celery_worker") - -if __name__ == "__main__": - logger.info("🌿 启动Celery Worker...") - - # 启动Celery worker - celery_app.start([ - 'worker', - '--loglevel=info', - '--concurrency=4', - '--prefetch-multiplier=1', - '--max-tasks-per-child=1000', - '--time-limit=300', # 5分钟任务超时 - '--soft-time-limit=240', # 4分钟软超时 - ]) \ No newline at end of file diff --git a/main.py b/main.py index cc8832a..9ea5ffb 100644 --- a/main.py +++ b/main.py @@ -3,16 +3,13 @@ from fastapi.middleware.cors import CORSMiddleware from contextlib import asynccontextmanager from sqlalchemy import text import redis.asyncio as redis -import threading import subprocess import sys -import os from app.config import settings from app.database import engine from app.routes import router from app.logger import get_logger, LoggerManager -from app.celery_app import celery_app # 初始化日志系统 @@ -25,20 +22,21 @@ celery_worker_process = None def start_celery_worker(): """启动 Celery Worker""" - global celery_worker_process try: logger.info("🌿 启动Celery Worker...") - # 启动Celery worker - celery_app.start([ - 'worker', + # 使用subprocess启动独立的celery worker进程 + subprocess.run([ + sys.executable, "-m", "celery", + "-A", "app.celery_app", # 指定celery应用模块 + "worker", '--loglevel=info', '--concurrency=4', '--prefetch-multiplier=1', '--max-tasks-per-child=1000', '--time-limit=300', # 5分钟任务超时 '--soft-time-limit=240', # 4分钟软超时 - ]) + ], check=True) except Exception as e: logger.error(f"❌ Celery Worker 启动失败: {e}") @@ -48,16 +46,45 @@ def start_celery_beat(): try: logger.info("📅 启动Celery Beat调度器...") - # 启动Celery beat - celery_app.start([ - 'beat', + # 使用subprocess启动独立的celery beat进程 + subprocess.run([ + sys.executable, "-m", "celery", + "-A", "app.celery_app", # 指定celery应用模块 + "beat", '--loglevel=info', '--schedule=/tmp/celerybeat-schedule', - ]) + ], check=True) except Exception as e: logger.error(f"❌ Celery Beat 启动失败: {e}") +def start_flower(): + """启动 Flower 监控服务""" + try: + logger.info("📊 启动Flower监控服务...") + + # 构建Flower启动命令 - 独立进程启动 + flower_cmd = [ + sys.executable, "-m", "celery", + "-A", "app.celery_app", # 指定celery应用模块 + f"--broker={settings.celery_broker_url}", + "flower", + f"--port={settings.flower_port}" + ] + + # 添加基础认证(如果配置了) + if settings.flower_basic_auth: + flower_cmd.append(f"--basic_auth={settings.flower_basic_auth}") + + # 添加URL前缀(如果配置了) + if settings.flower_url_prefix: + flower_cmd.append(f"--url_prefix={settings.flower_url_prefix}") + + subprocess.run(flower_cmd, check=True) + except Exception as e: + logger.error(f"❌ Flower 监控服务启动失败: {e}") + + @asynccontextmanager async def lifespan(app: FastAPI): # 启动时初始化 @@ -97,19 +124,7 @@ async def lifespan(app: FastAPI): logger.error(f"❌ Redis 连接验证失败: {redis_error}") raise - # 检查启动模式 - if len(sys.argv) > 1: - mode = sys.argv[1].replace("--mode=", "") - if mode == "all": - # 启动 Celery Worker (在后台线程中) - logger.info("🌿 启动Celery Worker...") - celery_worker_thread = threading.Thread(target=start_celery_worker, daemon=True) - celery_worker_thread.start() - - # 启动 Celery Beat 调度器 (在后台线程中) - logger.info("📅 启动Celery Beat调度器...") - celery_beat_thread = threading.Thread(target=start_celery_beat, daemon=True) - celery_beat_thread.start() + logger.info("✅ 服务初始化完成,FastAPI应用启动中...") logger.info(f"🎉 {settings.app_name} 启动完成!") yield @@ -187,21 +202,7 @@ async def health_check(): logger.debug("💓 健康检查接口被访问") - health_status = {"status": "healthy"} - - # 检查 Redis 连接状态 - try: - redis_client = getattr(app.state, 'redis_client', None) - if redis_client: - await redis_client.ping() - health_status["redis"] = "connected" - else: - health_status["redis"] = "disconnected" - except Exception as e: - health_status["redis"] = f"error: {str(e)}" - health_status["status"] = "degraded" - - return health_status + return {"status": "healthy"} if __name__ == "__main__": @@ -209,10 +210,34 @@ if __name__ == "__main__": import argparse parser = argparse.ArgumentParser(description="AI Talk Callback API") - parser.add_argument("--mode", choices=["api", "worker", "beat", "all"], - default="api", help="启动模式: api(仅API), worker(仅Celery Worker), beat(仅Celery Beat), all(全部)") + parser.add_argument("--mode", choices=["api", "worker", "beat", "flower"], + help="启动模式: api(仅API), worker(仅Celery Worker), beat(仅Celery Beat), flower(仅Flower监控)") args = parser.parse_args() + # 如果没有传递 mode 参数,输出完整的提示信息 + if not args.mode: + print("🚀 AI Talk Callback API 启动指南") + print("=" * 50) + print("\n📋 可用的启动模式:") + print(" api - 启动 FastAPI Web 应用服务 (端口: 8000)") + print(" worker - 启动 Celery Worker 任务处理器") + print(" beat - 启动 Celery Beat 定时任务调度器") + print(" flower - 启动 Flower 监控服务") + 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🌐 服务地址:") + print(" API 服务: http://localhost:8000") + print(" API 文档: http://localhost:8000/docs") + print(f" 任务监控界面: {settings.flower_url}") + print("\n💡 提示:") + print(" - 请确保 Redis 和 PostgreSQL 服务已启动") + print(" - 生产环境请根据需要调整配置文件") + print(" - 建议在多个终端中分别启动不同服务") + sys.exit(0) + if args.mode == "api": # 仅启动 FastAPI 应用 logger.info("🚀 启动FastAPI应用...") @@ -225,7 +250,7 @@ if __name__ == "__main__": # 仅启动 Celery Beat logger.info("📅 启动Celery Beat调度器...") start_celery_beat() - elif args.mode == "all": - # 启动 FastAPI + Celery Worker + Celery Beat - logger.info("🚀 启动完整服务栈(FastAPI + Celery Worker + Celery Beat)...") - uvicorn.run("main:app", host="0.0.0.0", port=8000, reload=settings.debug) + elif args.mode == "flower": + # 仅启动 Flower 监控服务 + logger.info("📊 启动Flower监控服务...") + start_flower() diff --git a/requirements.txt b/requirements.txt index 264e742..dc9248e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,6 +5,7 @@ asyncpg>=0.13.0 alembic>=1.17.2 celery>=5.3.0 redis>=4.5.0 +flower>=2.0.1 pydantic>=2.12.5 pydantic-settings>=2.12.0 python-multipart>=0.0.20