Skip to content
第 17 章 架构 ⏱ 11 分钟阅读

第 17 章:消息可靠性 ​

学习目标 ​

  • 掌握消息丢失的三大场景
  • 学会 RocketMQ/Kafka 的可靠性配置
  • 实现生产者、Broker、消费者的全链路保障
  • 避免消息重复、消息乱序的坑

一、消息丢失的三大场景 ​

text
生产端
  └── 消息没发出去

Broker
  └── 消息没存下来

消费端
  └── 消息没处理完就丢了
java
// 场景演示
producer.send(msg);    // 返回成功 = 已发到 Broker?
consumer.receive(msg);  // 拉到的 = 已处理?

二、生产端可靠性 ​

java
// RocketMQ 生产端
DefaultMQProducer producer = new DefaultMQProducer("order-producer");
producer.setNamesrvAddr("rocketmq:9876");
producer.start();

// 同步发送:等 Broker 确认
SendResult result = producer.send(msg);
if (result.getSendStatus() == SendStatus.SEND_OK) {
    // 已成功
}

// 异步发送:回调里判断
producer.send(msg, new SendCallback() {
    @Override
    public void onSuccess(SendResult sendResult) { ... }
    @Override
    public void onException(Throwable e) {
        // 重试
    }
});
java
// RocketMQ Spring Boot
rocketmq:
  producer:
    group: order-producer
    send-message-timeout: 10000
    retry-times-when-send-failed: 3          // 失败重试 3 次
    retry-times-when-send-async-failed: 3
java
// 发送失败,本地重试 + 落库
public void sendWithFallback(Message msg) {
    try {
        rocketMQTemplate.syncSend("order-topic", msg);
    } catch (Exception e) {
        // 写入本地重试表
        retryDao.insert(new RetryMessage(msg));
    }
}

⚠️ 坑 1:send-message-timeout 默认 3s,网络抖动时大量消息失败。调整到 10s + 失败重试。

三、Broker 端可靠性 ​

yaml
# RocketMQ Broker 配置
brokerRole: SYNC_MASTER          # 主从同步
flushDiskType: SYNC_FLUSH         # 同步刷盘(默认异步)
text
ASYNC_FLUSH(异步刷盘)
  └── 性能高,宕机丢消息
SYNC_FLUSH(同步刷盘)
  └── 性能略低,宕机不丢
java
// Kafka 配置
acks=all                          // 所有副本确认
replication.factor=3              // 副本数 3
min.insync.replicas=2              // 最少 2 个同步副本
enable.idempotence=true            // 幂等,防止重复
java
// 生产端配置
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);

⚠️ 坑 2:acks=1 性能高但只等主节点写完,主节点挂了丢消息。金融场景必须 acks=all。

四、消费端可靠性 ​

java
// 手动提交 ACK
@RocketMQMessageListener(
    topic = "order-topic",
    consumerGroup = "order-consumer",
    consumeMode = ConsumeMode.CONCURRENTLY,
    messageModel = MessageModel.CLUSTERING
)
public class OrderConsumer implements RocketMQListener<OrderMessage> {
    @Override
    public void onMessage(OrderMessage msg) {
        try {
            // 业务逻辑
            orderService.handle(msg);
        } catch (Exception e) {
            // 抛异常,RocketMQ 会重试
            throw new RuntimeException(e);
        }
    }
}
java
// 重试配置
rocketmq:
  consumer:
    group: order-consumer
    max-retry-times: 5             // 最多重试 5 次
    retry-delay-time: 1000         // 1s 后重试
    message-model: CLUSTERING
java
// 重试 5 次后进死信队列
public class DeadLetterConsumer implements RocketMQListener<OrderMessage> {
    @Override
    public void onMessage(OrderMessage msg) {
        // 人工处理
        alertService.send("订单死信:" + msg.getOrderId());
    }
}

五、消息幂等 ​

