FastAPI启动时等待Java服务启动完成的异步实现方案咨询
问题描述
我需要启动FastAPI服务并与Java服务通过HTTP通信,计划用subprocess调用Java启动脚本,通过监控日志判断启停状态。因为subprocess是阻塞的,改用asyncio.subprocess_shell异步调用,但它不会等待Java启动成功,存在潜在风险。尝试用asyncio.SubprocessProtocol监控日志直到出现成功标识,但直接写while循环会阻塞事件循环,导致进程僵死。请问如何让startup_event在Java服务启动成功后再标记启动完成?
原代码:
import asyncio import re from logging import logger from fastapi import FastAPI app = FastAPI() transport = asyncio.SubprocessTransport() protocal = asyncio.SubprocessProtocol() class MyProtocol(asyncio.SubprocessProtocol): startup_str = re.compile("Server - Started") is_startup = False is_exited = False def pipe_data_received(self, fd, data): logger.info(data) super().pipe_data_received(fd, data) if not self.is_startup: if re.search(self.startup_str, str(data)): self.is_startup = True def pipe_connection_lost(self, fd, exc): if exc is None: logger.debug(f"Pipe {fd} Closed") else: logger.error(exc) super().pipe_connection_lost(fd, exc) def process_exited(self): logger.info("Process Exited") super().process_exited() self.is_exited = True @app.on_event("startup") async def startup_event(): loop = asyncio.get_running_loop() global transport, protocal transport, protocal = await loop.subprocess_shell(MyProtocol, "/start_java_server.sh") # Waiting Application startup complete. # while not self.protocal.is_startup: # pass @app.on_event("shutdown") async def shutdown_event(): global transport, protocal transport.close()
解决方案
核心问题是直接用while循环会阻塞asyncio事件循环,导致pipe_data_received回调无法执行,永远等不到启动成功的信号。改用asyncio.Event实现异步等待,它能安全地在回调和协程间传递信号。
修改后的代码:
import asyncio import re from logging import logger from fastapi import FastAPI app = FastAPI() transport = None protocol = None # 修正原拼写错误:protocal → protocol class MyProtocol(asyncio.SubprocessProtocol): startup_str = re.compile("Server - Started") is_exited = False def __init__(self): self.startup_event = asyncio.Event() # 初始化异步事件 def pipe_data_received(self, fd, data): log_data = data.decode('utf-8') # 显式解码bytes为字符串,避免乱码 logger.info(log_data) super().pipe_data_received(fd, data) if not self.startup_event.is_set(): if re.search(self.startup_str, log_data): logger.info("Java服务启动成功") self.startup_event.set() # 设置事件,通知等待的协程 def pipe_connection_lost(self, fd, exc): if exc is None: logger.debug(f"Pipe {fd} Closed") else: logger.error(f"Pipe {fd} 异常关闭: {exc}") super().pipe_connection_lost(fd, exc) def process_exited(self): logger.info("Java进程已退出") super().process_exited() self.is_exited = True # 如果进程退出但启动事件未触发,设置事件避免协程永久等待 if not self.startup_event.is_set(): self.startup_event.set() logger.error("Java进程启动失败,已退出") @app.on_event("startup") async def startup_event(): global transport, protocol loop = asyncio.get_running_loop() protocol = MyProtocol() # 先实例化自定义Protocol,拿到startup_event transport, _ = await loop.subprocess_shell( lambda: protocol, # 用lambda传递已实例化的Protocol "/start_java_server.sh" ) # 异步等待启动事件,不会阻塞事件循环 await protocol.startup_event.wait() # 检查是否是因为进程退出导致事件触发 if protocol.is_exited: raise RuntimeError("Java服务启动失败,进程已退出") logger.info("FastAPI启动前准备完成") @app.on_event("shutdown") async def shutdown_event(): global transport, protocol if transport and not transport.is_closing(): transport.close() logger.info("FastAPI已关闭,Java服务已终止")
关键修改点说明
- 用
asyncio.Event替代轮询:startup_event.wait()会释放事件循环,让pipe_data_received回调可以正常处理日志输出,当检测到成功标识时设置事件,协程会自动恢复执行。 - 修正Protocol实例化方式:原代码直接传递类给
subprocess_shell,现在先实例化自定义Protocol,通过lambda传递,这样可以提前拿到startup_event对象用于等待。 - 显式解码日志数据:避免直接把bytes转字符串导致的乱码问题,确保正则匹配正常工作。
- 异常处理:在
process_exited中检查启动事件是否已触发,若未触发则设置事件并抛出异常,避免FastAPI一直卡在启动阶段。 - 拼写修正:原代码中
protocal是拼写错误,改为正确的protocol。
内容的提问来源于stack exchange,提问作者chetir
相关产品推荐
相关产品推荐

