这是《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());
}
适用:
- 订单事件;
- 支付事件;
- 关键审计;
- 不允许静默丢失的业务。
必须处理:
- 超时;
- 网络异常;
- Broker 拒绝;
- 消息过大;
- 本地降级策略。
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);
}
});
适合日志、通知、指标类消息。注意:
- 失败不能只打印;
- 回调线程不宜执行重逻辑;
- 需要设置发送缓冲区大小;
- 应用关闭前等待在途请求;
- 高可靠业务仍建议本地 Outbox。
4.5 单向发送
producer.sendOneway(message);
特点:
- 不等待响应;
- 不关心结果;
- 性能高;
- 可能丢失。
适合:
- 运行指标;
- 非关键日志;
- 采样埋点。
订单、支付、库存等关键消息禁止使用单向发送。
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 会选择同一队列。要求:
- Sharding Key 稳定;
- 消费端串行;
- 不随意扩缩队列;
- 生产失败重试不能改变 key;
- 异常消息进入重试后可能破坏严格顺序,需要业务补偿。
RocketMQ 顺序是分区内顺序,不是全局顺序。
4.7 发送重试
同步发送可设置:
producer.setRetryTimesWhenSendFailed(2);
重试考虑:
- Broker 短暂不可用可重试;
- 消息非法不应重试;
- 超时重试可能造成重复;
- 重试增加下游压力;
- 消费端必须幂等。
更可靠的模式:
业务本地事务
-> 写业务表
-> 写 outbox
-> 后台投递 RocketMQ
-> 失败重试
4.8 消息大小
默认单条消息大小有限制,具体默认值与版本有关。大消息策略:
- 精简字段;
- 压缩;
- 引用对象存储 key;
- 拆分批次;
- 使用流式通道。
错误:
订单事件携带完整商品详情、用户信息、日志和图片列表
正确:
订单事件只带 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
告警:
- 发送失败率;
- P99 发送耗时;
- 超时次数;
- 本地重试队列增长;
- 关键业务发送为 0。
4.10 常见问题
| 问题 | 排查 |
|---|---|
| No route info | Topic 未创建或路由未同步 |
| Send timeout | 网络、Broker 压力、消息过大 |
| Message too large | 检查大小限制 |
| Topic not exist | autoCreateTopicEnable 和权限 |
| RemotingTooMuchRequest | 发送并发或客户端压力 |
| 发送成功但消费不到 | Tag、队列、消费位点、权限 |
| 重复消息 | 超时重试,需要幂等 |
本章小结
生产者应根据业务可靠性选择同步、异步或单向发送。关键业务必须使用同步或 Outbox,并携带业务 key、事件版本和 traceId。顺序消息通过 Sharding Key 绑定队列,但消费端也必须串行。所有至少一次语义下的重复都可能发生,幂等是消费端必做工作。
思考题
- 哪些消息可以使用 oneway?
- 同步发送超时后消息是否一定失败?
- 顺序消息为什么不是全局顺序?
- 如何设计生产端失败兜底?
- 生产者需要记录哪些日志字段?