RocketMQNotes

第 04 章:生产者

zjc 于 2026-01-04 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 生产者的目标不是简单调用 send,而是在消息可靠性、延迟、顺序、幂等和故障降级之间做出明确选择。本章介绍 RocketMQ Java 客户端常用发送方式与生产实践。

4.1 初始化生产者

依赖:

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>5.3.1</version>
</dependency>

传统客户端:

DefaultMQProducer producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr("rocketmq-namesrv:9876");
producer.setSendMsgTimeout(3000);
producer.setRetryTimesWhenSendFailed(2);
producer.start();

关闭:

producer.shutdown();

生产者通常是单例,不要每次请求创建和销毁。

4.2 构建消息

Message message = new Message();
message.setTopic("OrderTopic");
message.setTags("OrderCreated");
message.setKeys(order.orderNo());
message.setBody(JsonUtils.toBytes(event));
message.putUserProperty("traceId", TraceContext.currentId());
message.putUserProperty("eventType", "OrderCreated");
message.putUserProperty("schemaVersion", "2");

字段建议:

字段 用途
Topic 业务域
Tag 事件类型
keys 查询和业务关联
traceId 链路追踪
schemaVersion 契约版本
occurredAt 业务发生时间

消息体应使用稳定契约对象,避免暴露内部实体。

4.3 同步发送

SendResult result = producer.send(message);

if (result.getSendStatus() != SendStatus.SEND_OK) {
    throw new MessageSendException(result.getSendStatus().name());
}

适用:

  1. 订单事件;
  2. 支付事件;
  3. 关键审计;
  4. 不允许静默丢失的业务。

必须处理:

  1. 超时;
  2. 网络异常;
  3. Broker 拒绝;
  4. 消息过大;
  5. 本地降级策略。

4.4 异步发送

producer.send(message, new SendCallback() {
    @Override
    public void onSuccess(SendResult result) {
        log.info("message sent, msgId={}, status={}", result.getMsgId(), result.getSendStatus());
    }

    @Override
    public void onException(Throwable e) {
        log.error("message send failed, orderNo={}", orderNo, e);
        saveToLocalRetryQueue(orderNo);
    }
});

适合日志、通知、指标类消息。注意:

  1. 失败不能只打印;
  2. 回调线程不宜执行重逻辑;
  3. 需要设置发送缓冲区大小;
  4. 应用关闭前等待在途请求;
  5. 高可靠业务仍建议本地 Outbox。

4.5 单向发送

producer.sendOneway(message);

特点:

  1. 不等待响应;
  2. 不关心结果;
  3. 性能高;
  4. 可能丢失。

适合:

  1. 运行指标;
  2. 非关键日志;
  3. 采样埋点。

订单、支付、库存等关键消息禁止使用单向发送。

4.6 顺序发送

SendResult result = producer.send(message, (mqs, msg, arg) -> {
    long index = Math.abs((long) arg) % mqs.size();
    return mqs.get((int) index);
}, order.orderNo());

同一个 orderId 会选择同一队列。要求:

  1. Sharding Key 稳定;
  2. 消费端串行;
  3. 不随意扩缩队列;
  4. 生产失败重试不能改变 key;
  5. 异常消息进入重试后可能破坏严格顺序,需要业务补偿。

RocketMQ 顺序是分区内顺序,不是全局顺序。

4.7 发送重试

同步发送可设置:

producer.setRetryTimesWhenSendFailed(2);

重试考虑:

  1. Broker 短暂不可用可重试;
  2. 消息非法不应重试;
  3. 超时重试可能造成重复;
  4. 重试增加下游压力;
  5. 消费端必须幂等。

更可靠的模式:

业务本地事务
  -> 写业务表
  -> 写 outbox
  -> 后台投递 RocketMQ
  -> 失败重试

4.8 消息大小

默认单条消息大小有限制,具体默认值与版本有关。大消息策略:

  1. 精简字段;
  2. 压缩;
  3. 引用对象存储 key;
  4. 拆分批次;
  5. 使用流式通道。

错误:

订单事件携带完整商品详情、用户信息、日志和图片列表

正确:

订单事件只带 orderId、版本、关键字段
消费者按需查询详情

4.9 生产者监控

指标:

rocketmq_producer_send_total
rocketmq_producer_success_total
rocketmq_producer_failure_total
rocketmq_producer_latency
rocketmq_producer_timeout_total
rocketmq_producer_retry_total
outbox_pending_total

日志必须包含:

topic
tag
keys / orderId
messageId
sendStatus
cost
broker
queueId
traceId

告警:

  1. 发送失败率;
  2. P99 发送耗时;
  3. 超时次数;
  4. 本地重试队列增长;
  5. 关键业务发送为 0。

4.10 常见问题

问题 排查
No route info Topic 未创建或路由未同步
Send timeout 网络、Broker 压力、消息过大
Message too large 检查大小限制
Topic not exist autoCreateTopicEnable 和权限
RemotingTooMuchRequest 发送并发或客户端压力
发送成功但消费不到 Tag、队列、消费位点、权限
重复消息 超时重试,需要幂等

本章小结

生产者应根据业务可靠性选择同步、异步或单向发送。关键业务必须使用同步或 Outbox,并携带业务 key、事件版本和 traceId。顺序消息通过 Sharding Key 绑定队列,但消费端也必须串行。所有至少一次语义下的重复都可能发生,幂等是消费端必做工作。

思考题

  1. 哪些消息可以使用 oneway?
  2. 同步发送超时后消息是否一定失败?
  3. 顺序消息为什么不是全局顺序?
  4. 如何设计生产端失败兜底?
  5. 生产者需要记录哪些日志字段?