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

asyncio的Event/Condition搭配aiohttp使用不生效如何解决?

核心问题排查

你的代码中Event不生效主要有四个原因:

  • 异常捕获范围过窄:连接断开时不止会抛出ConnectionResetError,还可能抛出BrokenPipeError、ConnectionAbortedError等多种异常,你只捕获了前者,绝大多数断开场景都不会触发EVENT.set()
  • Event未做重置处理:首次触发EVENT.set()后,Event会一直处于已触发状态,后续新的客户端连接进来后,handler里的await EVENT.wait()会直接返回,导致新连接立刻断开
  • 写入错误触发延迟:TCP协议的特性决定了单次写入如果内核缓冲区还能容纳数据,即使客户端已经断开,write和drain也不会立刻抛出异常,需要等缓冲区写满后才会报错,导致感知断开的时机非常晚
  • 缺失主动断开检测:aiohttp中客户端主动断开连接时,handler协程会被asyncio注入CancelledError,你没有利用这个机制同步连接状态,完全依赖写入报错的被动检测,效率极低
修复后示例代码
import asyncio
import logging
from aiohttp import web, StreamResponse, Request

routes = web.RouteTableDef()
log = logging.getLogger(__name__)
CLIENT = None
EVENT = None

@routes.get('/telemetry/json')
async def handler(request: Request):
    global CLIENT, EVENT
    # 新连接接入先重置Event
    if EVENT.is_set():
        EVENT.clear()
    resp = StreamResponse()
    resp.headers['Content-Type'] = 'application/json'
    CLIENT = await resp.prepare(request)
    log.debug(f"Client {request.remote} connected, wait for disconnect event")
    try:
        # 额外加读操作,客户端断开时会立刻返回EOF或者抛异常
        await request.content.read()
    finally:
        # 不管是读报错还是协程被取消,都触发Event
        EVENT.set()
        log.debug(f"Client {request.remote} disconnected")
    return resp

async def main():
    global EVENT
    EVENT = asyncio.Event()
    app = web.Application()
    app.add_routes(routes)
    runner = web.AppRunner(app)
    await runner.setup()
    await web.TCPSite(runner, port=8080).start()
    while True:
        await asyncio.sleep(1)
        if CLIENT is None:
            continue
        try:
            await CLIENT.write('FLUSH\n'.encode('utf-8'))
            await CLIENT.drain()
        # 捕获所有连接相关异常
        except (ConnectionResetError, BrokenPipeError, ConnectionAbortedError):
            log.debug("Connection broken, notify event")
            EVENT.set()
            # 重置CLIENT状态,避免后续循环写入
            CLIENT = None

log.addHandler(logging.StreamHandler())
log.setLevel(logging.DEBUG)
asyncio.run(main())
额外优化建议
  • 实际多客户端场景不要用全局变量存储连接,可在app实例上挂载一个线程安全的连接集合存储所有在线客户端
  • 推送数据时遍历连接集合逐个写入,单个连接报错时只移除对应连接,不影响其他客户端
  • 不需要额外用Event阻塞handler,只要保持handler协程不退出即可,读操作本身就会一直阻塞直到客户端断开

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 23:06:04