java
// 1. 业务唯一键
@Insert("INSERT INTO order_log(order_id, status, created_at) VALUES(#{orderId}, #{status}, NOW())")
int insertOrderLog(OrderLog log);

// 重复消息 → DuplicateKeyException
java
// 2. Redis 去重
public void onMessage(OrderMessage msg) {
    String dedupeKey = "msg:order:" + msg.getOrderId();
    Boolean first = redis.opsForValue().setIfAbsent(dedupeKey, "1", Duration.ofDays(7));
    if (Boolean.FALSE.equals(first)) {
        return;   // 已处理
    }
    orderService.handle(msg);
}
java
// 3. 状态机:每步只能前进一步
public class OrderStatusMachine {
    private static final Map<OrderStatus, Set<OrderStatus>> transitions = Map.of(
        CREATED, EnumSet.of(PAID, CANCELLED),
        PAID, EnumSet.of(SHIPPED, REFUNDED),
        SHIPPED, EnumSet.of(COMPLETED),
        CANCELLED, EnumSet.noneOf(OrderStatus.class)
    );

    public void transition(Order order, OrderStatus newStatus) {
        if (!transitions.get(order.getStatus()).contains(newStatus)) {
            throw new IllegalStateException("非法状态变更");
        }
        order.setStatus(newStatus);
    }
}

⚠️ 坑 3:消费者处理完业务再提交 ACK,失败重启可重试。业务和 ACK 必须在同一逻辑边界。

六、消息顺序性 ​

java
// 同一订单的消息必须顺序处理
// RocketMQ:同一订单路由到同一 queue
@RocketMQMessageListener(
    topic = "order-topic",
    consumeMode = ConsumeMode.ORDERLY   // 顺序消费
)
public class OrderConsumer implements RocketMQListener<OrderMessage> {
    @Override
    public void onMessage(OrderMessage msg) {
        // 单线程处理这条 queue
        orderService.handle(msg);
    }
}
java
// 投递:指定同一个 QueueSelector
rocketMQTemplate.syncSendOrderly(
    "order-topic",
    msg,
    String.valueOf(msg.getOrderId())   // hashKey 决定 queue
);

七、消息延迟 ​

java
// 定时消息(订单 30 分钟未支付自动关闭)
Message msg = new Message("order-topic", body);
msg.setDelayTimeLevel(4);     // 30s,1m,5m,10m,30m,1h,2h
rocketMQ.send(msg);

// 延迟级别
// 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
java
// 任意延迟(RocketMQ 5.x)
Message msg = new Message("order-topic", body);
msg.setDelayTimeSec(1800);     // 30 分钟

八、消息追踪 ​

java
// 生产端 TraceId
Message msg = new Message("order-topic", body);
msg.setKeys("order-" + orderId);
msg.setTraceId(MDC.get("traceId"));

// 消费端日志关联
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "order-consumer")
public class OrderConsumer implements RocketMQListener<OrderMessage> {
    @Override
    public void onMessage(OrderMessage msg) {
        MDC.put("traceId", msg.getTraceId());
        log.info("开始处理订单");
        try {
            orderService.handle(msg);
        } finally {
            MDC.remove("traceId");
        }
    }
}

本章小结 ​

环节措施
生产同步发送 + 失败重试 + 落库
Broker同步刷盘 + 主从同步 + 多副本
消费手动 ACK + 业务幂等 + 死信队列
顺序同一订单路由同 queue
排查TraceId + 链路日志

动手练习 ​

  1. 同步发送:用 producer.send 同步发送,验证失败重试
  2. 同步刷盘:Broker 配 SYNC_FLUSH,观察性能差异
  3. 幂等保证:用 Redis SETNX 实现消息去重
  4. 死信队列:故意让消费者抛异常,观察消息进入死信队列

下一章:第 18 章:链路追踪 →

本站基于 VitePress 构建 · 由 StackHub 团队维护