From 25a016977c3eaf5a39d068121bac2f83963a573a Mon Sep 17 00:00:00 2001 From: liangtianyu <124244236@qq.com> Date: Sun, 14 Dec 2025 20:20:30 +0800 Subject: [PATCH] =?UTF-8?q?=E4=B8=8D=E5=A4=84=E7=90=86=E9=80=80=E5=87=BA?= =?UTF-8?q?=E4=BF=A1=E5=8F=B7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- main.py | 120 +++++++++++++++++++++++++++++++++----------------------- 1 file changed, 71 insertions(+), 49 deletions(-) diff --git a/main.py b/main.py index 581d875..7edd994 100644 --- a/main.py +++ b/main.py @@ -1,4 +1,3 @@ - from datetime import datetime import os import traceback @@ -22,6 +21,7 @@ from sqlalchemy import text LoggerManager.setup_logging() logger = get_main_logger() + def global_exception_handler(exc_type, exc_value, exc_traceback): if issubclass(exc_type, KeyboardInterrupt): sys.__excepthook__(exc_type, exc_value, exc_traceback) @@ -32,7 +32,7 @@ def global_exception_handler(exc_type, exc_value, exc_traceback): error_msg += f"异常类型: {exc_type.__name__}\n" error_msg += f"异常信息: {exc_value}\n" error_msg += "堆栈跟踪:\n" - error_msg += ''.join(traceback.format_tb(exc_traceback)) + error_msg += "".join(traceback.format_tb(exc_traceback)) error_msg += f"异常时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}\n" error_msg += "-" * 50 + "\n" @@ -42,14 +42,16 @@ def global_exception_handler(exc_type, exc_value, exc_traceback): # 输出到文件 logger.error(error_msg) + # 注册全局异常处理器 sys.excepthook = global_exception_handler + @asynccontextmanager async def lifespan(app: FastAPI): # 启动时初始化 logger.info("🚀 应用启动中...") - + # Redis 连接对象 redis_client = None @@ -76,10 +78,10 @@ async def lifespan(app: FastAPI): # 执行 ping 命令验证连接 await redis_client.ping() logger.info("✅ Redis 连接验证成功") - + # 存储到应用状态中供其他组件使用 app.state.redis_client = redis_client - + except Exception as redis_error: logger.error(f"❌ Redis 连接验证失败: {redis_error}") raise @@ -96,7 +98,7 @@ async def lifespan(app: FastAPI): finally: # 关闭时清理 logger.info("🛑 应用关闭中...") - + # 关闭 Redis 连接 if redis_client: try: @@ -104,7 +106,7 @@ async def lifespan(app: FastAPI): logger.info("🔴 Redis 连接已关闭") except Exception as e: logger.warning(f"⚠️ 关闭 Redis 连接时出现警告: {e}") - + logger.info("👋 应用已关闭") @@ -161,17 +163,18 @@ async def health_check(): raise HTTPException(status_code=404, detail="Not Found") logger.debug("💓 健康检查接口被访问") - + return {"status": "healthy"} # 全局变量存储进程 processes = [] + def signal_handler(signum, frame): """信号处理器,用于优雅关闭所有服务""" print(f"\n🛑 接收到信号 {signum},正在关闭所有服务...") - + # 逆序关闭进程(最后启动的最先关闭) for i, process in enumerate(reversed(processes)): if process and process.poll() is None: # 进程仍在运行 @@ -185,86 +188,104 @@ def signal_handler(signum, frame): process.kill() # 强制杀死进程 except Exception as e: print(f"❌ 关闭进程时出错: {e}") - + print("👋 所有服务已关闭") sys.exit(0) + def start_all_services(): """启动所有服务""" global processes - + print("\n🚀 AI Talk Callback API 一键启动所有服务") print("=" * 60) - + # 注册信号处理器 - signal.signal(signal.SIGINT, signal_handler) # Ctrl+C - signal.signal(signal.SIGTERM, signal_handler) # 终止信号 - + # signal.signal(signal.SIGINT, signal_handler) # Ctrl+C + # signal.signal(signal.SIGTERM, signal_handler) # 终止信号 + try: # 1. 启动 Celery Worker print("🌿 启动 Celery Worker...") - worker_process = subprocess.Popen([ - sys.executable, "-m", "celery", - "-A", "app.celery_app", - "worker", - '--loglevel=info', - '--pool=solo', - '--concurrency=1', - '--time-limit=300', # 5分钟任务超时 - '--soft-time-limit=240' # 4分钟软超时 - ]) + worker_process = subprocess.Popen( + [ + sys.executable, + "-m", + "celery", + "-A", + "app.celery_app", + "worker", + "--loglevel=info", + "--pool=solo", + "--concurrency=1", + "--time-limit=300", # 5分钟任务超时 + "--soft-time-limit=240", # 4分钟软超时 + ] + ) processes.append(worker_process) time.sleep(2) # 等待 Worker 启动 - + # 2. 启动 Celery Beat print("\n📅 启动 Celery Beat...") - beat_process = subprocess.Popen([ - sys.executable, "-m", "celery", - "-A", "app.celery_app", - "beat", - '--loglevel=info', - f'--schedule={os.path.join(tempfile.gettempdir(), "celerybeat-schedule")}' - ]) + beat_process = subprocess.Popen( + [ + sys.executable, + "-m", + "celery", + "-A", + "app.celery_app", + "beat", + "--loglevel=info", + f'--schedule={os.path.join(tempfile.gettempdir(), "celerybeat-schedule")}', + ] + ) processes.append(beat_process) time.sleep(2) # 等待 Beat 启动 - + # 3. 启动 Flower 监控(如果启用) if settings.flower_enabled: print("\n📊 启动 Flower 监控服务...") flower_cmd = [ - sys.executable, "-m", "celery", - "-A", "app.celery_app", + sys.executable, + "-m", + "celery", + "-A", + "app.celery_app", f"--broker={settings.celery_broker_url}", "flower", - f"--port={settings.flower_port}" + f"--port={settings.flower_port}", ] - + if settings.flower_basic_auth: flower_cmd.append(f"--basic_auth={settings.flower_basic_auth}") if settings.flower_url_prefix: flower_cmd.append(f"--url_prefix={settings.flower_url_prefix}") - + flower_process = subprocess.Popen(flower_cmd) processes.append(flower_process) time.sleep(2) # 等待 Flower 启动 - + # 4. 启动 FastAPI 应用 print("\n🚀 启动 FastAPI 应用...") api_cmd = [ - sys.executable, "-m", "uvicorn", + sys.executable, + "-m", + "uvicorn", "main:app", - "--host", "0.0.0.0", - "--port", "8000" + "--host", + "0.0.0.0", + "--port", + "8000", ] - + # 添加调试模式(如果配置了) if settings.debug: api_cmd.append("--reload") - + api_process = subprocess.Popen(api_cmd) processes.append(api_process) time.sleep(2) # 等待 API 启动 - + print("\n" + "=" * 60) print("✅ 所有服务启动完成!") print("\n🌐 服务地址:") @@ -275,7 +296,7 @@ def start_all_services(): print(f" 📊 监控界面: {settings.flower_url}") print("\n💡 使用 Ctrl+C 可以优雅关闭所有服务") print("=" * 60) - + # 等待所有进程 while True: # 检查是否有进程异常退出 @@ -284,15 +305,16 @@ def start_all_services(): print(f"❌ 进程 {i+1} 异常退出,退出码: {process.returncode}") signal_handler(signal.SIGINT, None) return - + time.sleep(1) # 每秒检查一次 - + except KeyboardInterrupt: signal_handler(signal.SIGINT, None) except Exception as e: print(f"❌ 启动服务时出错: {e}") signal_handler(signal.SIGINT, None) + if __name__ == "__main__": # 直接启动所有服务 start_all_services()