拆分服务时怎样处理双写一致性
2026/8/21 9:58:12 网站建设 项目流程

拆分服务时怎样处理双写一致性

服务拆分涉及数据边界和双写一致性时,先记录目标、假设、回滚方式和核对口径。本文提供一种复盘模板,避免用一段戏剧化经历替代决策证据。

上个月,团队将原本挤在同一个 MySQL 大库中的“订单服务”与“支付服务”进行物理独立拆分。为了保证迁移过程中业务零停机,设计了“老库主写 -> 异步 Binlog 双写新库 -> 动态切流 -> 延迟对齐校验”的平滑过渡方案。

然而在流量切换的瞬间,主从延迟叠加高并发扣款请求,直接引发了 30 多笔订单双写状态不一致。用户扣了钱,订单服务状态却依然停留在“待支付”。现场排查与对齐数据花了整整一个通宵。

拆分分布式服务绝不是画几张微服务架构图那么简单。必须深刻复盘数据库双写平滑迁移中的锁竞争与数据补偿机制,并形成可复用的架构决策记录(ADR)模板。


1. 数据库平滑切流与双写一致性架构

在单体拆分为分布式微服务的过渡期,采用“Canal Binlog 增量监听 + 延迟异步补偿”双写架构。

业务写请求依然发往老库,Canal 实时捕获老库order_info表的 Binlog 变更日志,将其转化为 JSON 消息推送到 RocketMQ。拆分后的新订单服务订阅 MQ 消息并写入新库。同时后台运行一个 Async Reconciliation Task(异步对齐任务),对 5 分钟前产生的订单进行双向 Hash 校验与补缝。


2. 生产环境故障排查与 Binlog 延迟诊断命令

当切流期间收到数据不一致报警时,使用以下命令快速诊断 Canal 延迟与 Mysql 事务锁状态。

# 1. 检查 Canal 消费老库 Binlog 的延迟秒数 (Delay Metrics) curl -s "${CANAL_ADMIN_URL}/api/v1/canal/instance/delay?destination=order_instance" # 2. 监控 Mysql 慢日志与锁等待状态 mysql -h legacy-db.internal -u root -p'Pass0821!' -e " SELECT r.trx_id waiting_trx_id, r.trx_mysql_thread_id waiting_thread, r.trx_query waiting_query, b.trx_id blocking_trx_id, b.trx_mysql_thread_id blocking_thread, b.trx_query blocking_query FROM information_schema.innodb_lock_waits w INNER JOIN information_schema.innodb_trx b ON b.trx_id = w.blocking_trx_id INNER JOIN information_schema.innodb_trx r ON r.trx_id = w.waiting_trx_id;" # 3. 使用 pt-table-checksum 检查老库与新库表结构数据是否一致 pt-table-checksum --replicate=test.checksums h=legacy-db.internal,u=migrator,p='Pass0821!' \ --databases=order_db --tables=t_order # 4. 统计 RocketMQ 数据迁移 Topic 积压量 mqadmin consumerProgress -n rocketmq-namesrv.internal:9876 -g order_migration_consumer_group

排查诊断发现:由于老库在订单完成时触发了包含大字段更新的复杂事务,导致的 Binlog 变更日志体积激增。RocketMQ 消费线程在反序列化 JSON 时产生了阻塞,单消费节点积压超过 40,000 条消息,最终引发切流瞬间新库数据落后老库 1.2 秒。


3. 生产级数据双写对齐与延迟补偿代码

在 Spring Boot 体系中,实现后台定时数据一致性比对与自动补偿纠错器。代码采用分段分批 Fetch 机制,避免对老库造成二次读压力。

