使用Python nats.aio.client模块实现发布订阅时回调函数无法正常工作的问题咨询
我看了你的代码,发现几个关键问题导致回调函数没正常工作:
1. 类构造函数写错了
你的NAT类里定义的是def init(self):,但Python的类初始化方法必须是__init__(前后各两个下划线)。这个错误会导致self.nc根本没被初始化,后续的连接、订阅、发布操作都会抛出AttributeError,直接阻断流程。
2. 事件循环过早终止
你用run_until_complete依次执行三个协程:连接、订阅、发布。当publish_msg执行完成后,事件循环就直接退出了,但NATS的消息回调是异步触发的,此时回调函数还没来得及处理消息就被终止了。
修正后的完整代码
首先是nats_client.py的修正:
from nats.aio.client import Client as NATS class NAT: def __init__(self): # 改为双下划线的构造函数 self.nc = NATS() async def run(self): print("connection starts") await self.nc.connect("demo.nats.io:4222", connect_timeout=10, verbose=True) print("connection success") async def publish_msg(self): print("msg to publish") await self.nc.publish("Hello", b'Hellowelcome') async def subscribe_msg(self): async def message_handler(msg): print("Hello") subject = msg.subject reply = msg.reply print("Received a message on '{subject} {reply}'".format( subject=subject, reply=reply)) await self.nc.subscribe("Hello", cb=message_handler)
然后是主文件的修正:
import asyncio from nats_client import NAT async def main(): nat = NAT() await nat.run() await nat.subscribe_msg() await nat.publish_msg() # 留1秒时间让回调函数处理消息 await asyncio.sleep(1) # 最后关闭NATS连接 await nat.nc.close() if __name__ == "__main__": # 用asyncio.run统一管理事件循环(Python3.7+支持) asyncio.run(main())
这样修改后,事件循环会在发布消息后等待1秒,确保回调函数有足够时间执行。同时用asyncio.run代替手动管理事件循环,代码更简洁也更符合异步编程的最佳实践。
内容的提问来源于stack exchange,提问作者HARINI NATHAN
相关产品推荐
相关产品推荐

