十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

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

FastStream 与 NATS JetStream:利用 `in_progress()` 实现长任务消息心跳保活 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/faststreamNATS 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 把这一能力封装到了NatsMessage的in_progress()方法上形成本文要讲解的In-Progress sender模式。二、核心 APINatsMessage.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从源码可以看出两个关键点幂等保护方法首先检查底层raw_message._ackd标志——一旦消息已经被 ACK后续的in_progress()调用会被直接跳过避免在确认之后继续向服务端发送无意义的心跳。透传底层能力实际动作是调用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.mdimport 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 中subscriber的dependencies参数。FastStream 的依赖注入机制会这样驱动它每收到一条消息先执行progress_sender的前半段到yield之前然后把控制权交给handler执行真正的业务逻辑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否则心跳无法起到阻止重投的作用。在该示例中handler的asyncio.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(delayNone)raw_message.nak(delaydelay)处理失败可指定延迟后重新投递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 consumerFastStream 同样提供了对应支持。NatsBatchMessage封装一组nats.aio.msg.Msg的in_progress()实现于 faststream/nats/message.py#L74-L79async 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 模式时建议遵循以下要点心跳间隔要小于服务端的ack_wait。JetStream 消费者可配置 ACK 等待时长FastStream 的NatsBroker与订阅者均支持ack_policy相关配置如 faststream/nats/broker/broker.py 中所述ack_policy是所有订阅者的默认确认策略单个订阅者可覆盖。一般推荐心跳间隔设为ack_wait的 1/3 到 1/2留出足够的网络余量。务必在消息处理结束时取消心跳任务。利用依赖注入生成器的yield收尾段执行task.cancel()防止协程泄漏。in_progress()不能替代最终 ACK。任务成功结束仍需正常返回自动 ACK或显式调用ack()失败则调用nack()/reject()触发重投或丢弃。善用幂等保护。in_progress()内部的_ackd检查意味着在 ACK 之后再调用它是安全的你可以放心地在循环中反复调用。批量场景使用NatsBatchMessage。批量消费者同样具备逐条心跳能力逻辑与单条一致。测试无额外成本。内存测试模式下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),仅供参考
返回列表