第 16 章:最终一致性
学习目标
- 理解最终一致性的核心思想
- 掌握消息队列实现最终一致
- 学会对账机制与幂等设计
- 避免消息丢失、重复消费的坑
一、为什么选最终一致?
text
强一致事务(2PC/XA)
├── 性能差(TPS 百级)
├── 锁资源,影响可用
└── 协调者单点
最终一致(消息)
├── 性能高(TPS 万级)
├── 不锁资源
└── 容忍短暂不一致java
// 业务场景:下单后发积分
// 强一致:实时调用积分服务,失败回滚订单
// 最终一致:下单成功 + 发消息,积分服务异步加,中间有 1 秒延迟二、消息队列方案
java
// 1. 业务 + 消息同事务(本地消息表)
@Transactional
public void placeOrder(OrderDto dto) {
// 1. 业务操作
orderDao.insert(dto);
// 2. 同一事务插入消息表
outboxDao.insert(new OutboxMessage(
"order.placed",
JsonUtil.toJson(dto),
LocalDateTime.now()
));
}
// 2. 定时任务发消息
@Scheduled(fixedDelay = 1000)
public void publishMessages() {
List<OutboxMessage> pending = outboxDao.findByStatus("PENDING");
for (OutboxMessage msg : pending) {
try {
rocketMQ.send(msg.getTopic(), msg.getPayload());
outboxDao.updateStatus(msg.getId(), "SENT");
} catch (Exception e) {
// 重试
}
}
}java
// 3. 消费者
@RocketMQMessageListener(topic = "order.placed", consumerGroup = "points-group")
public class OrderPlacedListener implements RocketMQListener<OrderEvent> {
@Override
public void onMessage(OrderEvent event) {
// 加积分(可能重复消息,要做幂等)
pointsService.addPoints(event.getUserId(), event.getAmount());
}
}⚠️ 坑 1:消息表 + 业务表不同事务,先发消息还是先业务?先业务后消息,消息丢了还有对账兜底。
三、可靠消息投递
RocketMQ / Kafka 都有"事务消息"功能:
java
// RocketMQ 事务消息
public void placeOrder(OrderDto dto) {
Message msg = new Message("order.placed", JsonUtil.toJson(dto).getBytes());
rocketMQTemplate.sendMessageInTransaction(
"order-tx-group",
msg,
dto // 业务参数
);
}
// 事务回调
@RocketMQTransactionListener
public class OrderTxListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 执行业务
orderService.create((OrderDto) arg);
return RocketMQLocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 反查:检查业务是否成功
OrderDto dto = parseMessage(msg);
return orderService.exists(dto.id)
? RocketMQLocalTransactionState.COMMIT_MESSAGE
: RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
}
}四、幂等设计
消息可能重复,消费必须幂等:
java
// 方案 1:唯一索引
@Insert("INSERT INTO points_log(user_id, points, biz_id) VALUES(...)")
int insertPointsLog(PointsLog log);
// 重复插入报 DuplicateKeyException,直接跳过java
// 方案 2:Redis 幂等
public void onMessage(OrderEvent event) {
String dedupeKey = "points:order:" + event.getOrderId();
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(dedupeKey, "1", Duration.ofDays(7));
if (Boolean.FALSE.equals(first)) {
log.info("重复消息,已跳过:{}", event.getOrderId());
return;
}
pointsService.addPoints(event.getUserId(), event.getAmount());
}java
// 方案 3:版本号
public class Order {
private Long version; // 乐观锁
}
// 消费时检查版本
if (event.getVersion() > currentVersion) {
apply(event);
}⚠️ 坑 2:幂等键设置太短(分钟级),过期后重复消息又处理。幂等键有效期 > 业务查重周期。
五、对账机制
即使消息可靠,也要有对账:
java
// 每日 0:00 对账
@Scheduled(cron = "0 0 0 * * ?")
public void dailyReconciliation() {
// 1. 查订单表所有"已支付"订单
List<Order> paidOrders = orderDao.findByDateAndStatus(yesterday, "PAID");
// 2. 查积分表
Map<Long, BigDecimal> pointsByUser = pointsDao.findByDate(yesterday);
// 3. 比对
for (Order order : paidOrders) {
BigDecimal expected = calculatePoints(order);
BigDecimal actual = pointsByUser.getOrDefault(order.buyerId(), BigDecimal.ZERO);
if (expected.compareTo(actual) != 0) {
// 不一致,补发
pointsService.addPoints(order.buyerId(), expected.subtract(actual));
alert.send("对账补发:user=" + order.buyerId());
}
}
}java
// 对账结果写日志
public class ReconciliationLog {
private LocalDate date;
private int totalOrders;
private int mismatchedOrders;
private BigDecimal diffAmount;
}六、补偿任务
java
// 扫描超时未处理的订单
@Scheduled(fixedDelay = 60000)
public void scanTimeoutOrders() {
Instant deadline = Instant.now().minus(Duration.ofMinutes(30));
List<Order> pending = orderDao.findPendingOrders(deadline);
for (Order order : pending) {
// 1. 检查下游服务是否处理
if (pointsService.exists(order.id())) {
orderDao.updateStatus(order.id, "COMPLETED");
} else {
// 2. 重发消息
rocketMQ.send("order.placed", order);
}
}
}⚠️ 坑 3:补偿任务没有幂等,业务异常时反复触发,数据库被刷爆。补偿也要幂等,最坏情况是"补漏"而不是"补重复"。
七、消息顺序性
java
// 同一订单的消息要按顺序消费
// RocketMQ:Orderly 消费模式
@RocketMQMessageListener(
topic = "order.placed",
consumerGroup = "points-group",
consumeMode = ConsumeMode.ORDERLY // 顺序消费
)
public class OrderPlacedListener implements RocketMQListener<OrderEvent> {
@Override
public void onMessage(OrderEvent event) {
// 同一订单 id 的消息在同一个 queue,顺序处理
pointsService.addPoints(event.getUserId(), event.getAmount());
}
}java
// 投递:相同 orderId 路由到同一 queue
public MessageQueueOrderlySelector implements MessageQueueSelector {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Long orderId = (Long) arg;
int index = (int) (orderId % mqs.size());
return mqs.get(index);
}
}本章小结
| 方案 | 适用 |
|---|---|
| 本地消息表 | 简单可靠,RocketMQ 标配 |
| 事务消息 | 不引入额外表 |
| 幂等 | 必做,任意方案 |
| 对账 | 兜底,任何系统都需要 |
| 补偿 | 异常自动修复 |
动手练习
- 本地消息表:实现订单 + outbox 表,定时发送
- 幂等保证:用 Redis SETNX 防止重复消费
- RocketMQ 事务消息:用
sendMessageInTransaction改造下单 - 对账脚本:写一个对账任务,补发积分差异
下一章:第 17 章:消息可靠性 →