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

Tornado HTTPServer阻塞持续异步回调的问题求解

核心问题原因

  • 阻塞问题根源是sleep函数使用错误:如果代码中导入的是标准库time.sleep,该函数为同步阻塞调用,执行后会卡住整个Tornado IO 循环,导致协程调度中断,因此队列检查只会执行一次。需要替换为异步非阻塞的tornado.gen.sleep或asyncio.sleep,等待期间IO循环可正常处理WebSocket请求等其他任务。
  • 队列非空判断逻辑不规范:queue.not_empty是标准库同步队列的内部条件变量,不属于对外开放的判断API,正确的非空判断应使用not queue.empty()。同时建议在调用convert_oldest时增加异常捕获,避免多线程操作队列时出现空弹出错导致协程终止。
  • convert_oldest方法异常未捕获:如果该方法执行时抛出异常且没有被捕获,会直接终止check_queue协程,导致后续队列检查逻辑不再执行,需要补充异常捕获逻辑保证协程持续运行。

修复后的核心代码

首先调整sleep导入:

from tornado.gen import sleep

修改check_queue逻辑:

async def check_queue(executor, loop):
    while True:
        log(f"Checking queue: {MainLoop.queue.qsize()}")
        if not MainLoop.queue.empty():
            try:
                await loop.run_in_executor(executor, MainLoop.queue.convert_oldest)
            except Exception as e:
                log(f"Process file failed: {str(e)}")
        await sleep(1)

优化方案

当前每秒轮询队列的方案存在处理延迟和空转开销,可替换为asyncio.Queue异步队列,无需轮询即可实现新消息实时处理:

  1. 将MainLoop.queue替换为asyncio.Queue实例,在convert_oldest中适配异步队列的取值逻辑
  2. 修改check_queue为如下实现:
async def check_queue(executor, loop):
    while True:
        # 协程阻塞等待新消息,有消息存入时自动唤醒
        file_name = await MainLoop.queue.get()
        log(f"Start processing file: {file_name}, queue size: {MainLoop.queue.qsize()}")
        try:
            await loop.run_in_executor(executor, MainLoop.queue.convert_oldest, file_name)
        except Exception as e:
            log(f"Process file {file_name} failed: {str(e)}")
        finally:
            MainLoop.queue.task_done()
  1. WebSocket收到消息时直接使用await MainLoop.queue.put(msg)存入队列,注意on_message方法要加上async修饰符。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 22:36:08