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

如何在子进程中终止asyncio?解决进程挂起问题

解决WebSocket/RabbitMQ进程停止后挂起的问题

当前使用websockets.py的asyncio信号处理、资源清理逻辑存在缺陷,导致接收停止信号后进程无法正常退出,即使改用Docker内的supervisord管理,进程仍会挂起。以下是针对性修复方案:

核心问题分析

  1. 同步阻塞操作:stop_server中使用time.sleep(5)会卡住asyncio事件循环,导致后续清理操作无法完成。
  2. 未初始化全局变量:server、rabbitmq_connection等全局变量未提前初始化,触发停止时会抛出NameError。
  3. 信号处理逻辑不完善:lambda创建的任务未被正确关联到事件循环的退出逻辑,main函数仅依赖server.wait_closed(),无法响应其他资源清理完成的信号。
  4. 缺失模块导入:未导入logging模块却调用logging.info。

修复后的完整代码

# websockets.py
import asyncio
import websockets
import json
import sys
import signal
import logging

# 提前初始化全局变量,避免未定义错误
server = None
rabbitmq_connection_reader = None
rabbitmq_connection = None

logging.basicConfig(level=logging.INFO)

async def setup_rabbitmq():
    # 补充你的RabbitMQ初始化逻辑
    global rabbitmq_connection, rabbitmq_connection_reader
    # 示例:
    # rabbitmq_connection = await your_rabbitmq_connect_function()
    # rabbitmq_connection_reader = await rabbitmq_connection.create_channel_reader()
    pass

async def handler(websocket):
    # 你的WebSocket消息处理逻辑
    pass

async def stop_server(exit_event):
    global server, rabbitmq_connection_reader, rabbitmq_connection
    logging.info("Initiating service shutdown...")
    
    # 按顺序清理资源:先关闭RabbitMQ,再关闭WebSocket
    if rabbitmq_connection_reader:
        await rabbitmq_connection_reader.close()
        logging.info("RabbitMQ reader closed")
    if rabbitmq_connection:
        await rabbitmq_connection.close()
        logging.info("RabbitMQ connection closed")
    
    if server:
        server.close()
        await server.wait_closed()
        logging.info("WebSocket server closed")
    
    # 使用asyncio.sleep替代同步sleep,避免阻塞事件循环
    await asyncio.sleep(5)
    
    # 设置退出事件,通知main函数结束
    exit_event.set()
    logging.info("All services stopped successfully")

async def main():
    global server
    exit_event = asyncio.Event()
    loop = asyncio.get_running_loop()
    
    # 定义信号处理器:创建stop_server任务
    def handle_signal():
        asyncio.create_task(stop_server(exit_event))
    
    # 注册SIGTERM和SIGINT信号处理(兼容supervisord默认停止信号)
    for signame in ('SIGTERM', 'SIGINT'):
        try:
            loop.add_signal_handler(getattr(signal, signame), handle_signal)
        except ValueError:
            logging.warning(f"Failed to register signal handler for {signame} (non-main thread environment)")
    
    try:
        await setup_rabbitmq()
        logging.info("RabbitMQ setup finished")
        
        # 启动WebSocket服务器
        server = await websockets.serve(handler, sys.argv[1], sys.argv[2])
        logging.info(f"WebSocket server running on {sys.argv[1]}:{sys.argv[2]}")
        
        # 同时等待服务器关闭或退出事件触发,确保进程能正常结束
        await asyncio.gather(
            server.wait_closed(),
            exit_event.wait()
        )
        
    except Exception as e:
        logging.error(f"Runtime error: {e}")
        exit_event.set()

if __name__ == "__main__":
    if len(sys.argv) < 3:
        logging.error(f"Usage: python {sys.argv[0]} <HOST> <PORT>")
        sys.exit(1)
    asyncio.run(main())

关键修复点说明

  • 全局变量初始化:提前声明并初始化全局资源变量,避免清理时抛出未定义错误。
  • 退出事件管理:用asyncio.Event统一管理进程退出触发条件,确保main函数能响应资源清理完成的信号。
  • 替换同步sleep:用await asyncio.sleep(5)替代time.sleep(5),不阻塞asyncio事件循环,保证所有异步清理操作能执行完毕。
  • 信号处理优化:将信号处理器逻辑抽离为独立函数,避免lambda的变量捕获问题,确保stop_server任务被正确创建。
  • 完善错误处理:异常场景下主动触发退出事件,避免进程因未处理异常挂起。

内容的提问来源于stack exchange,提问作者user655355

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 04:44:54