You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 13:20:39