第 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: 3java
// 发送失败,本地重试 + 落库
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: CLUSTERINGjava
// 重试 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);
// 重复消息 → DuplicateKeyExceptionjava
// 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 2hjava
// 任意延迟(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 + 链路日志 |
动手练习
- 同步发送:用
producer.send同步发送,验证失败重试 - 同步刷盘:Broker 配
SYNC_FLUSH,观察性能差异 - 幂等保证:用 Redis SETNX 实现消息去重
- 死信队列:故意让消费者抛异常,观察消息进入死信队列
下一章:第 18 章:链路追踪 →