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

如何用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. 运行与测试步骤

  1. 启动NATS服务器
    • 使用Docker快速启动:
      docker run -p 4222:4222 nats
      
    • 或下载官方NATS服务器二进制包直接启动
  2. 运行Python代码
    • 执行脚本:python your_script_name.py
  3. 测试手动触发boo
    • 用NATS命令行工具发送空消息到boo主题:
      nats pub boo ""
      
    • 控制台会输出手动触发的boo执行结果
  4. 测试foo计算
    • 发送带参数的JSON消息到foo.calculate主题:
      nats pub foo.calculate '{"a": 10, "b": 20}'
      
    • 会收到返回的计算结果,控制台也会打印计算日志

4. 关键细节说明

  • NATS连接配置:如果使用远程NATS集群,修改nc.connect()中的URL即可
  • 定时间隔调整:修改scheduled_boo_task里的asyncio.sleep(30)参数,调整boo的定时执行频率
  • 消息格式:统一使用JSON传递参数,确保跨服务调用的兼容性
  • 错误处理:示例包含基础的参数校验和异常捕获,可根据业务需求扩展更复杂的错误处理逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 23:15:36