这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章以“电商订单事件链路”为例,把生产者、消费者、幂等、顺序、延迟、事务、监控和部署串联成一个可落地方案。示例可运行性优先,但连接地址、Topic 和数据库配置需要按环境替换。
28.1 业务场景
流程:
create order
-> publish OrderCreated
-> reserve inventory
-> wait payment
-> publish OrderPaid
-> shipping service
-> timeout 30m
-> close unpaid order
需求:
- 创建订单和事件不丢;
- 库存不重复扣减;
- 同一订单状态顺序处理;
- 超时可靠关闭;
- 全链路可追踪;
- 死信可治理;
- 核心链路可降级。
28.2 Topic 设计
| Topic | 类型 | 生产者 | 消费者 |
|---|---|---|---|
| order.order-created.v1 | 事务顺序 | order-api | inventory |
| order.order-paid.v1 | 事务顺序 | order-api | shipping |
| order.order-timeout.v1 | 延迟 | order-api | order-timeout-worker |
| order.governance.v1 | 普通 | consumer | ops 服务 |
命名建议:业务域 + 事件名 + 版本。
28.3 事件契约
OrderCreated:
{
"eventId": "01J8ZP7Q8M2A",
"eventType": "OrderCreated",
"version": 1,
"occurredAt": "2026-08-25T10:00:00+08:00",
"data": {
"orderNo": "O202608250001",
"userId": "U10001",
"items": [
{
"skuId": "SKU1001",
"quantity": 2,
"price": 9900
}
]
}
}
字段规范:
eventId全局唯一;version明确契约版本;occurredAt带时区;- 金额使用最小单位;
- 不放密码、token 和完整敏感信息。
28.4 数据表设计
订单表:
CREATE TABLE orders (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
order_no VARCHAR(64) NOT NULL,
user_id VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL,
amount BIGINT NOT NULL,
version INT NOT NULL DEFAULT 0,
created_at DATETIME(6) NOT NULL,
updated_at DATETIME(6) NOT NULL,
UNIQUE KEY uk_order_no (order_no),
KEY idx_user_time (user_id, created_at)
);
事件表:
CREATE TABLE outbox_event (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
event_id VARCHAR(64) NOT NULL,
transaction_id VARCHAR(64) NOT NULL,
event_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL,
payload JSON NOT NULL,
created_at DATETIME(6) NOT NULL,
updated_at DATETIME(6) NOT NULL,
UNIQUE KEY uk_event_id (event_id),
UNIQUE KEY uk_transaction_id (transaction_id),
KEY idx_status_time (status, created_at)
);
库存变更表:
CREATE TABLE inventory_change (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
change_id VARCHAR(64) NOT NULL,
sku_id VARCHAR(64) NOT NULL,
order_no VARCHAR(64) NOT NULL,
quantity INT NOT NULL,
created_at DATETIME(6) NOT NULL,
UNIQUE KEY uk_change_id (change_id),
KEY idx_sku_time (sku_id, created_at)
);
28.5 生产者封装
@Service
public class OrderEventProducer implements AutoCloseable {
private final TransactionMQProducer producer;
public OrderEventProducer(OrderTransactionListener listener) throws MQClientException {
producer = new TransactionMQProducer("order-producer-group");
producer.setNamesrvAddr("rocketmq-namesrv:9876");
producer.setTransactionListener(listener);
producer.setSendMsgTimeout(3000);
producer.start();
}
public void sendOrderCreated(Order order, OrderCreatedEvent event) throws Exception {
Message message = new Message(
"order.order-created.v1",
"OrderCreated",
Json.write(event));
message.setKeys(order.orderNo());
producer.sendMessageInTransaction(message, order.orderNo());
}
@Override
public void close() {
producer.shutdown();
}
}
顺序发送时使用稳定选择器:
public MessageQueue orderQueueSelector(List<MessageQueue> queues, Object arg) {
int index = Math.floorMod(String.valueOf(arg).hashCode(), queues.size());
return queues.get(index);
}
28.6 库存消费者
@Component
public class InventoryOrderConsumer {
private final InventoryReserveService reserveService;
@PostConstruct
public void start() throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("inventory-consumer");
consumer.setNamesrvAddr("rocketmq-namesrv:9876");
consumer.subscribe("order.order-created.v1", "OrderCreated");
consumer.setConsumeThreadMin(4);
consumer.setConsumeThreadMax(8);
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
for (MessageExt msg : msgs) {
OrderCreatedEvent event = Json.read(msg.getBody(), OrderCreatedEvent.class);
reserveService.reserve(event);
}
return ConsumeOrderlyStatus.SUCCESS;
});
consumer.start();
}
}
库存处理:
INSERT INTO inventory_change(change_id, sku_id, order_no, quantity, created_at)
VALUES (?, ?, ?, ?, NOW(6));
若唯一键冲突,说明事件重复,记录日志后确认。
28.7 超时关单
发送:
Message message = new Message(
"order.order-timeout.v1",
"OrderTimeout",
Json.write(new OrderTimeoutEvent(orderNo, expireAt)));
message.setKeys(orderNo);
message.setDeliverTimeMs(expireAt.toEpochMilli());
producer.send(message);
消费:
UPDATE orders
SET status = 'CLOSED', close_reason = 'TIMEOUT', version = version + 1
WHERE order_no = ? AND status = 'WAIT_PAY';
更新行数为 0 表示订单已进入其他状态,记录审计后确认消费。
28.8 监控埋点
核心指标:
order_created_total
order_event_send_success_total
order_event_send_failure_total
inventory_reserve_success_total
inventory_reserve_duplicate_total
order_timeout_fired_total
order_timeout_effective_total
order_event_lag
order_dead_letter_age
日志字段:
traceId, eventId, orderNo, topic, messageId, queueId, offset
告警:
- 创建订单成功但事件发送失败;
- 库存消费 lag 超阈值;
- 超时事件未按时触发;
- 死信年龄超阈值;
- 业务对账差异。
28.9 部署配置
资源建议:
| 组件 | 副本 | 资源 |
|---|---|---|
| order-api | 2+ | 2C4G |
| inventory-consumer | 按队列和吞吐 | 2C4G |
| timeout-worker | 2+ | 1C2G |
| rocketmq cluster | 2 broker 组以上 | 独立磁盘 |
配置中心:
rocketmq:
namesrv: rocketmq-namesrv:9876
producer-group: order-producer-group
topics:
created: order.order-created.v1
paid: order.order-paid.v1
timeout: order.order-timeout.v1
生产环境使用环境变量或密钥服务注入地址和凭据。
28.10 验收测试
用例:
- 创建订单成功且事件可查询;
- 本地事务回滚后事件不投递;
- 半消息回查能返回正确状态;
- 库存重复消费不重复扣减;
- 同一订单事件按状态顺序处理;
- 未支付订单 30 分钟关闭;
- 已支付订单不会被超时关闭;
- 死信重放不会产生重复业务;
- Broker 重启后链路恢复;
- 业务对账无差异。
本章小结
项目实战的关键是让订单数据库、本地事件表、RocketMQ 消息、库存变更表和监控对账形成闭环。生产端用事务消息和 Outbox,消费端用顺序、幂等和状态机,超时用延迟消息加条件更新,治理用死信和指标承接异常。
思考题
- 为什么事件表和订单同事务?
- 库存扣减的幂等键如何选择?
- 超时关单为什么要条件更新?
- 顺序消费失败时如何避免阻塞全部订单?
- 上线前哪些故障必须演练?