如何用Python结合NATS实现发布订阅,触发自定义函数并传参计算?求分步示例
NATS集成foo/boo函数分步实现示例
1. 安装依赖
- 安装NATS Python客户端及辅助库:
pip install nats-py python-dotenv
2. 完整实现代码
import asyncio import json from nats.aio.client import Client as NATS from nats.aio.errors import ErrConnectionClosed, ErrTimeout, ErrNoServers # 保留原有业务函数 def foo(a, b): return a + b def boo(): # 替换为实际映射逻辑 logic = True # 示例条件,按需修改 return 1 if logic else 0 async def run(): nc = NATS() # 连接NATS服务器(默认本地,远程需修改URL) await nc.connect("nats://localhost:4222") # --- 处理foo的参数计算请求 --- async def foo_msg_handler(msg): try: # 解析JSON格式的消息参数 req_data = json.loads(msg.data.decode()) a, b = req_data.get("a"), req_data.get("b") if a is None or b is None: await msg.respond(json.dumps({"error": "缺少a或b参数"}).encode()) return calc_result = foo(a, b) # 返回计算结果给请求方 await msg.respond(json.dumps({"result": calc_result}).encode()) print(f"执行foo计算:{a} + {b} = {calc_result}") except Exception as e: await msg.respond(json.dumps({"error": str(e)}).encode()) print(f"处理foo请求出错:{e}") # 订阅主题,收到消息即触发foo计算 await nc.subscribe("foo.calculate", cb=foo_msg_handler) # --- 定时自动触发boo --- async def scheduled_boo_task(): while True: # 每30秒执行一次(可修改定时间隔) await asyncio.sleep(30) boo_result = boo() print(f"定时执行boo,结果:{boo_result}") # 可选:将结果发布到NATS主题供其他服务消费 await nc.publish("boo.scheduled.result", json.dumps({"result": boo_result}).encode()) # 启动定时任务 asyncio.create_task(scheduled_boo_task()) # --- 手动触发boo(按需求,发送"boo"主题消息触发)--- async def boo_trigger_handler(msg): boo_result = boo() print(f"手动触发boo,结果:{boo_result}") await msg.respond(json.dumps({"result": boo_result}).encode()) await nc.subscribe("boo", cb=boo_trigger_handler) # 保持服务持续运行 while True: await asyncio.sleep(3600) if __name__ == "__main__": loop = asyncio.get_event_loop() try: loop.run_until_complete(run()) except KeyboardInterrupt: print("退出服务...") finally: loop.close()
3. 运行与测试步骤
- 启动NATS服务器
- 使用Docker快速启动:
docker run -p 4222:4222 nats - 或下载官方NATS服务器二进制包直接启动
- 使用Docker快速启动:
- 运行Python代码
- 执行脚本:
python your_script_name.py
- 执行脚本:
- 测试手动触发boo
- 用NATS命令行工具发送空消息到
boo主题:nats pub boo "" - 控制台会输出手动触发的boo执行结果
- 用NATS命令行工具发送空消息到
- 测试foo计算
- 发送带参数的JSON消息到
foo.calculate主题:nats pub foo.calculate '{"a": 10, "b": 20}' - 会收到返回的计算结果,控制台也会打印计算日志
- 发送带参数的JSON消息到
4. 关键细节说明
- NATS连接配置:如果使用远程NATS集群,修改
nc.connect()中的URL即可 - 定时间隔调整:修改
scheduled_boo_task里的asyncio.sleep(30)参数,调整boo的定时执行频率 - 消息格式:统一使用JSON传递参数,确保跨服务调用的兼容性
- 错误处理:示例包含基础的参数校验和异常捕获,可根据业务需求扩展更复杂的错误处理逻辑
内容的提问来源于stack exchange,提问作者SteveS
相关产品推荐
相关产品推荐