package com.example.migration.task; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.List; import java.util.Map; import java.util.Objects; @Component public class DataReconciliationScheduledTask { private static final Logger log = LoggerFactory.getLogger(DataReconciliationScheduledTask.class); private final JdbcTemplate legacyJdbcTemplate; private final JdbcTemplate newJdbcTemplate; public DataReconciliationScheduledTask(JdbcTemplate legacyJdbcTemplate, JdbcTemplate newJdbcTemplate) { this.legacyJdbcTemplate = legacyJdbcTemplate; this.newJdbcTemplate = newJdbcTemplate; } // 每 2 分钟扫描过去 5 分钟至 10 分钟之间产生的数据,规避正常的毫秒级 Binlog 延迟 @Scheduled(cron = "0 */2 * * * ?") public void reconcileOrderData() { log.info("Starting scheduled data reconciliation task..."); long startTime = System.currentTimeMillis() - (10 * 60 * 1000); long endTime = System.currentTimeMillis() - (5 * 60 * 1000); String sql = "SELECT order_id, status, amount, updated_at FROM t_order WHERE updated_at BETWEEN ? AND ?"; List<Map<String, Object>> legacyRecords = legacyJdbcTemplate.queryForList(sql, new java.util.Date(startTime), new java.util.Date(endTime)); int mismatchCount = 0; for (Map<String, Object> legacyRow : legacyRecords) { String orderId = (String) legacyRow.get("order_id"); Integer legacyStatus = (Integer) legacyRow.get("status"); // 查询新库对应记录 String newDbSql = "SELECT status, amount FROM t_order WHERE order_id = ?"; List<Map<String, Object>> newRecords = newJdbcTemplate.queryForList(newDbSql, orderId); if (newRecords.isEmpty()) { log.warn("Data missing in new DB for orderId: {}, initiating repair insertion", orderId); repairMissingData(legacyRow); mismatchCount++; } else { Integer newStatus = (Integer) newRecords.get(0).get("status"); if (!Objects.equals(legacyStatus, newStatus)) { log.warn("Status mismatch for orderId: {}. Legacy: {}, New: {}. Repairing...", orderId, legacyStatus, newStatus); repairStatusMismatch(orderId, legacyStatus); mismatchCount++; } } } log.info("Data reconciliation task completed. Scanned: {}, Mismatches Repaired: {}", legacyRecords.size(), mismatchCount); } private void repairMissingData(Map<String, Object> row) { String insertSql = "INSERT INTO t_order (order_id, status, amount, updated_at) VALUES (?, ?, ?, ?) " + "ON DUPLICATE KEY UPDATE status = VALUES(status), amount = VALUES(amount)"; newJdbcTemplate.update(insertSql, row.get("order_id"), row.get("status"), row.get("amount"), row.get("updated_at")); } private void repairStatusMismatch(String orderId, Integer correctStatus) { String updateSql = "UPDATE t_order SET status = ? WHERE order_id = ?"; newJdbcTemplate.update(updateSql, correctStatus, orderId); } }

4. 可复制的服务拆分项目复盘与决策记录模板

为了在后续其他业务线拆分时不再重复踩坑,团队必须把经验整理成可复用的规范模板。

架构决策记录模板 (ADR-Split-Pattern)

# [ADR-编号] 服务拆分与数据迁移架构决策记录 ## 1. 拆分背景与业务驱动力 - **原单体痛点**:[描述单体 DB/代码库的耦合瓶颈,如并发锁冲突、部署相互拖累] - **拆分目标**:[明确拆分后的微服务边界与预期 QPS 目标] ## 2. 数据库切流与迁移方案选择 - **方案 A (停机迁移)**:[优点:简单无一致性风险;缺点:业务不可接受 2 小时停机] - **方案 B (Binlog 双写 + 异步补偿)**:[最终选中方案,无感切流,需配套对齐校验任务] ## 3. 关键风险与应对防护 - **风险 1:Binlog 积压导致的读写不一致** - *防护手段*:设置 5 分钟对齐窗口,切流前检查 MQ Lag 必须为 0,否则拒绝切流。 - **风险 2:分布式事务跨库一致性** - *防护手段*:禁止跨库本地事务,全面改用 RocketMQ 事务消息或 Seata AT 模式。 ## 4. 验收与回滚标准 - **切流完成标准**:双写校验连续 48 小时 0 Mismatch 报告。 - **一键回滚预案**:Nacos 动态开关在 3 秒内切回老库读取,保留老库写双发。

5. 服务拆分复盘总结

分布式系统拆分不是一蹴而就的潇洒重构,而是带伤换引擎的精密工程。

通过 Canal 增量同步、双写补偿校验任务以及标准化 ADR 决策模板的落地,团队才真正掌握了平滑拆分单体架构的技术主动权。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询