FastStream 多 Broker 应用实战:跨 Kafka/NATS 的消息桥接、动态注册与端到端内存测试

发布时间:2026/9/18 21:23:03

FastStream 多 Broker 应用实战:跨 Kafka/NATS 的消息桥接、动态注册与端到端内存测试 FastStream 多 Broker 应用实战跨 Kafka/NATS 的消息桥接、动态注册与端到端内存测试【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststreamFastStream 应用默认围绕单个消息中间件构建但真实生产环境常常需要在同一进程内同时对接多个消息系统——例如将 Kafka 中的事件桥接到 NATS、在系统迁移期间让新旧两个 RabbitMQ 集群并行运行。本篇技术指南以 multiple_brokers.md 为核心系统讲解如何向FastStream应用传入多个 Broker、如何用add_broker动态注册、如何让每个 Broker 保持独立的订阅者/发布者并深入源码与测试带你掌握多 Broker 场景下的端到端内存测试技巧。为什么需要多个 Broker大多数 FastStream 应用围绕单个 Broker 构建但有些场景天然需要一个进程、多个消息系统桥接Bridging从一个 Broker 消费消息再重新发布到另一个 Broker实现异构消息系统之间的数据联通渐进式迁移Migration从一个 Broker 迁移到另一个 Broker 时两者必须并行运行一段时间逐步切换流量多集群/多环境并存同一应用需要同时连接多个 Kafka 集群、多个 NATS 或 RabbitMQ 实例。为了支持这些场景FastStream应用的构造函数接受多个 Broker。每个 Broker 都是一个完全独立的对象各自持有自己的订阅者subscriber和发布者publisher应用启动时统一连接所有 Broker关闭时统一断开。从源码看这一设计在 faststream/_internal/application.py 中落地_init_setupable_会遍历传入的 brokers 逐个调用add_broker注册_start_broker在启动时依次await b.start()stop时则逐个await broker.stop()保证多个 Broker 生命周期完全由应用统一管理。将多个 Broker 传给应用构造函数只需把想要运行的所有 Broker 实例作为位置参数传给FastStream(...)构造函数即可。下面是一个标准的Kafka → NATS桥接示例完整代码见 app.pyfrom faststream import FastStream from faststream.kafka import KafkaBroker from faststream.nats import NatsBroker kafka_broker KafkaBroker(localhost:9092) nats_broker NatsBroker(nats://localhost:4222) app FastStream(kafka_broker, nats_broker) kafka_broker.subscriber(incoming) nats_broker.publisher(outgoing) async def from_kafka(msg: str) - str: # Bridge the message from Kafka to NATS return msg nats_broker.subscriber(outgoing) async def from_nats(msg: str) - None: print(fReceived from NATS: {msg})当应用运行时两个 Broker 都会建立连接from_kafka处理器注册在kafka_broker上订阅主题incomingnats_broker.publisher(outgoing)装饰器把from_kafka的返回值路由到nats_broker的outgoing主题from_nats处理器注册在nats_broker上订阅outgoing最终收到来自 Kafka 的消息。于是一条来自 Kafka 的消息最终被投递到 NATS 的订阅者手中。关键点在于每个 Broker 都是独立对象你要把订阅者和发布者挂到你确切想要的那个 Broker 上系统之间的路由因此变得显式、可控。底层视角发布者装饰器与订阅者的绑定从实现上看nats_broker.publisher(...)装饰器会把处理器返回值写入 NATS 的outgoing主题这与kafka_broker.subscriber(incoming)注册的消费入口在同一个处理器上叠加。FastStream 的多 Broker 应用不会隐式地自动转发任何消息——只有当你显式地使用发布者装饰器、或在处理器内部调用某个 Broker 的publish跨 Broker 路由才会发生。这正是桥接逻辑清晰、可读性强的原因。用 add_broker 动态注册 Broker如果构造应用时并非所有 Broker 都已就绪例如某些 Broker 的地址在运行时才确定可以先用部分 Broker 创建应用之后再通过add_broker方法注册其余 Broker完整代码见 add_broker.pyfrom faststream import FastStream from faststream.kafka import KafkaBroker from faststream.nats import NatsBroker kafka_broker KafkaBroker(localhost:9092) nats_broker NatsBroker(nats://localhost:4222) app FastStream(kafka_broker) app.add_broker(nats_broker)add_broker与在构造函数中直接传入 Broker 完全等价。从 faststream/_internal/application.py 的源码可以看到其真实行为def add_broker(self, broker: BrokerUsecase[Any, Any, Any]) - None: if broker in self.brokers: msg fBroker {broker} is already added raise SetupError(msg) self.brokers.append(broker) self.schema.add_broker(broker) broker._update_fd_config(self.config)也就是说add_broker做了三件事去重校验如果同一个 Broker 实例被重复添加会抛出SetupError避免重复注册登记到 brokers 列表把 Broker 追加到self.brokers后续启动/停止时统一处理同步规格与依赖配置调用self.schema.add_broker(broker)把 Broker 纳入 AsyncAPI 文档生成范围同时用broker._update_fd_config(self.config)让 Broker 共享应用的依赖注入配置。注意第一个传入或添加的 Broker 会成为应用的默认Broker可通过app.broker访问所有已注册的 Broker 都在app.brokers列表中。这一约定在 application.py 中实现broker属性返回self.brokers[0] if self.brokers else None。多 Broker 应用的内存测试每个 Broker 都可以独立使用自己的TestBroker进行隔离测试。在多 Broker 场景下只需为应用使用的每一个Broker 包裹上对应的测试上下文管理器内存补丁就会覆盖所有 Broker桥接逻辑可以端到端地跑通——全程不需要任何真实运行的 Broker测试代码见 testing.pyimport pytest from faststream.kafka import TestKafkaBroker from faststream.nats import TestNatsBroker from .app import from_kafka, from_nats, kafka_broker, nats_broker pytest.mark.asyncio() async def test_bridge() - None: async with ( TestKafkaBroker(kafka_broker) as br, TestNatsBroker(nats_broker), ): await br.publish(Hi!, incoming) from_kafka.mock.assert_called_once_with(Hi!) from_nats.mock.assert_called_once_with(Hi!)测试流程非常直观TestKafkaBroker(kafka_broker)与TestNatsBroker(nats_broker)分别替换掉 Kafka 和 NATS 的底层连接与生产者进入纯内存模式通过br.publish(Hi!, incoming)向 Kafka 订阅者发布消息消息触发桥接逻辑from_kafka的返回值经nats_broker.publisher(outgoing)路由到 NATSfrom_nats处理器被调用最终断言from_kafka.mock与from_nats.mock各被精确调用一次。该测试在仓库中由 tests/docs/getting_started/multiple_brokers/test_app.py 通过require_aiokafka与require_nats标记自动执行保证文档示例与真实代码始终同步。同类型多 Broker共享一个 TestBroker如果应用运行了多个同类型的 Broker例如两个 Kafka 集群不需要为每个 Broker 单独准备上下文管理器。直接把所有 Broker 传给同一个TestBroker即可——它会同时修补传入的每个 Broker并在内存中完成消息路由示例应用见 same_type_app.py测试见 same_type_testing.pyfrom faststream import FastStream from faststream.kafka import KafkaBroker broker_1 KafkaBroker(localhost:9092) broker_2 KafkaBroker(localhost:9093) app FastStream(broker_1, broker_2) broker_1.subscriber(incoming) async def from_first(msg: str) - None: # Bridge the message from the first cluster to the second one await broker_2.publish(msg, outgoing) broker_2.subscriber(outgoing) async def from_second(msg: str) - None: print(fReceived on the second cluster: {msg})import pytest from faststream.kafka import TestKafkaBroker from .same_type_app import broker_1, broker_2, from_first, from_second pytest.mark.asyncio() async def test_bridge() - None: async with TestKafkaBroker(broker_1, broker_2) as (br1, _): await br1.publish(Hi!, incoming) from_first.mock.assert_called_once_with(Hi!) from_second.mock.assert_called_once_with(Hi!)这是官方推荐的同类型多 Broker 测试方式TestKafkaBroker(broker_1, broker_2)让两个集群保持连线状态从第一个集群桥接出去的消息能够投递到第二个集群的订阅者全程内存模拟无需真实 Kafka。从源码看同类型共享 TestBroker 之所以可行是因为 faststream/kafka/testing.py 中的TestKafkaBroker构造函数提供了针对单个 Broker 与多个 Broker 的重载签名其create_publisher_fake_subscriber会遍历self.brokers即传入的所有同类型 Broker上的全部订阅者来匹配发布目标FakeProducer的subscribers属性同样跨self.brokers聚合所有订阅者因此消息可以在多个同类型 Broker 之间完成内存路由。测试基础设施的基类定义在 faststream/_internal/testing/broker.pyTestBroker的_create_ctx会对每个传入 Broker 依次执行补丁、启动与停止保证多 Broker 上下文生命周期一致。混合类型的组合测试注意两种测试风格可以自由混用——每种 Broker 类型使用一个共享的TestBroker当应用同时组合不同类型时再把这些上下文管理器嵌套起来例如TestKafkaBroker(...)与TestNatsBroker(...)一起使用如前面 testing.py 中async with (...)的组合写法。从 FastStream 对象看多 Broker 的管理模型综合 faststream/app.py 与 faststream/_internal/application.py多 Broker 的管理模型可以归纳为能力实现方式源码位置构造时传入多个 BrokerFastStream(*brokers)位置参数展开faststream/app.py动态追加 Brokerapp.add_broker(broker)等价于构造传入faststream/_internal/application.py统一启动_start_broker依次await b.start()faststream/_internal/application.py统一关闭stop依次await broker.stop()faststream/_internal/application.py默认 Brokerapp.broker返回列表首个元素faststream/_internal/application.pyBroker 清单app.brokers列表faststream/_internal/application.py重复添加防护抛出SetupErrorfaststream/_internal/application.py另外值得注意多 Broker 也会统一进入应用的 AsyncAPI 规格生成流程add_broker中的self.schema.add_broker(broker)意味着你可以在同一份 AsyncAPI 文档中看到所有消息系统的通道定义便于整体查阅接口契约。实战要点小结路由是显式的FastStream 不会隐式转发消息跨 Broker 路由必须通过publisher装饰器或处理器内部的publish调用显式声明这让桥接逻辑一目了然生命周期统一无论构造时传入还是add_broker追加所有 Broker 都由应用统一启动、统一关闭无需手动管理连接默认 Broker 约定第一个 Broker 即默认 Broker可通过app.broker快速访问其余通过app.brokers索引访问测试两套配方异构多 Broker 用各自的TestBroker上下文管理器组合同类型多 Broker 用单个共享TestBroker一次传入多个实例内存测试全覆盖所有补丁均在内存中完成消息构造与投递详见 faststream/kafka/testing.py 中FakeProducer.publish的build_message与_find_handler流程桥接链路可以在无真实 Broker 的情况下端到端验证。无论你是要构建 Kafka↔NATS、RabbitMQ↔Redis 之类的消息桥接还是在多集群迁移阶段让新旧系统并行运行多 Broker 能力配合内存测试都能让你用一份清晰的代码和一套可复现的测试安全地把多个消息系统编排进同一个 FastStream 进程。【免费下载链接】faststreamAsynchronous Python framework for event-driven services. A thin client for Kafka, RabbitMQ, NATS, Redis and MQTT with full access to native broker features, plus AsyncAPI docs, in-memory tests and observability out of the box.项目地址: https://gitcode.com/GitHub_Trending/fa/faststream创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/18 21:23:03

多模态长期记忆Agent接模型,TaoToken 替换 Key 即可

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/18 21:23:03

C语言函数库手册PDF:man/groff导出与索引实践

简介:C语言函数库手册以PDF形式整理,面向正在学习或使用C语言做软件开发的学生、初学者与需随时查阅的工程师,用于解决函数名、参数及返回值记忆模糊、标准库分类不清晰的问题。全包仅1个PDF文件,约51KB,体积轻便&…

2026/9/18 21:23:03

教案结构化:用Python自动化生成家畜饲养学教案表格

简介:家畜饲养学教学教案文档(.doc)专为畜牧兽医专业师生设计,系统整理了家畜饲养学课程的核心教学框架。内容以绪论为起点,明确学习任务与研究方法,随后逐章展开畜禽营养原理,涉及植物性饲料与…

2026/9/18 22:23:06

VSCode背景美化实战:background-cover+自定义CSS配置指南

看腻了 VSCode 默认的深蓝黑灰界面?想让它更像自己的 IDE?说真的,这件事没有你想的那么玄乎。我试过好几个改背景的方案,最后稳定用下来的就是两样:background-cover 插件负责托底,自定义 CSS 样式负责精调…

2026/9/18 22:23:06

【NebulaGraph】在生产环境中,推荐的 NebulaGraph 集群部署拓扑结构是怎样的?Meta、Storage、Graph 节点应该如何分离?

NebulaGraph 3.8.0 生产部署拓扑权威指南:Meta、Storage、Graph 服务分离策略与最佳实践 用户问题原文:“在生产环境中,推荐的 NebulaGraph 集群部署拓扑结构是怎样的?Meta、Storage、Graph 节点应该如何分离?” 在金融反洗钱团伙挖掘场景中,图数据库集群需要7x24小时不间…

2026/9/18 22:23:06

虚幻引擎WebUI插件实战:从安装到跑通第一个网页界面

如果你用虚幻引擎做过带复杂界面的项目,应该能理解那种“UUMG 够用但很憋屈”的感觉。做按钮、列表、进度条还好,一旦牵扯到富文本、大数据表格、动态图表、后台管理面板,用 UMG 一个个拼控件简直是在给自己上刑。后来我在项目里引入了 WebUI…

2026/9/18 22:18:06

从汇编角度理解C语言篇 (三) —— C语言函数的实现

1. C语言函数组成// 返回类型 函数名 参数列表int add (int a, int b){// 函数体int ret a b;// 返回值return ret;}在C语言中,函数是执行特定任务的独立代码块。一个函数可以接收参数(如果有的话),执行一系列操作&#x…

2026/9/18 14:13:01

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/18 0:01:09

Google Colab 实战:运行模型、数据加载与报错排查

1. 为什么我劝你先搞懂 Colab 的运行模型1.1 Colab 到底是什么,跟本地跑代码差在哪Google Colab 简单说就是一台跑在浏览器里的 Linux 虚拟机,你打开一个 Notebook,背后就连上了一台带 GPU 的远程机器。你在单元格里敲的每一行 Python&#x…

2026/9/18 0:01:09

C语言数据类型与表达式详解

1. C语言数据与数据类型概述在C语言编程中,数据是程序处理的核心对象。理解数据的分类和特性是掌握C语言的基础。C语言中的数据主要分为四大类:常量、变量、表达式和函数。这些数据类型构成了C语言程序的基本元素,每种类型都有其独特的特性和…

2026/9/18 0:01:09

SQL时间字段指定时间段查询:区间语义、索引与时区避坑

上周排查一个线上问题&#xff0c;用户反馈"昨天的订单一条都没查到"&#xff0c;但数据库里明明躺着两千多条。最后定位下来&#xff0c;不是数据丢了&#xff0c;也不是接口挂了&#xff0c;而是那个查询条件把时间段写成了> 2024-05-20 00:00:00 AND < 2024…

2026/9/18 14:13:03

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行&#xff0c;Type-C接口算是典型的“看着简单&#xff0c;做起来全坑”的东西。光引脚就24个&#xff0c;高低速信号、电源、控制线全部塞在一个小小的连接器里&#xff0c;如果PCB布局不做规划&#xff0c;打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/18 14:13:02

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/18 14:13:02

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码