消费重试
消费重试是 RocketMQ 保证消息最终被处理的重要机制。它解决的问题是:消费者拿到消息后,业务处理可能因为网络抖动、数据库短暂不可用、第三方接口超时等原因失败,此时不能直接丢弃消息,而是应该过一段时间再处理。
为什么需要消费重试
假设订单支付成功后发送消息给积分系统,积分系统消费消息时数据库临时不可用。如果没有重试,这条积分变更消息就会丢失,用户支付成功但积分没有增加。
使用消费重试后,消费者可以先返回失败,RocketMQ 会在稍后重新投递消息。只要故障是短暂的,后续重试就能让业务最终成功。
flowchart TD
A[消费者拉取消息] --> B[执行业务逻辑]
B --> C{处理成功?}
C -->|成功| D[提交消费成功]
C -->|失败| E[返回消费失败]
E --> F[Broker 投递到重试队列]
F --> G[等待下一次重试]
G --> A重试适合解决什么问题
适合重试的问题一般有一个特点:过一会儿可能恢复。
- 数据库连接短暂失败。
- Redis、ES、第三方接口短时间不可用。
- 网络超时。
- 下游限流。
- 服务发布过程中短暂不可达。
不适合简单依赖重试的问题:
- 消息格式错误。
- 业务参数永远不合法。
- 下游接口已经废弃。
- 消费代码存在稳定 bug。
这类问题重试多少次都很难成功,应该尽快进入死信队列并报警。
消费返回结果
RocketMQ 消费端处理完成后,需要告诉 Broker 消费结果。
| 消费结果 | 含义 | 后续行为 |
|---|---|---|
| 成功 | 业务已经正确处理 | 提交消费进度,不再投递 |
| 失败 | 业务没有处理完成 | 进入重试流程 |
| 超时 | Broker 没有及时收到结果 | 可能重新投递 |
消费代码中不要在业务还没真正成功时就返回成功,否则 RocketMQ 会认为消息已经完成,后续不会再投递。
重试队列和死信队列
当消费失败时,RocketMQ 会把消息放入重试队列。超过最大重试次数后,消息会进入死信队列。
flowchart TD
A["原始 Topic"] --> B["ConsumerGroup 消费失败"]
B --> C["%RETRY% 重试队列"]
C --> D{"达到最大重试次数?"}
D -->|否| B
D -->|是| E["%DLQ% 死信队列"]重试队列
重试队列是 RocketMQ 为消费组维护的特殊队列。同一个 Topic 被不同消费组消费时,每个消费组都有自己的重试逻辑,互不影响。
这也意味着:A 消费组消费失败不会影响 B 消费组的进度。
死信队列
死信队列用于保存多次重试仍然失败的消息。进入死信队列后,通常不会再被普通消费者自动消费,需要人工排查或编写专门的补偿程序处理。
死信消息重点排查:
- 消息体是否符合当前代码预期。
- 业务数据是否已经被删除或状态不允许处理。
- 消费端是否有未修复异常。
- 是否需要人工修复数据后重新投递。
消费幂等
消费重试会带来重复消费的可能,所以消费端必须做幂等。
常见幂等方案:
- 数据库唯一索引:用订单号、业务流水号建立唯一约束。
- 状态机判断:只有当前状态允许流转时才处理。
- 去重表:消费前先插入消息 ID 或业务 ID,插入成功才执行。
- Redis
setnx:适合短期防重,但要注意过期时间和数据可靠性。
示例:积分增加消息可以使用 orderId + pointsType 做唯一键。如果这条记录已经存在,说明处理过,直接返回成功即可。
消费重试开发建议
- 消费失败时不要吞异常。吞掉异常并返回成功,会让失败数据永久丢失。
- 对可恢复异常返回失败,让 RocketMQ 重试。
- 对不可恢复异常记录清晰日志,并考虑直接进入人工补偿流程。
- 消费逻辑必须幂等。
- 日志里打印 Topic、MessageId、Key、ConsumerGroup、业务主键。
- 给死信队列配置监控和告警。
和事务消息的关系
事务消息 解决的是“生产者本地事务”和“消息发送”之间的一致性;消费重试解决的是“消息已经投递后,下游能否最终处理成功”。
两者配合后,一条典型链路是:
flowchart TD
A[订单本地事务] --> B[事务消息提交]
B --> C[积分系统消费]
C --> D{积分处理成功?}
D -->|成功| E[提交消费进度]
D -->|失败| F[消费重试]
F --> C所以不要误以为使用事务消息后消费端就不需要重试和幂等。事务消息只保证上游发送可靠,下游处理仍然要靠消费端自己保证。
代码 Demo:失败返回重试,成功再提交
下面示例演示消费者处理失败时返回 RECONSUME_LATER,让 RocketMQ 稍后重新投递。
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("POINTS_ORDER_PAID_CG");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("ORDER_EVENT", "ORDER_PAID");
consumer.registerMessageListener((MessageListenerConcurrently) (messages, context) -> {
for (MessageExt message : messages) {
String orderId = message.getKeys();
try {
if (alreadyProcessed(orderId)) {
continue;
}
addPoints(orderId);
markProcessed(orderId);
} catch (Exception e) {
System.err.println("消费失败,稍后重试,orderId=" + orderId);
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();幂等判断要放在业务执行前;业务成功后再记录处理完成。否则重复消息可能导致积分、库存、余额重复变更。
