本文主要介绍消息队列 TDMQ RocketMQ 版中消息重试与使用方法,涵盖 Remoting Java SDK 与 gRPC Java SDK 在不同消费模型下的重试行为差异。
功能介绍
当消息第一次被消费者消费后,没有得到正常的回应,或者用户主动要求服务端重投,TDMQ RocketMQ 版会通过消费重试机制自动重新投递该消息,直到该消息被成功消费,当重试达到一定次数后,消息仍未被成功消费,则会停止重试,将消息投递到死信队列中。
当消息进入到死信队列中,表示 TDMQ RocketMQ 版已经无法自动处理这批消息,一般这时就需要人为介入来处理这批消息。您可以通过编写专门的客户端来订阅死信 Topic,处理这批之前处理失败的消息。
说明:
只有当消费模式为集群消费模式时,Broker 才会自动进行重试,广播消费模式下不会进行重试。
出现以下三种情况会按照消费失败处理并会发起重试:
消费者返回
ConsumeResult.FAILURE。消费者返回
null。消费者主动/被动抛出异常。
消费重试与请求重试
本文重点讨论业务消费失败后的消息重投(用户业务返回失败后,消息是否再次投递给消费者),这与 SDK 内部的请求重试不是同一件事。
类型 | 说明 |
消费重试 | 用户业务处理失败后,消息再次投递给消费者 |
ACK 请求重试 | SDK 向服务端确认消费成功或失败后,重新发送 ACK 请求 |
ChangeInvisibleDuration 请求重试 | SDK 修改消息不可见时间失败后,重新发送修改请求 |
Receive 请求重试 | SDK 拉取消息失败后,重新发起拉取 |
最大重试次数与最大投递次数
RocketMQ 中需要区分两个概念:
名称 | 含义 |
最大重试次数 | 首次消费失败后,最多还能重新投递多少次 |
最大投递次数 / maxAttempts | 包含首次投递在内,一共最多投递多少次 |
通常关系为:
最大投递次数 = 首次投递 1 次 + 最大重试次数。例如消费组
retryMaxTimes =16,表示首次消费 1 次 + 最多重试 16 次 = 最多投递 17 次。在 gRPC 5.x 协议中,客户端使用 maxAttempts 表示总尝试次数,Proxy 会将消费组配置转换为 maxAttempts = retryMaxTimes + 1。重试次数
当消息需要重试时,TDMQ RocketMQ 中配置了如下的
messageDelayLevel 参数来设置重试次数与时间间隔。messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
由于 SDK 中第一次
delayLevel 为 3,所以重试次数与重试时间间隔关系如下:第几次重试 | 距离上一次重试的时间间隔 | 第几次重试 | 距离上一次重试的时间间隔 |
1 | 10秒 | 9 | 7分钟 |
2 | 30秒 | 10 | 8分钟 |
3 | 1分钟 | 11 | 9分钟 |
4 | 2分钟 | 12 | 10分钟 |
5 | 3分钟 | 13 | 20分钟 |
6 | 4分钟 | 14 | 30分钟 |
7 | 5分钟 | 15 | 1小时 |
8 | 6分钟 | 16 | 2小时 |
各消费方式默认重试次数总览
在 RocketMQ 5.x 架构下,不同 SDK、不同消费模型的默认重试行为存在差异,具体如下:
消费方式 | 默认最多重试次数 | 默认最多投递次数 | 最大重试次数控制位置 | 重试间隔控制位置 | 用户是否需要自己处理失败重试 |
Remoting PushConsumer 并发消费 | 16 | 17 | 客户端 setMaxReconsumeTimes | Broker delay level / setDelayLevelWhenNextConsume | 否 |
Remoting PushConsumer 顺序消费 | 近似无限 | 近似无限 | 客户端 setMaxReconsumeTimes | 本地队列短暂挂起后重试,默认 1000ms,可由 suspendCurrentQueueTimeMillis 控制 | 否 |
Remoting PushConsumer POP 模式 | 16 | 17 | 客户端 setMaxReconsumeTimes | 客户端指定 delay level;delay level 对应时长由服务端集群配置控制 | 否;仅作为代码路径说明,开源 Remoting SDK 不作为用户可直接使用的消费方式 |
gRPC PushConsumer 普通消息 | 16 | 17 | 消费组 retryMaxTimes | 消费组 groupRetryPolicy | 否 |
gRPC PushConsumer 顺序消息 | 16 | 17 | 消费组 retryMaxTimes | 消费组 groupRetryPolicy | 否 |
gRPC SimpleConsumer | 16 | 17 | 消费组 retryMaxTimes | 用户传入 invisibleDuration 或调用 changeInvisibleDuration | 是 |
核心结论:
Remoting PushConsumer:主要由客户端
setMaxReconsumeTimes 控制最大重试次数。gRPC PushConsumer:最大重试次数和重试间隔主要由消费组配置控制。
gRPC SimpleConsumer:最大重试次数由消费组控制,但每次失败后多久重新可见由用户控制。
使用方式
Remoting Java SDK PushConsumer
并发消费
Remoting Java SDK 并发消费通常使用
DefaultMQPushConsumer:DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GID_xxx");consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {try {// 业务处理return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;} catch (Throwable t) {return ConsumeConcurrentlyStatus.RECONSUME_LATER;}});
消费成功:业务返回
ConsumeConcurrentlyStatus.CONSUME_SUCCESS,SDK 认为消息消费成功,并更新消费进度。消费失败:业务返回
ConsumeConcurrentlyStatus.RECONSUME_LATER,SDK 会将消息送回重试链路,稍后重新投递给同一消费组下的消费者。最大重试次数:通过客户端控制
consumer.setMaxReconsumeTimes(16);,默认 16 次重试,即首次消费 1 次 + 最多重试 16 次 = 最多投递 17 次。超过最大重试次数后,消息进入死信队列 DLQ。重试间隔:可通过消费上下文设置下一次重试延迟级别:
context.setDelayLevelWhenNextConsume(3);return ConsumeConcurrentlyStatus.RECONSUME_LATER;
说明:
设置的是延迟级别,不是精确的秒数。若未显式设置,由服务端默认延迟策略控制。
顺序消费
顺序消费失败时通常返回
ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT。这类失败不是进入服务端 delay level 重试,而是客户端把消息重新放回当前 ProcessQueue,短暂挂起该队列后再次提交消费请求(默认约 1 秒后重试)。挂起时间默认
suspendCurrentQueueTimeMillis = 1000ms,最终调度时间会被限制在 10ms ~ 30000ms 之间。用户可通过 context.setSuspendCurrentQueueTimeMillis(1000) 或 consumer.setSuspendCurrentQueueTimeMillis(1000) 控制。顺序消费默认最大重试次数较特殊:若
maxReconsumeTimes = -1,默认近似 Integer.MAX_VALUE(近似无限重试)。可通过 consumer.setMaxReconsumeTimes(16) 改为有限重试。注意:
同一队列内,前面的消息一直失败会阻塞后续消息。顺序消费场景建议显式设置最大重试次数,并对不可恢复错误做好兜底处理。
POP 模式
POP 消息有不可见时间,如果没有 ACK 或主动修改不可见时间,消息会在不可见时间到期后重新可见。消费失败后通常通过
changePopInvisibleTime 安排下一次投递。最大重试次数由客户端 setMaxReconsumeTimes(16) 控制(默认 16 次、最多 17 次);POP 不可见时间 setPopInvisibleTime(60000) 默认 60000ms(有效范围 5000ms ~ 300000ms)。客户端设置的是 delay level,delay level 到具体延迟时间的映射由服务端集群的延迟级别配置控制。gRPC Java SDK PushConsumer
控制模型
gRPC PushConsumer 的消费重试次数和重试间隔主要由服务端消费组配置控制,关键配置包括:
配置 | 含义 |
retryMaxTimes | 最大重试次数 |
groupRetryPolicy | 消费失败后的退避策略 |
Proxy 会将消费组配置下发给 gRPC 客户端,转换关系为
maxAttempts = retryMaxTimes + 1。消费组默认 retryMaxTimes = 16,因此 gRPC PushConsumer 默认首次消费 1 次 + 最多重试 16 次 = 最多尝试 17 次。普通消息
PushConsumer consumer = provider.newPushConsumerBuilder().setClientConfiguration(clientConfiguration).setConsumerGroup("GID_xxx").setSubscriptionExpressions(subscriptionExpressions).setMessageListener(messageView -> {try {// 业务处理return ConsumeResult.SUCCESS;} catch (Throwable t) {return ConsumeResult.FAILURE;}}).build();
消费成功:业务返回
ConsumeResult.SUCCESS,SDK 自动发送 ACK。消费失败:业务返回
ConsumeResult.FAILURE,SDK 根据消费组 groupRetryPolicy 计算下一次重试延迟,调用 ChangeInvisibleDuration 修改消息不可见时间,到期后重新投递;达到 retryMaxTimes 后进入 DLQ。控制方式:
控制项 | 控制位置 |
最大重试次数 | 消费组 retryMaxTimes |
最大总尝试次数 | retryMaxTimes + 1 |
重试间隔 | 消费组 groupRetryPolicy |
业务是否成功 | listener 返回 SUCCESS / FAILURE |
用户通常不需要在 gRPC PushConsumer 代码里控制重试次数。
顺序消息
gRPC 顺序 PushConsumer 不是无限重试,同样受消费组配置控制:最大重试次数 =
retryMaxTimes,最大总尝试次数 = retryMaxTimes + 1,默认最多重试 16 次、最多尝试 17 次。顺序 消息失败后的处理路径和普通消息不同:普通消息失败后主要通过修改不可见时间等待服务端重新投递;顺序消息失败后,SDK 会优先在客户端本地延迟后重新调用 listener。达到
maxAttempts 后,SDK 调用 ForwardMessageToDeadLetterQueue,消息进入 DLQ。可理解为有限次数、本地串行延迟重试。顺序消息的重试间隔来自消费组
groupRetryPolicy。若使用默认自定义退避策略,第一次失败后通常约 10 秒后重试,后续按策略递增。gRPC Java SDK SimpleConsumer
SimpleConsumer 是用户主动拉取、主动 ACK 的消费模型,与 PushConsumer 最大区别是:PushConsumer 由 SDK 根据 listener 返回值自动处理 ACK、重试和 DLQ;SimpleConsumer 必须用户自己决定 ACK、不 ACK 或修改不可见时间。
List<MessageView> messages = simpleConsumer.receive(16, Duration.ofSeconds(30));for (MessageView message : messages) {try {// 业务处理simpleConsumer.ack(message);} catch (Throwable t) {// 用户自行决定:不 ack、修改不可见时间,或者 ack 丢弃}}
最大重试次数:仍受消费组
retryMaxTimes 控制,默认 16 次重试、最多 17 次投递,达到最大次数后进入 DLQ,并非无限重投。重试间隔:每次失败后多久重新可见主要由用户控制:
方式一:接收消息(receive)时传入
invisibleDuration,如 simpleConsumer.receive(16, Duration.ofSeconds(30)),未 ACK 则约 30 秒后重新可见。方式二:主动修改不可见时间
simpleConsumer.changeInvisibleDuration(message, Duration.ofMinutes(5)),约 5 分钟后重新可见。消费组 RetryPolicy 对 SimpleConsumer 的影响:
配置项 | 是否影响 SimpleConsumer | 说明 |
消费组 retryMaxTimes | 是 | 控制最终最多重投多少次,超过后进入 DLQ |
消费组 groupRetryPolicy | 不直接托管业务失败后的每次间隔 | 由用户通过不可见时间控制下一次可见时间 |
receive 的 invisibleDuration | 是 | 控制本次拉取后未 ACK 消息多久重新可见 |
changeInvisibleDuration | 是 | 主动控制单条消息下一次可见时间 |
4.x SDK / 5.x SDK
5.x SDK:无需特别处理,5.0 的 SDK 遵循上文各小节所述的重试规则。
4.x SDK:如果用户需要自行调整重试次数,可通过设置 consumer 的参数决定:
pushConsumer.setMaxReconsumeTimes(3);
消费组重试策略(gRPC)
默认自定义退避
消费组默认重试策略为自定义退避
CUSTOMIZED,默认间隔表为:1s, 5s, 10s, 30s, 1m, 2m, 3m, 4m, 5m, 6m, 7m, 8m, 9m, 10m, 20m, 30m, 1h, 2h
服务端兼容旧延迟级别时,映射为
index = reconsumeTimes + 2。若当前 reconsumeTimes = 0,第一次失败后的延迟通常映射到 10s,后续大致为 10s, 30s, 1m, 2m … 超过表长度后使用最后一个间隔 2h。指数退避
若消费组配置为指数退避
EXPONENTIAL,默认参数为 initial = 5s、max = 2h、multiplier = 2,计算方式为 delay = min(initial * multiplier^reconsumeTimes, max),典型间隔为 5s, 10s, 20s, 40s, 80s …,最大不超过 2h。配置建议
Remoting PushConsumer 用户:建议显式设置最大重试次数
consumer.setMaxReconsumeTimes(16);,必要时设置失败后的延迟级别 context.setDelayLevelWhenNextConsume(3); return ConsumeConcurrentlyStatus.RECONSUME_LATER;。gRPC PushConsumer 用户:建议通过消费组配置
retryMaxTimes、groupRetryPolicy 控制;用户代码只需根据业务结果返回 ConsumeResult.SUCCESS 或 ConsumeResult.FAILURE。gRPC SimpleConsumer 用户:需同时关注两层控制——消费组
retryMaxTimes(最多重试多少次)与 invisibleDuration / changeInvisibleDuration(每次失败后多久重新可见)。建议 invisibleDuration 应大于业务最大处理时间 + 网络抖动时间;若业务处理时间可能超过不可见时间,及时调用 simpleConsumer.changeInvisibleDuration(message, newDuration),避免消息尚未处理完成就被重新投递。常见问题
Q1:gRPC 顺序消息 PushConsumer 是无限重试吗?
不是。gRPC 顺序 PushConsumer 的最大重试次数由消费组
retryMaxTimes 控制,默认 16 次重试,即最多 17 次总尝试。达到最大次数后,消息进入 DLQ。Q2: gRPC 顺序消息 PushConsumer 的重试间隔在哪里配置?
由消费组
groupRetryPolicy 控制。默认是自定义退避策略,第一次失败后通常约 10 秒后重试,后续按策略递增,最长通常为 2 小时。Q3: gRPC SimpleConsumer 会自动按照消费组 RetryPolicy 重试吗?
不会完全自动托管。SimpleConsumer 的最大次数受消费组
retryMaxTimes 控制,但每次失败后多久重新可见主要由用户控制:receive 传入的 invisibleDuration,或主动调用 changeInvisibleDuration。Q4: Remoting PushConsumer 是否受消费组 retryMaxTimes 控制?
Remoting PushConsumer 建议以客户端
setMaxReconsumeTimes 为主要控制方式。Q4: 为什么必须做幂等?
无论使用哪种消费方式,都可能发生重复投递,例如:业务处理成功但 ACK 失败、客户端进程崩溃、网络超时、不可见时间设置过短、消费者重平衡。因此业务需要基于业务唯一键、订单号、流水号或消息 ID 做幂等处理。