ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

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),仅供参考
返回列表