RocketMQNotes

第 28 章:项目实战

zjc 于 2026-01-28 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章以“电商订单事件链路”为例,把生产者、消费者、幂等、顺序、延迟、事务、监控和部署串联成一个可落地方案。示例可运行性优先,但连接地址、Topic 和数据库配置需要按环境替换。

28.1 业务场景

流程:

create order
  -> publish OrderCreated
  -> reserve inventory
  -> wait payment
  -> publish OrderPaid
  -> shipping service
  -> timeout 30m
  -> close unpaid order

需求:

  1. 创建订单和事件不丢;
  2. 库存不重复扣减;
  3. 同一订单状态顺序处理;
  4. 超时可靠关闭;
  5. 全链路可追踪;
  6. 死信可治理;
  7. 核心链路可降级。

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
      }
    ]
  }
}

字段规范:

  1. eventId 全局唯一;
  2. version 明确契约版本;
  3. occurredAt 带时区;
  4. 金额使用最小单位;
  5. 不放密码、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

告警:

  1. 创建订单成功但事件发送失败;
  2. 库存消费 lag 超阈值;
  3. 超时事件未按时触发;
  4. 死信年龄超阈值;
  5. 业务对账差异。

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 验收测试

用例:

  1. 创建订单成功且事件可查询;
  2. 本地事务回滚后事件不投递;
  3. 半消息回查能返回正确状态;
  4. 库存重复消费不重复扣减;
  5. 同一订单事件按状态顺序处理;
  6. 未支付订单 30 分钟关闭;
  7. 已支付订单不会被超时关闭;
  8. 死信重放不会产生重复业务;
  9. Broker 重启后链路恢复;
  10. 业务对账无差异。

本章小结

项目实战的关键是让订单数据库、本地事件表、RocketMQ 消息、库存变更表和监控对账形成闭环。生产端用事务消息和 Outbox,消费端用顺序、幂等和状态机,超时用延迟消息加条件更新,治理用死信和指标承接异常。

思考题

  1. 为什么事件表和订单同事务?
  2. 库存扣减的幂等键如何选择?
  3. 超时关单为什么要条件更新?
  4. 顺序消费失败时如何避免阻塞全部订单?
  5. 上线前哪些故障必须演练?