在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())
问题分析
同步回调调用异步函数的错误方式:
on_message是Tornado的同步回调方法,直接调用异步的generate()只会返回协程对象,不会实际调度执行;如果强行用同步逻辑等待异步任务,会直接阻塞整个事件循环,导致on_connection_close等其他回调无法响应。AIOKafkaConsumer的生命周期问题:
AIOKafkaConsumer.start()只能被调用一次,重复调用会触发异常或未知行为。async for message in self.consumer是无限循环,会持续占用当前协程,只要Kafka有消息就会一直消费,后续调用read()时,同一个Reader实例的协程仍在运行,无法进入新的read()执行流程,所以后续不会打印“enter read...”。
事件循环整合错误:
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
相关产品推荐
相关产品推荐

