RocketMQNotes

第 05 章:消费者

zjc 于 2026-01-05 发布

这是《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 每个实例消费全量

集群消费:

  1. 适合业务事件;
  2. 位点保存在 Broker;
  3. 支持重试和死信;
  4. 实例扩容提升并行度。

广播消费:

  1. 适合本地缓存刷新;
  2. 每个实例独立进度;
  3. 新实例可能从最新或配置位点开始;
  4. 通常不使用重试队列语义;
  5. 消息丢失影响本地副本。

5.3 Push 与 Pull

RocketMQ Java 客户端常见 Push 消费本质是长轮询加内部拉取。

Push:

  1. API 简单;
  2. 自动管理拉取和位点;
  3. 常用于业务消费者。

Pull:

  1. 自己控制拉取节奏;
  2. 适合批处理和特殊回放;
  3. 位点管理更复杂;
  4. 需要处理队列分配和异常。

一般业务优先 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());
}

兜底:

  1. 业务表唯一约束;
  2. 状态机条件更新;
  3. 乐观锁;
  4. Redis 预判;
  5. 对账。

Redis 只能优化,不能作为唯一幂等依据。

5.7 消费线程与批量

相关参数:

参数 说明
consumeThreadMin 最小消费线程
consumeThreadMax 最大消费线程
consumeMessageBatchMaxSize 单次处理条数
pullBatchSize 拉取批次
consumeTimeout 消费超时

设置原则:

  1. 看下游可承受并发;
  2. 看数据库连接池;
  3. 看每条消息耗时;
  4. 看是否破坏顺序;
  5. 看消费实例数。

盲目增加线程可能压垮数据库或下游服务。

5.8 队列分配与扩容

Topic 8 queues
  Consumer A -> 4 queues
  Consumer B -> 4 queues

新增 Consumer C
  A -> 3
  B -> 3
  C -> 2

限制:

  1. 有效消费者并行度不超过队列数;
  2. 顺序消息扩容要考虑 key 分散变化;
  3. 重平衡期间可能有短暂抖动;
  4. 消费者实例数应与队列数匹配;
  5. 扩队列后需重新分配。

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

告警:

  1. lag 持续增长;
  2. 消费耗时 P99 超标;
  3. 重试增长;
  4. 死信增长;
  5. 消息年龄超过业务容忍;
  6. 消费线程持续满。

处理积压顺序:

定位慢原因
  -> 先降级非关键逻辑
  -> 提升下游容量
  -> 再增加消费者
  -> 最后评估扩队列

5.10 常见问题

问题 原因
消息不消费 Tag、位点、权限、队列不可读
重复消费 至少一次、超时重试
消费积压 处理慢、下游故障、实例不足
顺序乱 使用并发监听器或扩缩队列
消费卡住 顺序队列挂起、下游超时
死信增长 持续异常未处理
位点异常 手工重置、客户端版本或权限问题

本章小结

消费者要区分集群与广播、并发与顺序,业务成功后确认,失败按可恢复性重试,不可恢复异常记录后进入死信治理。所有消费逻辑必须幂等,线程数要匹配下游容量。消费 lag、耗时、重试、死信和消息年龄是核心监控指标。

思考题

  1. Push 消费为什么仍称为拉取模型?
  2. 参数错误消息为什么不应无限重试?
  3. 顺序消费卡住时如何避免阻塞所有队列?
  4. 消费幂等为什么不能只靠 Redis?
  5. 消费积压时为什么不能先盲目扩线程?