线上数据不一致排查实战:分布式事务与最终一致性
2026/7/3大约 21 分钟
线上数据不一致排查实战:分布式事务与最终一致性
周五下午 5 点,客服反馈大量用户投诉:明明支付成功了,订单状态却显示"待支付";优惠券已经核销了,但订单里还显示"未使用"。你打开数据库一看,订单库和支付库的数据对不上——在一个微服务架构的系统中,数据不一致是最让人头疼的问题之一。它不像 OOM 或 CPU 飙升那样有明确的告警,而是悄无声息地累积,直到用户投诉才被发现。本文将系统性梳理数据不一致的排查和修复全流程。
一、数据不一致的 5 种典型场景
1.1 场景全景图
┌──────────────────────────────────────────────────────────────────┐
│ 微服务数据不一致的 5 种典型场景 │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 场景1: 跨服务操作部分成功 │ │
│ │ 订单服务 ──▶ 创建订单(✓) │ │
│ │ 库存服务 ──▶ 扣减库存(✗ 网络超时) │ │
│ │ 优惠券 ──▶ 核销券(✓) │ │
│ │ 结果: 订单创建了,券核销了,但库存没扣 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 场景2: 消息消费部分失败 │ │
│ │ Producer ──▶ 发送消息(✓) ──▶ 事务提交(✓) │ │
│ │ Consumer ──▶ 消费消息 ──▶ 更新DB(✗) ──▶ ACK(✓) │ │
│ │ 结果: 消息被确认消费了,但数据库更新失败 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 场景3: 并发更新顺序不一致 │ │
│ │ 线程A: 读旧值 ──▶ 计算 ──▶ 写入 │ │
│ │ 线程B: 读旧值 ──▶ 计算 ──▶ 写入 │ │
│ │ 结果: 后写的覆盖先写的,但逻辑上应该 A 在前 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 场景4: 缓存与数据库不一致 │ │
│ │ DB ──▶ 更新(✓) ──▶ 删除缓存(✗ 网络故障) │ │
│ │ 结果: 缓存还是旧值,读到脏数据 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 场景5: 补偿事务执行失败 │ │
│ │ TCC: Try(✓) ──▶ Confirm(✓) ──▶ Cancel(✗ 节点宕机) │ │
│ │ 结果: Try 成功但 Cancel 没执行,资源锁定 │ │
│ └─────────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘1.2 真实案例:支付成功但订单未更新
// 问题代码:跨服务操作没有事务保障
@Service
public class OrderService {
@Autowired
private PaymentClient paymentClient;
@Autowired
private CouponClient couponClient;
@Autowired
private InventoryClient inventoryClient;
@Transactional // 只保障本地数据库事务
public OrderResult createOrder(OrderRequest request) {
// 1. 创建订单(本地数据库)
Order order = orderMapper.insert(request.toOrder());
// 2. 核销优惠券(远程调用)
couponClient.redeem(request.getCouponId()); // 如果这里超时...
// 3. 扣减库存(远程调用)
inventoryClient.deduct(request.getSkuId(), request.getQuantity()); // ...这里就不执行了
// 4. 发起支付(远程调用)
PaymentResult payment = paymentClient.pay(order.getId(), request.getAmount());
// 如果 2 成功、3 失败、4 未执行:
// 订单已创建,券已核销,库存未扣,未发起支付
// 用户看到:订单存在但无法支付,券也找不回来了
return OrderResult.success(order);
}
}二、排查方法:如何发现数据不一致
2.1 发现途径
┌──────────────────────────────────────────────────────────────────┐
│ 数据不一致的发现途径 │
├──────────────┬──────────────┬───────────────┬───────────────────┤
│ 用户投诉 │ 对账系统 │ 监控告警 │ 人工巡检 │
├──────────────┼──────────────┼───────────────┼───────────────────┤
│ 最晚发现 │ 最可靠 │ 最及时 │ 最耗人力 │
│ 影响最大 │ 定期执行 │ 需要建设 │ 适合抽查 │
│ 已经造成损失 │ 可能有延迟 │ 误报率高 │ 不能全量覆盖 │
└──────────────┴──────────────┴───────────────┴───────────────────┘
优先级:对账系统 > 监控告警 > 用户投诉 > 人工巡检2.2 对账系统设计
// 对账系统核心逻辑
@Service
public class ReconciliationService {
@Scheduled(cron = "0 0 1 * * ?") // 每天凌晨 1 点执行
public void dailyReconciliation() {
LocalDate date = LocalDate.now().minusDays(1); // 对昨天的数据对账
// 1. 拉取各系统的数据
List<OrderRecord> orders = orderService.getByDate(date);
List<PaymentRecord> payments = paymentService.getByDate(date);
List<CouponRecord> coupons = couponService.getByDate(date);
List<InventoryRecord> inventories = inventoryService.getByDate(date);
// 2. 按 orderId 关联对比
Map<String, OrderRecord> orderMap = orders.stream()
.collect(Collectors.toMap(OrderRecord::getOrderId, o -> o));
for (PaymentRecord payment : payments) {
OrderRecord order = orderMap.get(payment.getOrderId());
if (order == null) {
// 支付了但订单不存在
reportInconsistency("ORDER_NOT_FOUND", payment.getOrderId(),
"Payment exists but order not found");
continue;
}
// 3. 逐项对比
if (!"PAID".equals(order.getStatus()) && "SUCCESS".equals(payment.getStatus())) {
reportInconsistency("STATUS_MISMATCH", order.getOrderId(),
String.format("Order status=%s but payment status=%s",
order.getStatus(), payment.getStatus()));
}
if (order.getAmount().compareTo(payment.getAmount()) != 0) {
reportInconsistency("AMOUNT_MISMATCH", order.getOrderId(),
String.format("Order amount=%s but payment amount=%s",
order.getAmount(), payment.getAmount()));
}
}
// 4. 检查订单是否都有对应的支付记录
Set<String> paymentOrderIds = payments.stream()
.map(PaymentRecord::getOrderId)
.collect(Collectors.toSet());
for (OrderRecord order : orders) {
if ("PAID".equals(order.getStatus()) && !paymentOrderIds.contains(order.getOrderId())) {
reportInconsistency("PAYMENT_MISSING", order.getOrderId(),
"Order is PAID but no payment record found");
}
}
// 5. 检查优惠券状态
for (OrderRecord order : orders) {
if (order.getCouponId() != null) {
CouponRecord coupon = couponService.getById(order.getCouponId());
if ("USED".equals(coupon.getStatus()) && !"PAID".equals(order.getStatus())) {
reportInconsistency("COUPON_ORDER_MISMATCH", order.getOrderId(),
String.format("Coupon %s is USED but order status is %s",
coupon.getCouponId(), order.getStatus()));
}
}
}
}
private void reportInconsistency(String type, String orderId, String message) {
// 记录到对账差异表
InconsistencyRecord record = new InconsistencyRecord();
record.setType(type);
record.setOrderId(orderId);
record.setMessage(message);
record.setCreatedAt(LocalDateTime.now());
record.setStatus("PENDING");
inconsistencyMapper.insert(record);
// 发送告警
alertService.send("DATA_INCONSISTENCY",
String.format("[%s] %s: %s", type, orderId, message));
log.error("Data inconsistency detected: type={}, orderId={}, msg={}",
type, orderId, message);
}
}2.3 实时监控:基于 Binlog 的 CDC 对比
// 使用 Canal/Debezium 监听 binlog,实时对比数据
@Component
public class CdcReconciliationListener {
@KafkaListener(topics = "order-db-binlog")
public void onOrderChange(ChangeEvent event) {
if (event.getOperation() != ChangeEvent.Operation.UPDATE) return;
String orderId = event.getTable() + ":" + event.getPrimaryKey();
String orderStatus = event.getAfter().get("status");
// 异步检查支付状态
CompletableFuture.runAsync(() -> {
PaymentRecord payment = paymentService.getByOrderId(orderId);
if (payment == null) return;
if ("PAID".equals(orderStatus) && !"SUCCESS".equals(payment.getStatus())) {
// 订单已支付但支付记录不是成功状态
alertService.send("CDC_INCONSISTENCY",
String.format("Order %s status=PAID but payment status=%s",
orderId, payment.getStatus()));
}
});
}
}2.4 Binlog 对比工具
# 使用 MySQL binlog 解析工具对比数据变化
# 工具 1: mysqlbinlog(MySQL 自带)
mysqlbinlog --start-datetime="2026-07-03 14:00:00" \
--stop-datetime="2026-07-03 15:00:00" \
/var/lib/mysql/mysql-bin.001234 > /tmp/binlog.sql
# 提取特定表的变更
grep -A 5 "UPDATE orders" /tmp/binlog.sql
# 工具 2: Canal(阿里巴巴开源)
# Canal 模拟 MySQL slave 协议,实时解析 binlog
# 部署 Canal server,将变更推送到 Kafka
# 应用消费 Kafka 消息做实时对账
# 工具 3: Debezium(基于 Kafka Connect)
# 配置 Debezium MySQL connector
# 自动将 binlog 变更推送到 Kafka topic
# 下游消费做对账三、消息对账:基于 MQ 的最终一致性
3.1 消息驱动的数据一致性问题
┌──────────────────────────────────────────────────────────────────┐
│ 消息驱动架构的一致性问题 │
│ │
│ Producer MQ Consumer │
│ ┌──────┐ ┌──────┐ ┌──────┐ │
│ │ 订单 │ ──1. 发送 ───▶ │ │ ──2. 投递──▶ │ 支付 │ │
│ │ 服务 │ │ Kafka│ │ 服务 │ │
│ │ │ ◀──3. 确认 ──── │ │ ◀──4. ACK── │ │ │
│ └──────┘ └──────┘ └──────┘ │
│ │
│ 问题点: │
│ 1. 发送成功但本地事务失败 → 消息发出去了但不该有 │
│ 2. 本地事务成功但发送失败 → 该发的消息没发 │
│ 3. 消费成功但处理失败 → ACK 了但数据库没更新 │
│ 4. 消费超时 → 消息重新投递,可能重复消费 │
│ 5. Consumer 宕机 → 消息积压,恢复后大量消费 │
└──────────────────────────────────────────────────────────────────┘3.2 本地消息表方案
// 方案:本地消息表 + 定时补偿
@Service
public class LocalMessageTableService {
@Autowired
private MessageMapper messageMapper;
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
// 第一步:业务操作 + 消息记录在同一个本地事务中
@Transactional
public void createOrder(OrderRequest request) {
// 1. 创建订单
Order order = orderMapper.insert(request.toOrder());
// 2. 在同一事务中写入消息表
MessageRecord message = new MessageRecord();
message.setTopic("order-created");
message.setKey(order.getId());
message.setPayload(JSON.toJSONString(order));
message.setStatus("PENDING"); // 待发送
message.setRetryCount(0);
message.setCreatedAt(LocalDateTime.now());
messageMapper.insert(message);
// 事务提交后,订单和消息表记录同时存在
// 即使后续发送失败,消息也不会丢
}
// 第二步:定时任务扫描未发送的消息
@Scheduled(fixedDelay = 5000) // 每 5 秒扫描一次
public void sendPendingMessages() {
List<MessageRecord> pending = messageMapper.selectByStatus("PENDING", 100);
for (MessageRecord msg : pending) {
try {
// 发送到 Kafka
kafkaTemplate.send(msg.getTopic(), msg.getKey(), msg.getPayload()).get(10, TimeUnit.SECONDS);
// 发送成功,更新状态
msg.setStatus("SENT");
msg.setSentAt(LocalDateTime.now());
messageMapper.updateStatus(msg);
} catch (Exception e) {
log.warn("Failed to send message {}: {}", msg.getId(), e.getMessage());
msg.setRetryCount(msg.getRetryCount() + 1);
if (msg.getRetryCount() >= 10) {
msg.setStatus("FAILED");
alertService.send("MESSAGE_SEND_FAILED",
"Message " + msg.getId() + " failed after 10 retries");
}
messageMapper.updateStatus(msg);
}
}
}
}3.3 消息消费幂等性保障
// 消费端必须保证幂等性,防止重复消费导致数据不一致
@Component
public class OrderMessageConsumer {
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private OrderService orderService;
@KafkaListener(topics = "order-created")
public void handle(OrderMessage message, Acknowledgment ack) {
String messageId = message.getMessageId();
try {
// 1. 幂等检查:Redis SET NX
Boolean processed = redisTemplate.opsForValue()
.setIfAbsent("msg:processed:" + messageId, "1", 24, TimeUnit.HOURS);
if (Boolean.FALSE.equals(processed)) {
// 已处理过,直接 ACK
log.info("Message {} already processed, skip", messageId);
ack.acknowledge();
return;
}
// 2. 执行业务逻辑
orderService.processOrderCreated(message);
// 3. ACK
ack.acknowledge();
} catch (Exception e) {
// 4. 处理失败,删除幂等标记,让消息重试
redisTemplate.delete("msg:processed:" + messageId);
log.error("Failed to process message {}: {}", messageId, e.getMessage());
// 不 ACK,让 Kafka 重新投递
// 注意:如果一直失败,需要死信队列
throw e;
}
}
}
// 更健壮的幂等方案:数据库唯一索引
@Component
public class IdempotentConsumer {
@Transactional
public void handle(OrderMessage message) {
// 利用数据库唯一索引保证幂等
try {
// 插入幂等记录(message_id 有唯一索引)
idempotentMapper.insert(new IdempotentRecord(message.getMessageId()));
} catch (DuplicateKeyException e) {
log.info("Message {} already processed", message.getMessageId());
return; // 已处理,跳过
}
// 执行业务逻辑
orderService.processOrderCreated(message);
}
}四、定时校验:全量对账系统
4.1 对账系统架构
┌──────────────────────────────────────────────────────────────────┐
│ 对账系统架构 │
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ 订单库 │ │ 支付库 │ │ 优惠券库 │ │ 库存库 │ │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ └────┬─────┘ │
│ │ │ │ │ │
│ ▼ ▼ ▼ ▼ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 数据抽取层 (Data Fetcher) │ │
│ │ 分页拉取、增量对比、内存对齐 │ │
│ └────────────────────────┬────────────────────────────────┘ │
│ ▼ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 对比引擎 (Comparison Engine) │ │
│ │ 规则匹配、差异识别、分类标记 │ │
│ └────────────────────────┬────────────────────────────────┘ │
│ ▼ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 差异处理层 (Discrepancy Handler) │ │
│ │ 自动修复、人工审核、告警通知 │ │
│ └─────────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘4.2 对账规则定义
// 对账规则配置
@Configuration
public class ReconciliationRules {
// 规则 1: 订单状态 vs 支付状态
public static final ReconciliationRule ORDER_PAYMENT_RULE =
ReconciliationRule.builder()
.name("ORDER_PAYMENT_STATUS")
.leftSource("ORDER_DB", "orders")
.rightSource("PAYMENT_DB", "payments")
.joinKey("order_id")
.rules(Arrays.asList(
// 订单已支付但支付记录不是成功
new FieldRule("status", "PAID",
FieldRule.Operator.EQUALS, "payment_status", "SUCCESS"),
// 金额必须一致
new FieldRule("amount", null,
FieldRule.Operator.EQUALS, "amount", null),
// 订单创建时间必须早于支付时间
new FieldRule("create_time", null,
FieldRule.Operator.LESS_THAN, "pay_time", null)
))
.build();
// 规则 2: 订单优惠券状态 vs 优惠券系统状态
public static final ReconciliationRule ORDER_COUPON_RULE =
ReconciliationRule.builder()
.name("ORDER_COUPON_STATUS")
.leftSource("ORDER_DB", "orders")
.rightSource("COUPON_DB", "coupons")
.joinKey("coupon_id")
.rules(Arrays.asList(
// 订单已使用优惠券,优惠券状态必须是已核销
new FieldRule("coupon_status", "USED",
FieldRule.Operator.EQUALS, "status", "REDEEMED")
))
.build();
// 规则 3: 订单库存扣减 vs 库存系统记录
public static final ReconciliationRule ORDER_INVENTORY_RULE =
ReconciliationRule.builder()
.name("ORDER_INVENTORY_DEDUCT")
.leftSource("ORDER_DB", "order_items")
.rightSource("INVENTORY_DB", "inventory_logs")
.joinKey("order_id:sku_id")
.rules(Arrays.asList(
// 讘认单已支付,库存必须有扣减记录
new FieldRule("order_status", "PAID",
FieldRule.Operator.REQUIRES, "deduct_status", "SUCCESS"),
// 扣减数量必须一致
new FieldRule("quantity", null,
FieldRule.Operator.EQUALS, "deduct_quantity", null)
))
.build();
}4.3 全量对账执行
@Service
public class FullReconciliationService {
@Autowired
private ReconciliationRules rulesConfig;
public ReconciliationReport reconcile(LocalDate date) {
ReconciliationReport report = new ReconciliationReport();
report.setDate(date);
// 执行所有规则
for (ReconciliationRule rule : rulesConfig.getAllRules()) {
log.info("Running reconciliation rule: {}", rule.getName());
ReconciliationResult result = runRule(rule, date);
report.addResult(rule.getName(), result);
}
// 汇总报告
report.setTotalChecked(report.getResults().stream()
.mapToInt(r -> r.getTotalChecked()).sum());
report.setTotalInconsistencies(report.getResults().stream()
.mapToInt(r -> r.getInconsistencies().size()).sum());
// 发送报告
if (report.getTotalInconsistencies() > 0) {
alertService.send("RECONCILIATION_REPORT",
String.format("Found %d inconsistencies on %s",
report.getTotalInconsistencies(), date));
}
return report;
}
private ReconciliationResult runRule(ReconciliationRule rule, LocalDate date) {
ReconciliationResult result = new ReconciliationResult();
result.setRuleName(rule.getName());
// 分页拉取左边数据
int pageSize = 1000;
int pageNum = 0;
while (true) {
// 拉取左表数据
List<Map<String, Object>> leftData = dataFetcher.fetch(
rule.getLeftSource(), date, pageNum, pageSize);
if (leftData.isEmpty()) break;
// 根据左表数据拉取右表数据
List<String> joinKeys = leftData.stream()
.map(row -> row.get(rule.getJoinKey()).toString())
.collect(Collectors.toList());
List<Map<String, Object>> rightData = dataFetcher.fetchByKeys(
rule.getRightSource(), joinKeys);
Map<String, Map<String, Object>> rightMap = rightData.stream()
.collect(Collectors.toMap(
row -> row.get(rule.getJoinKey()).toString(),
row -> row));
// 逐条对比
for (Map<String, Object> leftRow : leftData) {
result.incrementTotalChecked();
String key = leftRow.get(rule.getJoinKey()).toString();
Map<String, Object> rightRow = rightMap.get(key);
for (FieldRule fieldRule : rule.getRules()) {
if (!fieldRule.check(leftRow, rightRow)) {
Inconsistency inc = new Inconsistency();
inc.setRuleName(rule.getName());
inc.setJoinKey(key);
inc.setLeftValue(JSON.toJSONString(
fieldRule.extractLeft(leftRow)));
inc.setRightValue(rightRow == null ? "null" :
JSON.toJSONString(fieldRule.extractRight(rightRow)));
inc.setDescription(fieldRule.getDescription());
inc.setCreatedAt(LocalDateTime.now());
result.addInconsistency(inc);
}
}
}
pageNum++;
}
return result;
}
}五、修复方案:补偿、重试、告警闭环
5.1 自动修复策略
┌──────────────────────────────────────────────────────────────────┐
│ 数据不一致修复策略 │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 策略1: 自动修复(安全等级高,无需人工审批) │ │
│ │ - 缓存与DB不一致 → 删除缓存,重新加载 │ │
│ │ - 消息未发送 → 重新发送 │ │
│ │ - 状态机顺序错误 → 按正确顺序更新 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 策略2: 人工审批后修复(涉及资金/库存) │ │
│ │ - 金额不一致 → 生成修复单,财务审批后执行 │ │
│ │ - 库存数量不一致 → 运营审批后调整 │ │
│ │ - 订单状态需要回退 → 技术主管审批 │ │
│ └─────────────────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ 策略3: 告警 + 人工排查(根因不明或影响面大) │ │
│ │ - 大量数据不一致 → 先告警,人工分析根因 │ │
│ │ - 涉及外部系统 → 需要与外部系统协调 │ │
│ │ - 数据丢失不可恢复 → 需要评估业务影响 │ │
│ └─────────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────────┘5.2 自动修复代码实现
@Service
public class AutoRepairService {
// 修复 1: 缓存与 DB 不一致
public void repairCache(String key) {
// 删除缓存,让下次读时从 DB 重新加载
redisTemplate.delete(key);
log.info("Cache repaired for key: {}", key);
}
// 修复 2: 订单状态与支付状态不一致
@Transactional
public void repairOrderPaymentStatus(String orderId) {
PaymentRecord payment = paymentService.getByOrderId(orderId);
Order order = orderService.getById(orderId);
if (payment == null) {
log.warn("Cannot repair: no payment record for order {}", orderId);
alertService.send("REPAIR_FAILED",
"No payment record for order " + orderId);
return;
}
if ("SUCCESS".equals(payment.getStatus()) && !"PAID".equals(order.getStatus())) {
// 支付成功但订单不是已支付状态 → 更新订单状态
order.setStatus("PAID");
order.setPaidAt(payment.getSuccessTime());
orderMapper.updateByPrimaryKey(order);
log.info("Order {} status repaired: → PAID", orderId);
recordRepair("ORDER_STATUS_REPAIR", orderId,
"Updated order status to PAID based on payment success");
}
if ("FAILED".equals(payment.getStatus()) && "PAID".equals(order.getStatus())) {
// 支付失败但订单是已支付状态 → 需要人工处理
alertService.send("CRITICAL_INCONSISTENCY",
String.format("Order %s is PAID but payment FAILED — manual review required",
orderId));
}
}
// 修复 3: 优惠券状态不一致
@Transactional
public void repairCouponStatus(String orderId) {
Order order = orderService.getById(orderId);
if (order.getCouponId() == null) return;
Coupon coupon = couponService.getById(order.getCouponId());
if ("PAID".equals(order.getStatus()) && !"REDEEMED".equals(coupon.getStatus())) {
// 订单已支付但优惠券未核销 → 核销优惠券
coupon.setStatus("REDEEMED");
coupon.setRedeemedAt(order.getPaidAt());
couponMapper.updateByPrimaryKey(coupon);
log.info("Coupon {} status repaired: → REDEEMED", coupon.getId());
recordRepair("COUPON_STATUS_REPAIR", orderId,
"Updated coupon status to REDEEMED based on order payment");
}
if (!"PAID".equals(order.getStatus()) && "REDEEMED".equals(coupon.getStatus())) {
// 优惠券已核销但订单未支付 → 需要人工处理
alertService.send("CRITICAL_INCONSISTENCY",
String.format("Coupon %s is REDEEMED but order %s is %s — manual review required",
coupon.getId(), orderId, order.getStatus()));
}
}
// 修复 4: 库存扣减不一致
@Transactional
public void repairInventoryDeduct(String orderId, String skuId, int expectedQuantity) {
Order order = orderService.getById(orderId);
if ("PAID".equals(order.getStatus())) {
InventoryLog log = inventoryService.getLog(orderId, skuId);
if (log == null) {
// 订单已支付但库存没有扣减记录 → 补扣库存
try {
inventoryService.deduct(skuId, expectedQuantity, orderId);
log.info("Inventory deducted for order {} sku {}", orderId, skuId);
recordRepair("INVENTORY_DEDUCT_REPAIR", orderId,
"Deducted inventory: sku=" + skuId + ", qty=" + expectedQuantity);
} catch (Exception e) {
alertService.send("INVENTORY_REPAIR_FAILED",
String.format("Failed to deduct inventory for order %s: %s",
orderId, e.getMessage()));
}
} else if (log.getQuantity() != expectedQuantity) {
// 扣减数量不一致 → 需要人工处理
alertService.send("INVENTORY_MISMATCH",
String.format("Order %s sku %s: expected qty=%d, actual=%d",
orderId, skuId, expectedQuantity, log.getQuantity()));
}
}
}
private void recordRepair(String type, String orderId, String description) {
RepairRecord record = new RepairRecord();
record.setType(type);
record.setOrderId(orderId);
record.setDescription(description);
record.setRepairedAt(LocalDateTime.now());
repairMapper.insert(record);
}
}5.3 告警闭环
@Service
public class ConsistencyAlertService {
@Autowired
private AlertGateway alertGateway;
@Autowired
private RepairRecordMapper repairMapper;
// 告警发送 + 记录 → 处理 → 确认 → 复查
public void alertAndTrack(String type, String orderId, String message) {
// 1. 创建告警工单
AlertTicket ticket = new AlertTicket();
ticket.setType(type);
ticket.setOrderId(orderId);
ticket.setMessage(message);
ticket.setStatus("OPEN");
ticket.setCreatedAt(LocalDateTime.now());
ticketMapper.insert(ticket);
// 2. 发送告警(钉钉/飞书/邮件)
alertGateway.send(AlertLevel.P1,
String.format("[数据不一致] %s\n订单: %s\n描述: %s\n工单号: %d",
type, orderId, message, ticket.getId()));
// 3. 如果是自动可修复的,触发自动修复
if (isAutoRepairable(type)) {
autoRepairService.repair(type, orderId);
ticket.setStatus("AUTO_REPAIRED");
ticketMapper.updateStatus(ticket);
}
}
// 定时复查未解决的告警
@Scheduled(fixedDelay = 60000) // 每分钟检查
public void checkUnresolvedAlerts() {
List<AlertTicket> openTickets = ticketMapper.selectByStatus("OPEN", 50);
for (AlertTicket ticket : openTickets) {
// 检查是否仍然存在不一致
boolean stillInconsistent = checkConsistency(ticket);
if (!stillInconsistent) {
// 数据已一致,可能是人工修复了
ticket.setStatus("RESOLVED");
ticket.setResolvedAt(LocalDateTime.now());
ticketMapper.updateStatus(ticket);
log.info("Ticket {} resolved", ticket.getId());
} else {
// 升级告警
if (Duration.between(ticket.getCreatedAt(), LocalDateTime.now()).toHours() > 1) {
alertGateway.send(AlertLevel.P0,
String.format("[未解决] 工单 %d 超过 1 小时未处理: %s",
ticket.getId(), ticket.getMessage()));
}
}
}
}
}六、分布式事务方案选型
6.1 方案对比
┌──────────────────────────────────────────────────────────────────┐
│ 分布式事务方案对比 │
├──────────┬──────────┬──────────┬──────────┬─────────────────────┤
│ 方案 │ 一致性 │ 性能 │ 复杂度 │ 适用场景 │
├──────────┼──────────┼──────────┼──────────┼─────────────────────┤
│ 2PC │ 强一致 │ 低 │ 中 │ 传统数据库跨库 │
│ TCC │ 强一致 │ 中 │ 高 │ 资金/库存 │
│ Saga │ 最终一致 │ 高 │ 中 │ 长流程业务 │
│ 本地消息表│ 最终一致 │ 高 │ 低 │ 异步通知 │
│ MQ事务 │ 最终一致 │ 高 │ 中 │ 消息驱动 │
│ 最大努力 │ 最终一致 │ 最高 │ 最低 │ 非核心数据 │
└──────────┴──────────┴──────────┴──────────┴─────────────────────┘6.2 TCC 方案实现
// TCC (Try-Confirm-Cancel) 实现
@Service
public class TccOrderService {
@Autowired
private TccTransactionManager tccManager;
// Try 阶段:预留资源
public TccResult tryCreateOrder(OrderRequest request) {
String xid = tccManager.begin();
try {
// 预创建订单
orderService.tryCreate(xid, request);
// 预扣库存
inventoryService.tryDeduct(xid, request.getSkuId(), request.getQuantity());
// 预核销优惠券
couponService.tryRedeem(xid, request.getCouponId());
// 预冻结金额
accountService.tryFreeze(xid, request.getUserId(), request.getAmount());
tccManager.registerConfirm(xid, () -> confirmCreateOrder(xid, request));
tccManager.registerCancel(xid, () -> cancelCreateOrder(xid, request));
return TccResult.success(xid);
} catch (Exception e) {
tccManager.rollback(xid);
return TccResult.fail(e.getMessage());
}
}
// Confirm 阶段:确认所有操作
@Transactional
public void confirmCreateOrder(String xid, OrderRequest request) {
orderService.confirmCreate(xid);
inventoryService.confirmDeduct(xid);
couponService.confirmRedeem(xid);
accountService.confirmDeduct(xid);
}
// Cancel 阶段:回滚所有操作
@Transactional
public void cancelCreateOrder(String xid, OrderRequest request) {
orderService.cancelCreate(xid);
inventoryService.cancelDeduct(xid);
couponService.cancelRedeem(xid);
accountService.unfreeze(xid);
}
}
// TCC 事务管理器:处理 Confirm/Cancel 失败的重试
@Service
public class TccTransactionManager {
@Scheduled(fixedDelay = 10000) // 每 10 秒检查
public void retryPendingConfirm() {
List<TccTransaction> pending = tccMapper.selectPendingConfirm(100);
for (TccTransaction tx : pending) {
try {
tx.getConfirmAction().run();
tx.setStatus("CONFIRMED");
tccMapper.updateStatus(tx);
} catch (Exception e) {
tx.setRetryCount(tx.getRetryCount() + 1);
if (tx.getRetryCount() >= 10) {
alertService.send("TCC_CONFIRM_FAILED",
"TCC " + tx.getXid() + " confirm failed 10 times");
}
tccMapper.updateStatus(tx);
}
}
}
@Scheduled(fixedDelay = 10000)
public void retryPendingCancel() {
List<TccTransaction> pending = tccMapper.selectPendingCancel(100);
for (TccTransaction tx : pending) {
try {
tx.getCancelAction().run();
tx.setStatus("CANCELLED");
tccMapper.updateStatus(tx);
} catch (Exception e) {
tx.setRetryCount(tx.getRetryCount() + 1);
if (tx.getRetryCount() >= 10) {
alertService.send("TCC_CANCEL_FAILED",
"TCC " + tx.getXid() + " cancel failed 10 times");
}
tccMapper.updateStatus(tx);
}
}
}
}6.3 Saga 模式实现
// Saga 模式:编排式补偿事务
@Service
public class SagaOrderService {
@Autowired
private SagaOrchestrator orchestrator;
public SagaResult createOrder(OrderRequest request) {
SagaBuilder saga = SagaBuilder.create("create-order")
// 步骤 1: 创建订单
.step("create-order")
.action(() -> orderService.create(request))
.compensate(orderId -> orderService.cancel(orderId))
// 步骤 2: 扣减库存
.step("deduct-inventory")
.action(() -> inventoryService.deduct(request.getSkuId(), request.getQuantity()))
.compensate(() -> inventoryService.restore(request.getSkuId(), request.getQuantity()))
// 步骤 3: 核销优惠券
.step("redeem-coupon")
.action(() -> couponService.redeem(request.getCouponId()))
.compensate(() -> couponService.restore(request.getCouponId()))
// 步骤 4: 发起支付
.step("initiate-payment")
.action(() -> paymentService.initiate(request))
.compensate(paymentId -> paymentService.cancel(paymentId))
.build();
return orchestrator.execute(saga);
}
}
// Saga 编排器
@Service
public class SagaOrchestrator {
public SagaResult execute(Saga saga) {
List<SagaStep> completedSteps = new ArrayList<>();
for (SagaStep step : saga.getSteps()) {
try {
step.getAction().run();
completedSteps.add(step);
} catch (Exception e) {
// 正向执行失败,逆向补偿已完成的步骤
log.error("Saga {} failed at step {}, compensating...",
saga.getName(), step.getName(), e);
Collections.reverse(completedSteps);
for (SagaStep completed : completedSteps) {
try {
completed.getCompensate().run();
log.info("Compensated step: {}", completed.getName());
} catch (Exception ce) {
log.error("Compensation failed for step: {}", completed.getName(), ce);
alertService.send("SAGA_COMPENSATION_FAILED",
String.format("Saga %s compensation failed at %s: %s",
saga.getName(), completed.getName(), ce.getMessage()));
}
}
return SagaResult.fail(e.getMessage());
}
}
return SagaResult.success();
}
}七、缓存与数据库不一致排查
7.1 Cache Aside 模式的问题
// 问题代码:先更新 DB,再删除缓存
public void updateUser(User user) {
db.update(user);
redis.delete("user:" + user.getId()); // 如果这一步失败?
}
// 并发场景下的问题:
// 线程 A: 更新 DB → 删除缓存(失败)
// 线程 B: 读缓存(还是旧值)→ 返回旧值
// 结果:DB 已更新,缓存还是旧值
// 更严重的情况:
// 线程 A: 读 DB 旧值 → 准备写缓存
// 线程 B: 更新 DB → 删除缓存
// 线程 A: 写入缓存(旧值) ← 覆盖了删除操作
// 结果:DB 是新值,缓存是旧值,永远不一致!7.2 解决方案
// 方案 1: 延迟双删
public void updateUser(User user) {
// 第一次删除
redis.delete("user:" + user.getId());
// 更新 DB
db.update(user);
// 延迟第二次删除(确保读到旧值的线程不会覆盖)
scheduledExecutor.schedule(() -> {
redis.delete("user:" + user.getId());
}, 500, TimeUnit.MILLISECONDS); // 延迟 500ms
}
// 方案 2: 订阅 binlog 删除缓存(推荐)
@Component
public class BinlogCacheInvalidator {
@KafkaListener(topics = "user-db-binlog")
public void onUserChange(ChangeEvent event) {
if (!"users".equals(event.getTable())) return;
String userId = event.getPrimaryKey();
String cacheKey = "user:" + userId;
// 删除缓存
redis.delete(cacheKey);
log.info("Cache invalidated via binlog: {}", cacheKey);
}
}
// 方案 3: 设置合理的 TTL(兜底方案)
// 即使删除失败,缓存也会在 TTL 后过期
redis.opsForValue().set("user:" + userId, user, 30, TimeUnit.MINUTES);八、面试要点
Q1: 分布式系统数据不一致怎么排查?
参考回答:
- 确定不一致的范围:是两个系统之间不一致,还是缓存与 DB 不一致?
- 对账比对:按业务主键关联两个系统的数据,逐条对比关键字段
- 分析 binlog:通过 binlog 查看数据变更时间线,找到不一致的时间点
- 排查根因:
- 跨服务操作部分成功 → 引入分布式事务
- 消息消费失败 → 保证幂等 + 重试
- 缓存不一致 → 延迟双删或 binlog 删除
- 并发覆盖 → 加锁或乐观锁
- 修复数据:能自动修复的自动修复,涉及资金的走人工审批
Q2: 分布式事务有哪些方案?各自适用什么场景?
参考回答:
- 2PC:强一致,性能差,适合传统数据库跨库
- TCC:强一致,需要写 Try/Confirm/Cancel 三个方法,适合资金/库存
- Saga:最终一致,正向+补偿,适合长流程业务
- 本地消息表:最终一致,简单可靠,适合异步通知
- MQ 事务消息:RocketMQ 事务消息,半消息+回查,适合消息驱动
选型原则:优先选最终一致性方案(本地消息表/Saga),只在资金类场景用 TCC。
Q3: 如何保证消息消费的幂等性?
参考回答:
三种方案:
- Redis SET NX:消费前用
SETNX msg:processed:{id} 1标记,24 小时过期- 数据库唯一索引:插入幂等记录表,message_id 有唯一索引,DuplicateKeyException 表示重复
- 乐观锁:在业务表加 version 字段,
UPDATE ... SET version=version+1 WHERE version=?推荐数据库唯一索引方案,因为和业务事务在同一事务中,一致性更好。
Q4: 缓存和数据库不一致怎么办?
参考回答:
最佳实践是 Cache Aside + binlog 删除:
- 写时:先更新 DB,再删除缓存
- 读时:先读缓存,miss 后读 DB 并写回缓存
- 兜底:binlog 订阅删除缓存(解决删除失败的问题)
- 最终兜底:设置 TTL,确保即使删除失败也会在过期后自愈
如果对一致性要求极高,可以用延迟双删:先删 → 更新 → 延迟再删。
九、避坑指南
坑 1:对账系统本身的数据不一致
// 对账系统拉取数据时,两个系统的时间基准可能不同
// 订单系统按 create_time 拉取,支付系统按 pay_time 拉取
// 可能导致同一天的数据在两个系统中数量不对
// 解决:对账时间窗口加大,比如对 7 月 3 日的数据,拉取时间范围为 7/3 00:00 - 7/4 06:00
LocalDateTime startTime = date.atStartOfDay();
LocalDateTime endTime = date.plusDays(1).atTime(6, 0); // 第二天早上 6 点坑 2:自动修复导致数据被覆盖
// 如果用户正在操作,自动修复同时执行,可能互相覆盖
// 修复前检查:如果数据在此期间被修改过,跳过自动修复
public void repairOrderStatus(String orderId) {
Order current = orderService.getById(orderId);
// 检查最后更新时间,如果最近被修改过,说明用户在操作,跳过修复
if (current.getUpdatedAt().isAfter(LocalDateTime.now().minusMinutes(5))) {
log.info("Order {} recently updated, skip auto repair", orderId);
return;
}
// 执行修复
orderService.updateStatus(orderId, "PAID");
}坑 3:消息重试导致的重复消费
// Kafka 手动 ACK 场景下,ACK 之前业务逻辑执行了但 ACK 失败
// 消息会被重新投递,导致重复执行
// 解决:消费端必须幂等,用唯一消息 ID 去重
// 注意 Redis SETNX 方案的问题:如果 Redis 挂了,标记丢失,会重复消费
// 最可靠的方案:数据库唯一索引
@Transactional
public void consume(OrderMessage message) {
// 先插入幂等记录(同一事务)
idempotentMapper.insert(new IdempotentRecord(message.getId()));
// 执行业务
orderService.process(message);
// 事务提交:幂等记录和业务同时成功或同时失败
}坑 4:TCC 的空回滚和悬挂
// 空回滚:Cancel 在 Try 之前执行(网络延迟导致 Try 还没到,先触发了 Cancel)
// 悬挂:Try 在 Cancel 之后执行(Cancel 已经执行了,Try 才到,导致资源被锁定无法释放)
// 解决:在 Try 和 Cancel 时检查事务状态
public void tryDeduct(String xid, String skuId, int qty) {
// 检查是否已经 Cancel 了
if (tccMapper.isCancelled(xid)) {
log.warn("TCC {} already cancelled, skip try", xid);
return; // 拒绝 Try
}
// 正常执行 Try
inventoryService.preDeduct(skuId, qty);
tccMapper.markTried(xid);
}
public void cancelDeduct(String xid, String skuId, int qty) {
// 检查是否 Try 过
if (!tccMapper.isTried(xid)) {
log.warn("TCC {} not tried yet, empty cancel", xid);
tccMapper.markCancelled(xid); // 标记 Cancel 已执行
return; // 空回滚
}
// 正常执行 Cancel
inventoryService.restore(skuId, qty);
tccMapper.markCancelled(xid);
}坑 5:对账全量拉取 OOM
// 对账时一次性拉取全量数据到内存,数据量大时直接 OOM
// 解决:分页拉取 + 流式处理
public void reconcile(LocalDate date) {
int pageSize = 1000;
int pageNum = 0;
while (true) {
// 每次只拉 1000 条
List<Order> batch = orderService.getByDate(date, pageNum, pageSize);
if (batch.isEmpty()) break;
// 流式处理,处理完就 GC
for (Order order : batch) {
checkConsistency(order);
}
pageNum++;
}
}
// 更优方案:使用游标查询,避免 OFFSET 越来越慢
public void reconcileWithCursor(LocalDate date) {
Long lastId = 0L;
while (true) {
List<Order> batch = orderService.getByDateAfterId(date, lastId, 1000);
if (batch.isEmpty()) break;
for (Order order : batch) {
checkConsistency(order);
}
lastId = batch.get(batch.size() - 1).getId(); // 游标推进
}
}十、预防体系建设
10.1 数据一致性保障架构
┌──────────────────────────────────────────────────────────────────┐
│ 数据一致性保障体系 │
│ │
│ 事前预防: │
│ ├── 分布式事务方案(TCC/Saga/本地消息表) │
│ ├── 幂等性设计(唯一索引/Redis 去重) │
│ ├── 状态机约束(状态流转规则校验) │
│ └── Code Review 检查跨服务操作 │
│ │
│ 事中监控: │
│ ├── CDC 实时对账(binlog → 对比 → 告警) │
│ ├── 关键指标监控(成功率、失败率、延迟) │
│ └── 消息积压告警 │
│ │
│ 事后修复: │
│ ├── 每日全量对账 │
│ ├── 自动修复(缓存、状态) │
│ ├── 人工审核修复(金额、库存) │
│ └── 告警闭环(告警 → 处理 → 确认 → 复查) │
└──────────────────────────────────────────────────────────────────┘10.2 一致性检查清单
| 检查项 | 说明 | 风险等级 |
|---|---|---|
| 跨服务操作是否有事务保障 | 是否使用 TCC/Saga/消息表 | 高 |
| 消息消费是否幂等 | 是否有去重机制 | 高 |
| 缓存是否有 TTL 兜底 | 即使删除失败也能自愈 | 中 |
| 是否有定时对账 | 每日/每小时对账 | 高 |
| 对账差异是否有告警 | 发现差异能及时通知 | 高 |
| 修复操作是否记录日志 | 可追溯修复历史 | 中 |
| 状态变更是否有状态机校验 | 防止非法状态流转 | 高 |
| 并发更新是否有乐观锁 | 防止覆盖更新 | 中 |
总结
数据不一致排查的核心流程:用户投诉/对账发现 → 按 orderId 关联各系统数据 → 逐项对比 → 分析 binlog 时间线 → 定位根因 → 修复数据 → 补偿机制建设。
三个核心经验:
- 对账系统是最后一道防线——即使有完美的分布式事务方案,也需要对账系统兜底
- 消息消费必须幂等——这是最终一致性的基石,用数据库唯一索引最可靠
- 缓存不一致用 binlog 删除方案——比延迟双删更可靠,是业界最佳实践