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

在Tornado重写方法中调用异步函数遇阻塞及重复调用失效问题求助

问题描述

在Tornado的WebSocketHandler重写方法中调用异步函数时遇到以下问题:

  • 最初将generate()设为同步方法时,会阻塞整个程序(包括on_connection_close方法)。
  • 改成异步的generate()后,调用read()方法首次会打印“enter read...”,但后续调用不再输出该内容;若注释掉read方法中Kafka消费相关代码,功能恢复正常。

相关代码如下:

Reader类代码

class Reader:
  
    def __init__(self, topic, group_id):
        self.topic = topic
        self.consumer = AIOKafkaConsumer(topic, bootstrap_servers="kafka:9092",
                                         group_id=group_id)

    async def read(self, sender):
        print("enter read...")
        # 疑似阻塞代码
        await self.consumer.start()
        async for message in self.consumer:      
            sender(json.dumps(consumed_message))
            await self.consumer.commit()

SampleSocketHandler类代码

class SampleSocketHandler(tornado.websocket.WebSocketHandler):
    def on_connection_close(self):
        self.app.close()

    def on_message(self, message):
        generate(message)
        # await generate(message) 无法直接在同步方法中使用await

    async def generate(self, message):
        await self.reader.read(message)

主函数代码

async def main():
    app = Application()
    tornado.ioloop.IOLoop.instance().start()

if __name__ == "__main__":
    asyncio.run(main())
问题分析
  1. 同步回调调用异步函数的错误方式:on_message是Tornado的同步回调方法,直接调用异步的generate()只会返回协程对象,不会实际调度执行;如果强行用同步逻辑等待异步任务,会直接阻塞整个事件循环,导致on_connection_close等其他回调无法响应。

  2. AIOKafkaConsumer的生命周期问题:

    • AIOKafkaConsumer.start()只能被调用一次,重复调用会触发异常或未知行为。
    • async for message in self.consumer是无限循环,会持续占用当前协程,只要Kafka有消息就会一直消费,后续调用read()时,同一个Reader实例的协程仍在运行,无法进入新的read()执行流程,所以后续不会打印“enter read...”。
  3. 事件循环整合错误:tornado.ioloop.IOLoop.instance().start()是阻塞调用,会卡住asyncio.run(main())的执行流程,Tornado与asyncio的事件循环未正确整合,导致异步任务调度异常。

解决方案

1. 正确在同步回调中启动异步任务

修改on_message方法,通过Tornado的IOLoop调度异步任务:

class SampleSocketHandler(tornado.websocket.WebSocketHandler):
    def on_connection_close(self):
        self.app.close()
        if hasattr(self, 'reader'):
            self.reader.stop_read()

    def on_message(self, message):
        # 用IOLoop调度异步任务,避免阻塞同步方法
        tornado.ioloop.IOLoop.current().add_callback(self.generate, message)

    async def generate(self, message):
        # 传递WebSocket的消息发送方法作为sender
        await self.reader.read(self.send_message)

    def send_message(self, msg):
        # 确保连接未关闭时再发送消息
        if self.ws_connection is not None:
            self.write_message(msg)

2. 重构Reader类,管理消费生命周期

修改Reader类,确保consumer仅启动一次,并支持终止消费循环:

class Reader:
    def __init__(self, topic, group_id):
        self.topic = topic
        self.consumer = AIOKafkaConsumer(topic, bootstrap_servers="kafka:9092",
                                         group_id=group_id)
        self.consumer_started = False
        self.running = False  # 标记是否正在消费

    async def _start_consumer(self):
        if not self.consumer_started:
            await self.consumer.start()
            self.consumer_started = True

    async def _stop_consumer(self):
        if self.consumer_started:
            await self.consumer.stop()
            self.consumer_started = False

    async def read(self, sender):
        print("enter read...")
        await self._start_consumer()
        self.running = True
        try:
            async for message in self.consumer:
                if not self.running:
                    break
                # 正确处理Kafka消息,原代码的consumed_message未定义
                msg_content = json.dumps({
                    "topic": message.topic,
                    "value": message.value.decode()
                })
                sender(msg_content)
                await self.consumer.commit()
        finally:
            await self._stop_consumer()

    def stop_read(self):
        # 设置终止标志,让消费循环退出
        self.running = False

3. 修复事件循环整合问题

修改main函数,正确整合Tornado与asyncio的事件循环:

async def main():
    app = Application()
    # 监听指定端口
    app.listen(8888)
    # 安装AsyncIOMainLoop,让Tornado使用asyncio的事件循环
    tornado.platform.asyncio.AsyncIOMainLoop().install()
    # 用asyncio的事件等待,避免阻塞主线程
    await asyncio.Event().wait()

if __name__ == "__main__":
    asyncio.run(main())
关键注意点
  • 禁止在Tornado的同步回调(如on_message、get)中直接阻塞等待异步任务,必须通过IOLoop.add_callback或asyncio.create_task调度。
  • AIOKafkaConsumer的启动/停止操作必须成对出现,避免重复启动引发异常。
  • 异步无限循环(如Kafka消费)必须设置终止条件,否则会持续占用协程,导致后续任务无法执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 13:31:06