FastStream 与 NATS JetStream:利用 `in_progress()` 实现长任务消息心跳保活
2026/9/18 12:25:23 网站建设 项目流程

FastStream 与 NATS JetStream:利用in_progress()实现长任务消息心跳保活

【免费下载链接】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

NATS JetStream 默认采用 "at least once"(至少一次)投递语义:只要消息尚未收到 ACK,服务端就会持续重试投递,即使你的处理器需要很长时间才能完成消息处理。本文基于 FastStream 官方 HowTo 文档中的In-Progress sender模式(对应仓库文件 docs/docs/en/howto/nats/in-progress.md),讲解如何借助NatsMessage.in_progress()向 JetStream 周期性上报"消息仍在处理中"的状态,从而在不丢失消息、不产生重复投递的前提下优雅地运行耗时业务逻辑。读完本文,你将掌握 in-progress 心跳模式的完整实现、其底层源码调用链、批次消息场景下的用法以及测试时的模拟机制。

一、问题背景:JetStream 的至少一次投递与 ACK 语义

在使用 NATS JetStream 的消费者时,一个重要的默认行为是at least once(至少一次)投递原则。这意味着服务端会持续尝试把消息投递给消费者,直到它收到消费者返回的ACK(确认)状态为止。如果消费者的处理进程在确认前崩溃、超时或断开,消息会被重新投递,从而保证消息不丢失。

这种设计带来的直接后果是:

  • 消息一定会被送达,除非显式确认;
  • 消息可能被重复处理(重新投递),如果没有妥善处理幂等或确认逻辑;
  • 处理耗时长的消息时,如果迟迟不确认,可能导致 JetStream 判定投递失败并重试,造成同一消息被并发或重复消费。

in_progress()(NATS 术语中也称progress/ heartbeat ACK)正是 JetStream 为上述场景提供的官方机制:它允许消费者在不完成整个消息处理的情况下,向服务端发送一个"我还在处理,请别超时、别重投"的心跳信号。FastStream 把这一能力封装到了NatsMessagein_progress()方法上,形成本文要讲解的In-Progress sender模式。

二、核心 API:NatsMessage.in_progress()的源码实现

在 FastStream 中,NATS 消息的底层封装位于 faststream/nats/message.py。对于普通(单条)JetStream 消息,其in_progress()实现非常简洁:

async def in_progress(self) -> None: if not self.raw_message._ackd: await self.raw_message.in_progress()

