如何保证消息的幂等性
为什么要实现消息幂等
MQ通常只保证至少一次投递,不保证不重复:生产者发送重试、消费者处理成功但ACK丢失导致重投、消费超时重新投递,都可能导致同一业务消息被消费多次,所以消费者端必须做幂等
-
对于支付接口,如果用户因为网络波动问题,多次点击支付按钮,发送多条消息后,就可能导致用户账户被多次扣款
-
对于订单保存,假设发送了多条消息,就会导致数据库对同一个订单保存多条数据。最终导致数据的错误
-
消息队列发送消息是需要经过网络传输的,如果多次发送消息,可能会占用网络带宽,影响性能
对于上述问题,我们需要在消费者端实现消息幂等性处理
首先,设置一个唯一标识的业务ID,发送消息时带上这个唯一的业务ID。消费者端处理时,根据业务ID来进行处理。处理消息时,可以尝试插入这个唯一ID,如果插入成功,则标识这个消息是第一次到来,可以直接处理。插入失败,则表示之前已经处理过这条消息了,直接丢弃消息即可
如何解决消息丢失的问题
消息队列主要就分为三大块:生产者,中间件,消费者。所以只需要保证这三个环节不丢失消息即可
-
生产者端,对于RocketMQ来说,这一阶段主要是保证消息被可靠地发送到Broker
-
可以使用同步发送,相较于异步发送,同步发送能够立刻得到Broker的响应结果,明确知道消息是否发送成功
-
重试机制:同步消息发送失败,客户端会自动重试。可以设置合理的重试次数
-
失败兜底策略:如果重试后消息发送依旧失败,业务系统需要进行兜底处理,例如将消息存入数据库中,再由后台定时任务补发
-
事务消息:解决"本地事务执行与消息发送"的一致性问题。发送半消息后,Broker 对未反馈最终结果的半消息主动回查生产者本地事务执行结果,据此提交或回滚
-
Broker端:确保已接收的消息被安全、持久地保存
-
配置同步刷盘:Broker默认是异步刷盘(先写入内存缓存,再异步刷入磁盘)。为了防止宕机导致内存数据丢失,可以改为同步刷盘。消息写入磁盘后才返回成功响应
-
配置主从同步:RocketMQ支持主从架构。为了防止Master宕机造成数据丢失,可以配置同步复制,即等待Master消息同步到Slave后才返回成功
同步刷盘和同步复制都会降低系统吞吐量,增加延迟。这是在可靠性和性能之间的取舍
-
消费者端:确保消息被成功处理,避免因为消费失败导致消息丢失
-
采用手动ACK:Consumer应禁用自动提交消费位点。务必在业务逻辑处理成功后,才手动向Broker返回消费成功
-
消息重试和死信队列:如果Consumer处理失败,Broker会按策略重新投递该消息。当重试次数耗尽,仍然失败,消息会被转入死信队列,需要人工介入处理
怎么保证消息的可靠性、顺序性
可靠性
-
消息持久化:确保消息队列能够持久化消息。例如在系统崩溃、重启或网络故障的情况下,未处理的消息不应该丢失。这样在服务重启后,消息依然能被重新读取和处理
-
消息确认机制:消费者成功处理消息后,应该向消息队列发送确认。消息队列只有收到确认后,才会将消息从消息队列中移除。如果没有收到确认,消息队列可能会在一定时间后重新发送消息给其他消费者或再次发送给同一个消费者
-
消息重试策略:消费者处理消息失败后,需要有合理的重试策略。可以设置重试次数和重试间隔时间。若重试次数耗尽后仍然失败,发送到死信队列中的消息,可以让后续人工介入排查
顺序性
-
有序消息处理场景识别:对于需要顺序处理的消息,需要确保消息队列和消费者能够按照特定的顺序进行处理
-
消息队列对顺序性的支持:部分消息队列本身提供了顺序性保障功能。RocketMQ可以将消息发送到同一个消息队列中,在同一个消费者组内,一个队列同一时间只能被一个消费者实例消费
-
消费者顺序处理策略:消费者在顺序处理消息时,应该避免并发处理可能导致顺序打乱的情况。例如,可以通过单线程或者使用线程池并对顺序消息进行串行化处理,确保消息按照正确的顺序被消费
消息积压该怎么处理
消息积压是因为生产者的生产速度大于消费者的消费速度
-
首先要确认消息积压发生在哪个环节。查看监控,确认是哪个消费者组出现积压,检查消费者日志是否有异常,确认下游依赖是否正常
-
最直接的止损手段就是紧急扩容,提升消费能力。但是注意消费者的总数不能超过Topic的分区数,否则新增的消费者会限制。或者在消费者进程中适当调达消费线程池的大小。如果业务允许,也可以开启批量消费,提升吞吐量
-
降级非核心业务,如果扩容后积压仍然未缓解,可以考虑暂时关闭一些非核心业务的消费进程,把资源全部让给核心业务。
-
如果上述方法都无效,可以考虑先将消息全部导出到其他地方,然后情况队列,先恢复正常。之后再通过另一个程序慢慢处理这些导出的消息。
Redis数据持久化
Redis将内存中的数据持久化到磁盘中,一共有三种方案
-
AOF:记录每一次Redis执行写命令的内容,在日志文件中,将每次的命令追加写入到文件末尾。刷盘频率由appendfsync控制,生产常用everysec,每秒刷盘,最多只会丢失1s的数据;always每次写命令都fsync,最安全但是最慢;no刷盘时机交给操作系统,可能会丢失很多数据。但是由于记录的是命令,在一些多写的场景中,可能会导致aof日志文件过大。并且数据恢复时,需要一条条命令逐个执行,所以恢复速度较慢。所以redis利用了aof重写机制,对每个值的更新,只保留最新一次的更新。例如多次对name属性更新,只保留最后一次更新的命令。这样就能大大减少aof日志文件的体积
-
RDB:记录的是Redis数据库快照。在指定的时间间隔内,保存当前Redis的数据快照。RDB日志文件是二进制格式,恢复数据时,速度较快。可以直接加载二进制文件,就能恢复数据。但是可靠性较低,因为如果两次记录快照的时间较长,就会丢失大量的数据。RDB由fork出的子进程以Copy-On-Write的方式生成快照,生成期间不阻塞主线程命令。但是fork瞬间主线程会出现短暂阻塞,且写时复制会带来额外内存和CPU开销
-
Redis4.0+推荐混合持久化:RDB做全量基底+AOF做增量,兼顾恢复速度和数据安全



