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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 18:02:45