RocketMQ 如何保障消息的可靠性
消息在 RocketMQ 中会经历三个阶段:生产阶段、存储阶段和消费阶段。因此,为了保证消息的可靠性,需要综合考虑这三个阶段。
生产
通过请求确认机制,保证消息正确发送。
- 同步发送的时候,要注意响应结果和异常。如果响应失败或发生了异常,需要进行重试
- 异步发送可以在回调函数里面检查响应结果和异常
- 如果发生了超时,可以通过查询日志的 API,检查是否在 Broker 中存储成功
步骤 1 就是通过同步发送保证在生产阶段消息的可靠性。
存储
配置可靠性优先的 Broker 参数来避免因为宕机丢失消息。
- 消息持久化到 CommitLog 中,这样即使 Broker 宕机,未消费的消息也能重新恢复
- Broker 采用同步刷盘机制,在 Producer 发送消息后,等消息持久化到硬盘之后再返回响应给 Producer
消费
确保消费逻辑执行完成之后,再发送消费确认信息。如果在消费逻辑执行完成到发送消费确认消息之间,服务出现问题了,那就会导致消息重复消费的问题。RocketMQ 在保证消息一定投递的情况下,有可能造成消息重复。
要处理消息重复的问题,需要靠业务端去保证,主要有两种方式:业务幂等和消息去重。
业务幂等
保证消费逻辑的幂等性,即无论多少次调用消费逻辑,对业务都不会产生影响。
音频预标注异步架构 通过业务幂等来解决消息重复消费的问题,因为音频无论经过多少次重复标注,得到的结果都是幂等的。
消息去重
- 利用
MessageId去重:Consumer 维护一张已处理消息 ID 列表,如果消费的消息 ID 已经存在这张表里了,那么跳过该消息。注意这里需要使用事务来确保消息 ID 入库和业务逻辑处理完成这两个操作具有原子性 - 利用业务唯一标志去重:Consumer 在接收到消息之后,从消息中提取业务唯一标识符,然后查询业务系统,判断该业务是否完成。此方法即需要消息中包含业务唯一标识符,又需要业务系统具有业务是否完成的查询功能