梳理main.py中的服务

This commit is contained in:
mark.tian
2025-12-08 12:31:05 +08:00
parent 43ff4aa5f9
commit 4b0ab3260a
6 changed files with 97 additions and 77 deletions

View File

@@ -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

View File

@@ -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应用配置完成")

View File

@@ -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 = ""

View File

@@ -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分钟软超时
])

117
main.py
View File

@@ -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()

View File

@@ -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