这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 消费者决定消息最终是否真正被业务处理。核心原则:业务成功后确认,失败有限重试,异常进入死信,消费幂等,处理有超时,指标可观测。
5.1 初始化消费者
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("order-consumer-group");
consumer.setNamesrvAddr("rocketmq-namesrv:9876");
consumer.subscribe("OrderTopic", "OrderCreated || OrderPaid");
consumer.setConsumeThreadMin(4);
consumer.setConsumeThreadMax(8);
consumer.setConsumeTimeout(15);
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
handle(msg);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
关闭:
consumer.shutdown();
消费组命名应包含:
业务域-服务-用途
order-inventory-consumer
5.2 集群消费与广播消费
| 模式 | 行为 |
|---|---|
| CLUSTERING | 同组分摊队列 |
| BROADCASTING | 每个实例消费全量 |
集群消费:
- 适合业务事件;
- 位点保存在 Broker;
- 支持重试和死信;
- 实例扩容提升并行度。
广播消费:
- 适合本地缓存刷新;
- 每个实例独立进度;
- 新实例可能从最新或配置位点开始;
- 通常不使用重试队列语义;
- 消息丢失影响本地副本。
5.3 Push 与 Pull
RocketMQ Java 客户端常见 Push 消费本质是长轮询加内部拉取。
Push:
- API 简单;
- 自动管理拉取和位点;
- 常用于业务消费者。
Pull:
- 自己控制拉取节奏;
- 适合批处理和特殊回放;
- 位点管理更复杂;
- 需要处理队列分配和异常。
一般业务优先 Push,只有特殊回放、大批量迁移才考虑 Pull。
5.4 消费监听器
并发消费:
(MessageListenerConcurrently) (msgs, context) -> {
try {
handle(msgs);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (BusinessRejectException e) {
log.warn("business reject, key={}", e.getKey());
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (Exception e) {
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
顺序消费:
(MessageListenerOrderly) (msgs, context) -> {
try {
handle(msgs);
return ConsumeOrderlyStatus.SUCCESS;
} catch (Exception e) {
return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT;
}
}
顺序消费挂起会阻塞当前队列,必须设置重试上限和告警。
5.5 消费结果处理
| 场景 | 返回 |
|---|---|
| 处理成功 | CONSUME_SUCCESS |
| 业务明确拒绝但已记录 | CONSUME_SUCCESS |
| 下游临时异常 | RECONSUME_LATER |
| 参数错误 | 建议记录死信并成功确认 |
| 未知异常 | 有限重试 |
不要无条件返回 RECONSUME_LATER。不可修复的消息会反复进入重试,最终仍进入死信,期间浪费资源并制造噪声。
5.6 幂等消费
事件:
public record OrderCreatedEvent(
String eventId,
String orderNo,
String userId,
int version) {
}
处理:
@Transactional
public void handle(OrderCreatedEvent event) {
if (consumedEventRepository.exists(event.eventId())) {
return;
}
inventoryService.reserve(event.orderNo());
consumedEventRepository.save(event.eventId());
}
兜底:
- 业务表唯一约束;
- 状态机条件更新;
- 乐观锁;
- Redis 预判;
- 对账。
Redis 只能优化,不能作为唯一幂等依据。
5.7 消费线程与批量
相关参数:
| 参数 | 说明 |
|---|---|
| consumeThreadMin | 最小消费线程 |
| consumeThreadMax | 最大消费线程 |
| consumeMessageBatchMaxSize | 单次处理条数 |
| pullBatchSize | 拉取批次 |
| consumeTimeout | 消费超时 |
设置原则:
- 看下游可承受并发;
- 看数据库连接池;
- 看每条消息耗时;
- 看是否破坏顺序;
- 看消费实例数。
盲目增加线程可能压垮数据库或下游服务。
5.8 队列分配与扩容
Topic 8 queues
Consumer A -> 4 queues
Consumer B -> 4 queues
新增 Consumer C
A -> 3
B -> 3
C -> 2
限制:
- 有效消费者并行度不超过队列数;
- 顺序消息扩容要考虑 key 分散变化;
- 重平衡期间可能有短暂抖动;
- 消费者实例数应与队列数匹配;
- 扩队列后需重新分配。
5.9 消费监控
指标:
consumer_lag
consumer_latency
consumer_success_total
consumer_failure_total
consumer_retry_total
dead_letter_total
message_age
consume_thread_active
downstream_latency
告警:
- lag 持续增长;
- 消费耗时 P99 超标;
- 重试增长;
- 死信增长;
- 消息年龄超过业务容忍;
- 消费线程持续满。
处理积压顺序:
定位慢原因
-> 先降级非关键逻辑
-> 提升下游容量
-> 再增加消费者
-> 最后评估扩队列
5.10 常见问题
| 问题 | 原因 |
|---|---|
| 消息不消费 | Tag、位点、权限、队列不可读 |
| 重复消费 | 至少一次、超时重试 |
| 消费积压 | 处理慢、下游故障、实例不足 |
| 顺序乱 | 使用并发监听器或扩缩队列 |
| 消费卡住 | 顺序队列挂起、下游超时 |
| 死信增长 | 持续异常未处理 |
| 位点异常 | 手工重置、客户端版本或权限问题 |
本章小结
消费者要区分集群与广播、并发与顺序,业务成功后确认,失败按可恢复性重试,不可恢复异常记录后进入死信治理。所有消费逻辑必须幂等,线程数要匹配下游容量。消费 lag、耗时、重试、死信和消息年龄是核心监控指标。
思考题
- Push 消费为什么仍称为拉取模型?
- 参数错误消息为什么不应无限重试?
- 顺序消费卡住时如何避免阻塞所有队列?
- 消费幂等为什么不能只靠 Redis?
- 消费积压时为什么不能先盲目扩线程?