(对应 faststream/nats/message.py#L36-L38)

从源码可以看出两个关键点:

  1. 幂等保护:方法首先检查底层raw_message._ackd标志——一旦消息已经被 ACK,后续的in_progress()调用会被直接跳过,避免在确认之后继续向服务端发送无意义的心跳。
  2. 透传底层能力:实际动作是调用nats.aio.msg.Msg原生对象上的in_progress(),即完全使用nats-py客户端提供的 JetStream heartbeat ACK 能力,FastStream 只负责将其优雅地暴露在统一的消息接口上。

类似的幂等守卫逻辑也存在于ack()ack_sync()nack()reject()等方法中(见 faststream/nats/message.py#L11-L34),它们共同保证了消息确认状态机的正确性。

三、完整示例:In-Progress sender 模式逐行解析

原文档给出了一个可直接运行的完整示例(docs/docs/en/howto/nats/in-progress.md):

import asyncio from faststream import Depends, FastStream from faststream.nats import NatsBroker, NatsMessage broker = NatsBroker() app = FastStream(broker) async def progress_sender(message: NatsMessage): async def in_progress_task(): while True: await asyncio.sleep(10.0) await message.in_progress() task = asyncio.create_task(in_progress_task()) yield task.cancel() @broker.subscriber("test", dependencies=[Depends(progress_sender)]) async def handler(): await asyncio.sleep(20.0)

这个看似简短的模式实际上包含了几个非常值得深挖的设计点。

3.1 以依赖注入(Depends)承载消息生命周期

progress_sender是一个异步生成器,它通过Depends(progress_sender)被注入到订阅者handler的依赖链中(对应 faststream/nats/broker/registrator.py 中subscriberdependencies参数)。FastStream 的依赖注入机制会这样驱动它:

  1. 每收到一条消息,先执行progress_sender的前半段(到yield之前);
  2. 然后把控制权交给handler执行真正的业务逻辑;
  3. handler返回后,再回到生成器中执行yield之后的收尾代码(这里是task.cancel())。

因此,progress_sender天然拥有与消息处理完全一致的生命周期:消息开始处理时启动心跳任务,消息处理结束时取消心跳任务。这种写法把"保活"逻辑从业务代码中彻底剥离出来,handler本身只需要关心业务,无需感知任何 ACK 细节。

3.2 后台心跳任务:每 10 秒上报一次处理中状态

async def in_progress_task(): while True: await asyncio.sleep(10.0) await message.in_progress()

心跳任务的核心逻辑是:

  • 通过asyncio.create_task创建独立的协程任务;
  • 循环中先sleep(10.0)再调用message.in_progress(),即每 10 秒向 JetStream 发送一次"仍在处理"的心跳;
  • 心跳间隔(10 秒)应当明显小于JetStream 服务端配置的 ACK 等待/超时时间(ack_wait),否则心跳无法起到阻止重投的作用。

在该示例中,handlerasyncio.sleep(20.0)模拟了一个耗时 20 秒的长任务——远大于心跳间隔,因此整个处理过程中 JetStream 会持续收到心跳,不会因为"长时间未 ACK"而判定投递失败并重投。

3.3 为什么需要 yield 之后取消任务

task = asyncio.create_task(in_progress_task()) yield task.cancel()

如果不在handler处理完毕后调用task.cancel(),心跳任务会变成"幽灵任务"继续无限循环,产生以下问题:

  • 消息已被 ACK 后仍持续发送无意义的心跳(虽然_ackd守卫会拦截,但依然是无谓的开销);
  • 任务泄漏,长时间运行的服务中协程数量会持续增长,最终耗尽资源。

yield前后的对称结构(启动 ↔ 取消)正是依赖注入生成器模式的精髓:保证每条消息的副作用都能被精确清理

四、深入底层:心跳与 ACK 状态机的关系

要正确使用in_progress(),必须理解它与其他确认方法的关系。在 faststream/nats/message.py 中,JetStream 消息的完整确认工具集包括:

方法底层行为语义
ack()raw_message.ack()处理成功,确认消息,不再重投
ack_sync()raw_message.ack_sync()同步确认(阻塞等待服务端响应)
nack(delay=None)raw_message.nak(delay=delay)处理失败,可指定延迟后重新投递
reject()raw_message.term()终止消息,立即丢弃,不重投
in_progress()raw_message.in_progress()上报处理中状态,延长 ACK 等待窗口

可以看到,in_progress()ack/nack/reject属于完全不同性质的信号:它不是最终裁决,而是"还活着"的中间状态上报。实践中应当:

  • 在长任务运行期间周期性调用in_progress()保活;
  • 在任务结束时,仍需要依据业务结果调用ack()(成功)或nack()(失败重试)完成最终确认——in_progress()不能替代最终 ACK。

FastStream 的自动确认机制会与这些显式调用协同:默认情况下,handler正常返回会自动 ACK,抛异常则自动 NACK/REJECT;当你手动调用确认方法后,自动确认会依据内部状态机跳过重复操作(这正是_ackd守卫存在的意义)。

五、批次消息场景:NatsBatchMessage.in_progress()

如果使用批量订阅(batch consumer),FastStream 同样提供了对应支持。NatsBatchMessage(封装一组nats.aio.msg.Msg)的in_progress()实现于 faststream/nats/message.py#L74-L79:

async def in_progress(self) -> None: for m in filter( lambda m: not m._ackd, self.raw_message, ): await m.in_progress()

与单条消息版本的区别在于:批量版本会遍历批次内所有尚未确认的消息,逐条发送心跳。这样即使你正在批量处理一大批消息,JetStream 也能收到批次中每条消息的"处理中"信号,避免长批次处理被服务端误判为超时。

六、测试支持:内存模式下的 no-op 实现

FastStream 提供免依赖的测试模式(in-memory testing),在测试时无需连接真实 NATS 服务。为了保证测试环境中调用in_progress()不会出错,测试框架在 faststream/nats/testing.py#L311-L312 中为PatchedMessage提供了空实现:

async def in_progress(self) -> None: pass

这意味着:

  • TestNatsBroker等测试工具下,心跳调用被安全地静默忽略,不会因缺少真实连接而抛异常;
  • 你可以在测试中对包含 in-progress 心跳的订阅者直接进行消息投递与断言,无需 mock 底层心跳逻辑。

关于如何为使用该模式的消费者编写测试,可以参考仓库中 NATS 相关的测试用例目录 tests/brokers/nats,并结合 FastStream 的TestNatsBroker测试辅助类进行验证。

七、最佳实践与注意事项

综合原文档与源码实现,在使用 In-Progress sender 模式时建议遵循以下要点:

  1. 心跳间隔要小于服务端的ack_wait。JetStream 消费者可配置 ACK 等待时长(FastStream 的NatsBroker与订阅者均支持ack_policy相关配置,如 faststream/nats/broker/broker.py 中所述:ack_policy是所有订阅者的默认确认策略,单个订阅者可覆盖)。一般推荐心跳间隔设为ack_wait的 1/3 到 1/2,留出足够的网络余量。
  2. 务必在消息处理结束时取消心跳任务。利用依赖注入生成器的yield收尾段执行task.cancel(),防止协程泄漏。
  3. in_progress()不能替代最终 ACK。任务成功结束仍需正常返回(自动 ACK)或显式调用ack();失败则调用nack()/reject()触发重投或丢弃。
  4. 善用幂等保护in_progress()内部的_ackd检查意味着在 ACK 之后再调用它是安全的,你可以放心地在循环中反复调用。
  5. 批量场景使用NatsBatchMessage。批量消费者同样具备逐条心跳能力,逻辑与单条一致。
  6. 测试无额外成本。内存测试模式下in_progress()是 no-op,无需特殊 mock。

八、总结

FastStream 的 In-Progress sender 模式为 NATS JetStream 的耗时消息处理提供了一套优雅的解决方案:借助依赖注入把"周期性心跳"从业务逻辑中解耦,通过NatsMessage.in_progress()(底层透传nats-py的 JetStream heartbeat ACK)向服务端持续上报处理状态,从而在 at-least-once 投递语义下稳定运行长任务。其核心实现位于 faststream/nats/message.py,测试模拟位于 faststream/nats/testing.py,而本篇指南本身作为官方 HowTo 系列(见 docs/docs/en/howto/nats/index.md)的一部分,可直接复制到你的服务中作为参考实现。

【免费下载链接】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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询