JetStream与消息可靠性
1435597771 ·
JetStream是什么
JetStream时NATs这个消息队列中内置的持久化流层,他在CoreNats的基础上增加了消息持久化、历史消息消费以及失败重投能力。
与CoreNats最多一次的消息语义相反,JetStream采用至少一次投递语义,即只要消息没有被确认,系统会尝试重新投递消息,使消息有机会被消费者处理。之所以说是一定程度是因为他的消息重放次数达到上限后不再重新重放。
JetStream的组成部分
Stream:Stream用来持久化存储消息,他订阅一个对应的Subject,接收消息,然后存储起来,并为这个消息分配序列号。
Consumer:虽然Consumer翻译过来是消费者,但是他实际上是JetStream 中保存消费状态的组件,它记录消息投递进度、确认状态、投递次数等信息。
Worker:实际执行业务逻辑的应用程序,它通过 Consumer 获取消息,并完成订单处理、积分增加等业务操作。
PubAck和Ack:PubAck 是 JetStream 返回给 Publisher 的确认响应,用来表示消息已经被成功接收并存入 Stream,同时返回 Stream 名称和消息序列号。Ack 是 Worker 返回给 Consumer 的确认信号,用来表示这条消息对应的处理已经完成。在一般情况下是等到业务正确处理完后才返回Ack
Advisory:在JetStream里,Advisory是服务器主动发布的一种特殊系统事件消息,用来通知外部“发生了一些重要的事情”,他的Subject前缀是:$JS.EVENT.ADVISORY.>,比如当消息达到最大投递次数上限,那么他就会发送$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.<Stream>.<Consumer>,用来表示某个Stream的某个Consumer达到了最大投递次数上限,需要处理。
消息可靠性
消息可靠性指的就是: 在消息从发送方传到接收方的过程中,系统能不能保证消息被送达。 它核心主要解决的就是:消息会不会丢?会不会重复?
NATs提供的JetStream就是为了保障消息可靠性,我们接下来一个问题一个问题分析,看JetStream是如何保障消息可靠性的
消息会不会丢:
从消息生产与消费粗略来看,需要有生产者,Server,消费者这三者,那么消息丢失的情况有:
生产者与Server之间
消费者与Server之间
JetStream提供了PubAck与Ack来解决消息丢失的问题。
只有当Stream返回给生产者PubAck的时候,我们才能认定消息被Stream持久化地存储起来了。
只有 Worker 返回 Ack 后,Consumer 才认为这条消息已经完成处理。并不是说业务被正确处理了
那么在这个情况下,JetStream是如何实现的呢?
在文章开头的时候讲过,JetStream的消息语义是至少一次。JetStream使用Consumer这个组件来记录消息的状态。消息的状态有:
未投递:消息已经存在于 Stream 中,但当前 Consumer 还没有将它投递给 Worker
已投递+等待确认:当某个Worker接收了这个消息,那么他必须给Consumer返回一个Ack,否则Consumer的状态会一直在等待确认的状态中,等到Consumer的等待时间到了,那么这个消息就会被重新投送到某个Worker,让Worker重新处理这个消息
已确认:当Woker返回Ack的时候,Consumer接收到这个Ack,则说明这个消息被Worker处理完成。Ack改变的是Consumer的消费状态,不代表消息立即从Stream中删除,消息是否删除还取决于Stream的保留策略。
JetStream通过Consumer记录消息的状态流转,来确保消息至少一次投递。
会不会重复
重复有两种情况:
发布重复:JetStream 可以通过消息 ID(Nats-Msg-Id)进行发布去重,这种去重只针对消息发布阶段的重复提交,并不能解决消费阶段业务重复执行的问题。在存储消息的时候,如果消息的id相同,那么JetStream则不会存储。但是如果消息的id不同内容相同,JetStream还是会给他当作不同的消息来处理。
消费重复:由于至少一次投递,业务仍可能重复执行,需要业务幂等。
所以通常来讲消息会不会重复,指的是业务幂等性,也就是消息会不会重复消费,以及消息重复消费带来的不良影响。
前面讲过JetStream的消息语义是至少一次,所以消息是在很大程度上会被重新消费的。如果不能保证业务幂等性,那么会导致业务被重复处理。特别设计金额相关,可能会导致用户重复扣费的严重问题。
但是JetStream并不难解决消息重复消费的问题,需要我们在业务场景上额外处理。
一般的处理方法是:业务唯一键+原子判断。
业务唯一键:业务唯一键不是消息的id,他是根据业务选择被包括在消息中的数据字段。通过不同的业务场景,选择不同的字段并把它存储在数据库中来判断有没有被处理过。比如支付场景的消息体可能是:
{
"payment_id": "pay-456",
"order_id": "order-123",
"amount": 100
}那么我们就可以选择payment_id来充当业务唯一键。
原子判断:原子性是一种性质,用来保证这个操作要么成功要么失败,不会有中间状态,而在平常使用中,则是通过事务来实现。
所以在保证消息幂等性场景里,用业务唯一键作为唯一约束,用事务来保证原子判断,从而来保证业务会不会重复
那么业务唯一键的存储时机是什么?
worker处理业务过程中,他会给这个payment_id尝试插入正在处理表中,如果插入失败,需要根据错误类型判断。如果是业务唯一键冲突,则说明该业务已经处理过,可以直接 Ack;如果是数据库异常,则不能认为业务已经完成。 如果插入成功,则说明没有worker处理过这个消息的业务,他会在事务中处理业务。提交后返回Ack。
综上所述,消息可靠性并不能仅仅通过JetStream单独来实现,他需要根据不同的业务来选择不同的唯一性,但是JetStream 提供了实现可靠消费所需要的机制,但最终的业务可靠性仍然需要业务层通过幂等、事务和失败处理共同保证。