FastStream 中的 NATS Direct Subject 消息路由:默认主题模型、队列组负载均衡与并发处理实战
FastStream 中的 NATS Direct Subject 消息路由默认主题模型、队列组负载均衡与并发处理实战【免费下载链接】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本篇技术指南围绕 FastStream 异步事件驱动框架对 NATSDirect Subject直接主题的支持展开讲解这一 NATS 最基本的路由方式在 FastStream 中的默认地位、声明方式、queue group 负载均衡机制、水平扩展与并发消费配置。读完本文你将掌握如何用broker.subscriber(subject, queue)声明 Direct Subject 消费者理解消息在多个消费者之间如何分配并能利用max_workers与多实例部署应对高吞吐场景。本文以仓库内 direct.md 为主体并结合 registrator.py 等源码实现加以印证。什么是 NATS Direct SubjectDirect Subject是 NATS 中路由消息的基本方式其本质非常简单一个subject主题会把消息发送给所有订阅了它的消费者。用一句话概括就是发布 - 订阅模型——发布者向某个subject发送消息NATS 服务器负责将该消息投递给当前所有订阅该subject的在线消费者。在 NATS 的路由体系中这是最基础的一类主题。与之相对的是Pattern Subject模式主题后者允许消费者通过通配符模式如*.info批量订阅一类主题。NATS 本身不提供复杂的路由规则配置能力唯一的路由实体就是subject——要么按名字直接订阅要么按正则表达式模式订阅详见 NATS 路由规则。这两种方式在 FastStream 中均有对应实现本文聚焦于第一种Direct Subject。Direct Subject 在 FastStream 中的默认地位FastStream 的 NATS 支持构建在nats-py客户端之上。在 FastStream 中Direct Subject 是默认使用的主题类型当你调用broker.subscriber(test_subject)而没有指定任何模式通配符时FastStream 就会创建一个 Direct Subject 订阅。这也意味着它是大多数 FastStream NATS 应用无需额外配置即可直接上手的路由方案。broker.handler(test_subject) async def handler(): ...说明broker.handler与broker.subscriber在 FastStream 中均可用于声明消息处理器前者是后者的通用别名形式。在 NATS 模块中官方示例统一使用broker.subscriber(...)语法。完整示例三个消费者两个主题下面是从仓库 direct.py 中提取的完整可运行示例它演示了 Direct Subject 在 FastStream 中最典型的使用方式from faststream import FastStream, Logger from faststream.nats import NatsBroker broker NatsBroker() app FastStream(broker) broker.subscriber(test-subj-1, workers) async def base_handler1(logger: Logger): logger.info(base_handler1) broker.subscriber(test-subj-1, workers) async def base_handler2(logger: Logger): # pragma: no branch logger.info(base_handler2) broker.subscriber(test-subj-2, workers) async def base_handler3(logger: Logger): logger.info(base_handler3) app.after_startup async def send_messages(): await broker.publish(, test-subj-1) # handlers: 1 or 2 await broker.publish(, test-subj-1) # handlers: 1 or 2 await broker.publish(, test-subj-2) # handlers: 3示例中FastStream(broker)将 NATS broker 与应用生命周期绑定app.after_startup保证在服务启动后自动向对应subject发布测试消息便于观察消息分配行为。消费者声明Consumer Announcement示例首先为两个subject声明了多个消费者test-subj-1和test-subj-2。注意subscriber装饰器的签名第一个位置参数是subject第二个位置参数是queue即队列组名称broker.subscriber(test-subj-1, workers) async def base_handler1(logger: Logger): logger.info(base_handler1) broker.subscriber(test-subj-1, workers) async def base_handler2(logger: Logger): logger.info(base_handler2) broker.subscriber(test-subj-2, workers) async def base_handler3(logger: Logger): logger.info(base_handler3)此处所有消费者都订阅在同一个queueworkers下这是模拟同一服务内多个消费者实例、演示负载均衡的写法。在真实单实例场景下同一服务内多个处理器共用同一队列组没有实际意义——因为消息会轮流进入这些处理器。这里只是为了模拟多消费者实例时 NATS 如何在其间做负载均衡。从源码签名看FastStream 的subscriber对queue参数的官方解释是Subscribers NATS queue name. Subscribers with same queue name will be load balanced by the NATS server订阅者的 NATS 队列名具有相同队列名的订阅者将由 NATS 服务器进行负载均衡定义位于 registrator.py。消息分配Message Distribution当上述消费者就绪后示例中发布的三条消息会这样分配await broker.publish(, test-subj-1) # handlers: 1 or 2消息 1会被投递给handler1或handler2中的任意一个——因为这两个消费者在同一个queue groupworkers下监听着同一个test-subj-1subject。NATS 会在组内随机选择一个消费者进行投递。await broker.publish(, test-subj-1) # handlers: 1 or 2消息 2的分发方式与消息 1 完全一致同样由test-subj-1队列组内的handler1或handler2随机承接即每次投递都是独立随机的可能连续落在同一个处理器上。await broker.publish(, test-subj-2) # handlers: 3消息 3则必然投递给handler3——因为它是唯一订阅test-subj-2的消费者。水平扩展与负载均衡ScalingDirect Subject 与 queue group 的组合是 NATS 实现消费者水平扩展的核心手段。仓库内的 scaling.md 明确描述了这一机制如果同一个subject被多个具有相同queue group的消费者监听那么每条消息每次都会随机进入其中一个消费者。这意味着NATS 可以独立地在队列消费者之间均衡负载无需应用层实现任何分发逻辑想要提升某个队列的消息处理速度只需额外启动消费者服务的更多实例无需修改现有基础设施配置——NATS 会自动负责消息在各服务实例之间的分配启动N个同组消费者该subject的处理速度理论上可提升约N倍见 NATS 路由规则 中的相关说明。这种加实例即扩容的特性正是 Direct Subject queue group 在微服务场景下被广泛采用的原因扩容成本极低且对已有应用完全透明。单个订阅者的并发处理max_workers默认情况下所有 NATS 订阅者都以阻塞模式消费同一subject的消息即同一时间无法并发处理来自同一subject的多条消息相当于每个subject有一个隐式的处理锁。这会限制单实例内的吞吐上限。为此FastStream 为所有NatsBroker订阅者提供了max_workers参数允许以每个订阅者一个工作池的方式并发消费broker.subscriber(test-subj-1, workers, max_workers10) async def handler(logger: Logger): ...如上配置后单个订阅者可在同一时间并发处理最多10条消息见 scaling.md。从实现上看max_workers在 registrator.py 中被接收并传入create_subscriber内部以max_workers or 1兜底为串行随后由ConcurrentCoreSubscriber这类并发订阅者使用ConcurrentMixin开启独立消费任务、通过self._put_msg将消息放入并发队列。也就是说不传max_workers时得到的是CoreSubscriber——同一 subject 串行处理传入max_workers时框架自动选择ConcurrentCoreSubscriber——按工作池并发处理。并发消费的相关行为在 test_consume.py 中有对应测试覆盖例如以max_workers2声明订阅者并验证并发接收消息。从源码看 Direct 订阅的底层实现为了更深入理解 Direct Subject 订阅的机制可以查看 core_subscriber.py 中的CoreSubscriber实现async def _create_subscription(self) - None: Create NATS subscription and start consume task. if self.subscription: return self.subscription await self.connection.subscribe( subjectself.subject.broker_address, queueself.queue, cbself.consume, **self.extra_options, )关键点在于queue参数被直接透传给底层nats-py的subscribe()调用由 NATS 服务器端完成同组消费者的负载均衡FastStream 本身不介入分发决策CoreSubscriber构造时使用的是subject.broker_address——对于 Direct Subject这就是你声明的字面主题名如test-subj-1与 Pattern 主题使用通配符模板地址的方式不同在解析器方面CoreSubscriber为 NATS Core 场景设置了is_ack_disabledTrue核心订阅不涉及 ack 策略这也呼应了 NATS 非持久化、无消息确认机制的特性详见 NATS 优劣势。此外subscriber装饰器还支持pending_msgs_limit、pending_bytes_limit等参数在 NATS Core 模式下若未确认消息数超出限制会触发服务器的 Slow Consumer 错误提示当前工作进程无法承担全部负载——这同样是判断是否需要扩容或调高max_workers的信号。Direct 与 Pattern如何选择维度Direct Subject本文Pattern Subject对照订阅方式按字面主题名精确订阅按通配符模式如*.info订阅路由依据主题名完全匹配消息键与模式匹配典型场景点对点、简单发布订阅日志分级、多租户主题分类声明示例broker.subscriber(test-subj-1, workers)broker.subscriber(*.info, workers)官方文档direct.mdpattern.mdPattern Subject 的完整示例与消息分配演示见 pattern.md其配套源码 pattern.py 中handler1/handler2监听*.info、handler3监听*.error发布logs.info时消息在 1、2 之间随机分配发布logs.error时固定进入 3。两者都支持 queue group 负载均衡选择哪一种取决于业务是需要精确主题还是分类主题。使用注意事项与限制基于仓库文档与源码使用 FastStream 的 NATS Direct Subject 时需要注意以下几点默认串行消费同一订阅者在未设置max_workers时按 subject 阻塞消费高吞吐场景请显式配置max_workers或多实例部署消息不持久化NATS Core 模式下消息不落盘消费者离线期间发布的消息会丢失需要持久化与投递保证时应改用 NATS JetStreamNatsJS能力FastStream 中对应stream、durable等订阅参数见 registrator.py无消息确认机制Core 订阅没有 ack 语义CoreSubscriber在实现中即通过is_ack_disabledTrue显式关闭了 ack 策略同一服务内避免重复组在单实例内为同一主题声明多个同队列组消费者并无负载均衡收益消息只是轮流进入不同处理器多实例部署才是 queue group 的正确打开方式queue group 与主题一一对应不同subject之间的队列组相互独立如上例中test-subj-2的组内只有handler3消息不会跨主题投递。通过 Direct Subject 这一基础但强大的路由原语结合 queue group 与max_workers你可以在 FastStream 中快速搭建出可水平扩展、无额外基础设施改动的 NATS 事件驱动服务如需更强的交付保证再平滑升级到 JetStream 即可。【免费下载链接】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),仅供参考