消息可靠性终极方案
概述
消息可靠性是指消息从生产 → 存储 → 消费全链路不丢失、不重复。本文总结各消息队列的可靠性机制,给出可落地的全链路保障方案。
全链路可靠性全景
生产者 → Broker → 消费者
│ │ │
├ 重试机制 ├ 持久化 ├ 幂等消费
├ 发送确认 ├ 集群 ├ 手动 ACK
├ 事务消息 ├ 副本 ├ 重试队列
└ 回调校验 └ 刷盘 └ 死信队列
│
├ 消息轨迹(链路追踪)
└ 补偿对账(定时任务)| 阶段 | 风险 | 方案 |
|---|---|---|
| 生产 | 网络超时、Broker 宕机 | 重试机制 + 发送确认 + 事务消息 |
| 存储 | 磁盘损坏、主从切换丢数据 | 同步刷盘 + ISR 副本 + 集群 |
| 消费 | 消费失败、重复消费 | 手动 ACK + 重试/死信 + 幂等 |
| 全链路 | 未知异常 | 消息轨迹 + 补偿对账 |
一、生产者可靠性
1.1 发送确认
yaml
# RocketMQ
rocketmq:
producer:
retry-times-when-send-failed: 3 # 同步发送重试次数
retry-times-when-send-async-failed: 3 # 异步发送重试次数
send-message-timeout: 3000 # 发送超时
# Kafka
spring:
kafka:
producer:
acks: all # 等待所有副本确认
retries: 3
enable-idempotence: true # 幂等生产者java
// RocketMQ 同步发送 + 回调确认
@Slf4j
@Component
public class ReliableProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public SendResult sendWithRetry(String topic, Object message) {
try {
// 同步发送:阻塞等待确认
SendResult result = rocketMQTemplate.syncSend(topic, message);
if (result.getSendStatus() != SendStatus.SEND_OK) {
// 发送失败,可以记录到数据库等待补偿
saveToCompensationTable(topic, message);
log.warn("Send not OK: {}, saving to compensation", result);
}
return result;
} catch (MessagingException e) {
// 异常重试耗尽,保存到补偿表
saveToCompensationTable(topic, message);
log.error("Send failed after retries, saving to compensation", e);
throw new BusinessException("MESSAGE_SEND_FAILED", "消息发送失败");
}
}
// 补偿表
@Transactional
public void saveToCompensationTable(String topic, Object message) {
MessageCompensation record = new MessageCompensation();
record.setTopic(topic);
record.setMessage(JSON.toJSONString(message));
record.setStatus("PENDING");
record.setCreateTime(LocalDateTime.now());
compensationMapper.insert(record);
}
}1.2 事务消息(RocketMQ)
java
@Component
public class OrderTransactionProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
// 事务消息:半消息 → 执行本地事务 → 提交/回滚
public void createOrderWithTransaction(OrderRequest request) {
Message<String> message = MessageBuilder
.withPayload(JSON.toJSONString(request))
.setHeader("order_id", request.getOrderId())
.build();
// 发送事务消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-tx-topic",
message,
request // 传递本地事务需要的参数
);
log.info("Transaction send result: {}", result);
}
}
// 事务监听器
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderMapper orderMapper;
// 执行本地事务
@Override
@Transactional
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
OrderRequest request = (OrderRequest) arg;
// 创建订单
orderMapper.insert(request.toOrder());
// 本地事务成功,提交消息
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("Local transaction failed", e);
// 本地事务失败,回滚消息
return RocketMQLocalTransactionState.ROLLBACK;
}
}
// 回查事务状态(半消息长时间未确认时 Broker 回查)
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderId = (String) msg.getHeaders().get("order_id");
Order order = orderMapper.selectById(orderId);
if (order != null) {
return RocketMQLocalTransactionState.COMMIT;
}
return RocketMQLocalTransactionState.UNKNOWN;
}
}二、Broker 可靠性
2.1 持久化配置
properties
# RocketMQ broker.conf
# 同步刷盘(性能低,可靠性高)
flushDiskType=SYNC_FLUSH
# 异步刷盘(性能高,可靠性中)
# flushDiskType=ASYNC_FLUSH
# 主从同步方式
brokerRole=SYNC_MASTER # 同步双写
# brokerRole=ASYNC_MASTER # 异步复制
# Kafka server.properties
# 最小 ISR 副本数
min.insync.replicas=2
# unclean leader 选举(禁止非 ISR 副本成为 Leader)
unclean.leader.election.enable=false2.2 集群高可用
text
RocketMQ:
┌─────────┐ ┌─────────┐
│ Master1 │◄──►│ Master2 │ ← 多主集群
└────┬────┘ └────┬────┘
│ │
┌────▼────┐ ┌────▼────┐
│ Slave1 │ │ Slave2 │
└─────────┘ └─────────┘
Kafka:
Partition 0: Leader=Broker1, ISR=[Broker1, Broker2, Broker3]
Partition 1: Leader=Broker2, ISR=[Broker1, Broker2, Broker3]
Partition 2: Leader=Broker3, ISR=[Broker1, Broker2, Broker3]
RabbitMQ 镜像队列:
- 主节点宕机,从节点自动晋升
- ha-mode: all(全节点同步)
- ha-sync-mode: automatic三、消费者可靠性
3.1 手动 ACK + 重试
java
@Component
public class ReliableConsumer {
// RocketMQ:重试 16 次后进入死信队列
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "order-group",
consumeMode = ConsumeMode.ORDERLY, // 顺序消费
maxReconsumeTimes = 16 // 最大重试次数
)
public class OrderConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
try {
OrderRequest request = JSON.parseObject(message, OrderRequest.class);
processOrder(request);
} catch (Exception e) {
log.error("Consume failed: {}", e.getMessage());
// 抛出异常 → RocketMQ 自动重试
throw e;
}
}
}
// Kafka:手动 ACK + 错误处理
@KafkaListener(topics = "order-topic", groupId = "order-group")
public void onMessage(ConsumerRecord<String, String> record,
Acknowledgment ack,
@Header(KafkaHeaders.DELIVERY_ATTEMPT) int delivery) {
try {
OrderRequest request = JSON.parseObject(record.value(), OrderRequest.class);
processOrder(request);
ack.acknowledge(); // 手动确认
} catch (Exception e) {
if (delivery >= 3) {
// 重试 3 次仍失败,发送到死信队列
kafkaTemplate.send("order-dlt", record.value());
ack.acknowledge(); // 确认原消息(已投递到死信)
} else {
throw new KafkaException("Retry later", e); // 触发重试
}
}
}
}3.2 消费幂等
java
@Component
public class IdempotentConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
// 基于消息唯一 ID 的幂等方案
public boolean consume(String messageId, Consumer<String> consumer) {
// SETNX:消息 ID 不存在时返回 true(首次消费)
Boolean success = redisTemplate.opsForValue()
.setIfAbsent("msg:" + messageId, "consumed",
Duration.ofHours(24)); // 24h 过期
if (Boolean.TRUE.equals(success)) {
try {
consumer.accept(messageId);
return true;
} catch (Exception e) {
// 消费失败,删除幂等标记(允许重试)
redisTemplate.delete("msg:" + messageId);
throw e;
}
}
// 已消费过,跳过
log.debug("Duplicate message: {}", messageId);
return false;
}
}四、死信队列 & 重试机制
4.1 RocketMQ 重试 & 死信
properties
# RocketMQ 重试队列
%RETRY%order-group # 重试队列(消费者组级别)
# 重试级别(每次重试间隔递增):
# 1:10s, 2:30s, 3:1m, 4:2m, 5:3m, 6:4m
# 7:5m, 8:6m, 9:7m, 10:8m, 11:9m, 12:10m
# 13:20m, 14:30m, 15:1h, 16:2h
# 超过 16 次后进入死信队列
%DLQ%order-group # 死信队列java
// 死信队列消费者(人工介入处理)
@RocketMQMessageListener(
topic = "%DLQ%order-group",
consumerGroup = "order-dlq-group"
)
public class DeadLetterConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
// 死信消息需要人工介入
log.error("Dead letter message: {}", message);
// 保存到数据库,运维人员检查
deadLetterMapper.insert(new DeadLetter(message, LocalDateTime.now()));
// 发送告警
alertService.sendAlert("订单消息进入死信队列");
}
}4.2 Kafka 死信
java
@Configuration
public class KafkaDeadLetterConfig {
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String>
kafkaListenerContainerFactory(ConsumerFactory<String, String> cf) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
factory.setConsumerFactory(cf);
factory.setCommonErrorHandler(new DefaultErrorHandler(
// 重试 3 次,间隔 1s
new FixedBackOff(1000L, 3)
) {
@Override
public void handleRemaining(Exception e, List<ConsumerRecord<?, ?>> records,
Consumer<?, ?> consumer, MessageListenerContainer container) {
// 重试耗尽,发送到死信
for (ConsumerRecord<?, ?> record : records) {
kafkaTemplate.send("order-dlt", record.value().toString());
}
// 确认原消息
ack.acknowledge();
}
});
return factory;
}
}五、消息轨迹
5.1 RocketMQ 消息轨迹
yaml
rocketmq:
producer:
enable-msg-trace: true # 开启生产者轨迹
consumer:
enable-msg-trace: true # 开启消费者轨迹java
// 自定义消息轨迹
@Component
public class MessageTracer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendWithTrace(String topic, Object message) {
// 注入 TraceId
Map<String, Object> headers = new HashMap<>();
headers.put("traceId", MDC.get("traceId"));
Message<?> traceMsg = MessageBuilder
.withPayload(message)
.copyHeaders(headers)
.build();
rocketMQTemplate.syncSend(topic, traceMsg);
// 记录投递轨迹
log.info("Message sent: topic={}, traceId={}, time={}",
topic, headers.get("traceId"), Instant.now());
}
}5.2 消息轨迹查询
sql
-- RocketMQ 消息轨迹表结构
CREATE TABLE msg_trace (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
msg_id VARCHAR(64) NOT NULL, -- 消息 ID
trace_id VARCHAR(32) NOT NULL, -- 链路 TraceId
topic VARCHAR(128) NOT NULL,
producer_host VARCHAR(64), -- 生产者 IP
produce_time DATETIME, -- 生产时间
consumer_host VARCHAR(64), -- 消费者 IP
consume_time DATETIME, -- 消费时间
status VARCHAR(16), -- SUCCESS / FAILED
cost_time BIGINT, -- 耗时
INDEX idx_msg_id (msg_id),
INDEX idx_trace_id (trace_id)
);六、补偿对账
6.1 定时补偿
java
@Component
public class MessageCompensationJob {
@Autowired
private CompensationMapper compensationMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
// 每 5 分钟扫描重发未确认的消息
@Scheduled(fixedDelay = 300000)
@Transactional
public void compensate() {
// 查询 5 分钟前发送且未确认的消息
List<MessageCompensation> pending = compensationMapper.selectPending(
LocalDateTime.now().minusMinutes(5));
for (MessageCompensation record : pending) {
try {
// 重新发送
SendResult result = rocketMQTemplate.syncSend(
record.getTopic(), record.getMessage());
if (result.getSendStatus() == SendStatus.SEND_OK) {
record.setStatus("COMPENSATED");
record.setCompensateTime(LocalDateTime.now());
compensationMapper.updateById(record);
}
} catch (Exception e) {
log.error("Compensation failed: {}", record.getId(), e);
// 超过最大补偿次数 → 人工介入
if (record.getRetryCount() >= 10) {
record.setStatus("MANUAL_REQUIRED");
compensationMapper.updateById(record);
alertService.sendAlert("消息补偿失败,需要人工介入");
} else {
record.setRetryCount(record.getRetryCount() + 1);
compensationMapper.updateById(record);
}
}
}
}
}6.2 对账方案
java
// 消息对账:生产者 vs 消费者
// 场景:订单系统 vs 支付系统对账
@Component
public class MessageReconciliation {
// 每 30 分钟核对一次
@Scheduled(fixedDelay = 1800000)
public void reconcile() {
// 1. 从消息轨迹表查询 30 分钟前的消息
List<MessageTrace> sent = traceMapper.selectByTimeRange(
LocalDateTime.now().minusMinutes(60),
LocalDateTime.now().minusMinutes(30));
// 2. 从业务表查询消费结果
for (MessageTrace trace : sent) {
Order order = orderMapper.selectByMsgId(trace.getMsgId());
if (order == null) {
// 消息丢失!重新发送
log.warn("Message lost, re-sending: msgId={}", trace.getMsgId());
rocketMQTemplate.syncSend(trace.getTopic(), trace.getPayload());
}
}
}
}七、故障注入测试
java
// 使用 Chaos Engineering 思想验证可靠性
// 场景:Broker 宕机、网络分区、磁盘故障
@SpringBootTest
@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
public class ReliabilityTest {
@Autowired
private ReliableProducer producer;
@Test
@Order(1)
void testBrokerDown() {
// 1. 发送消息
producer.sendWithRetry("order-topic", new OrderRequest("1001"));
// 2. 模拟 Broker 宕机(SSH 执行 kill)
// ssh user@broker "kill -9 \$(pidof java)"
// 3. 再发送一条消息
producer.sendWithRetry("order-topic", new OrderRequest("1002"));
// 4. 恢复 Broker
// ssh user@broker "./bin/mqbroker -c broker.conf"
// 5. 验证两条消息都被消费
// assertThat(consumer.getConsumedCount()).isEqualTo(2);
}
@Test
@Order(2)
void testNetworkPartition() {
// 模拟网络隔离:iptables 封禁端口
// iptables -A INPUT -p tcp --dport 9876 -j DROP
// Thread.sleep(5000)
// iptables -D INPUT -p tcp --dport 9876 -j DROP
// 验证:网络恢复后消息自动补发
}
}八、总结
| 阶段 | 风险点 | 最佳方案 |
|---|---|---|
| 生产 | 发送失败 | 同步发送 + 重试 + 事务消息 + 补偿表 |
| 生产 | 消息丢失 | acks=all + 同步刷盘 |
| Broker | 主从切换 | 同步双写 + ISR >= 2 |
| 消费 | 处理失败 | 手动 ACK + 重试队列 + 死信队列 |
| 消费 | 重复消息 | 幂等表(消息 ID SETNX) |
| 全链路 | 未知问题 | 消息轨迹 + 定时对账 + 补偿 |
| 验证 | 可靠性 | 故障注入测试 |
参考链接: