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

发布时间:2026/8/21 9:58:27
拆分服务时怎样处理双写一致性 拆分服务时怎样处理双写一致性服务拆分涉及数据边界和双写一致性时先记录目标、假设、回滚方式和核对口径。本文提供一种复盘模板避免用一段戏剧化经历替代决策证据。上个月团队将原本挤在同一个 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?destinationorder_instance # 2. 监控 Mysql 慢日志与锁等待状态 mysql -h legacy-db.internal -u root -pPass0821! -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 --replicatetest.checksums hlegacy-db.internal,umigrator,pPass0821! \ --databasesorder_db --tablest_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 ?; ListMapString, Object legacyRecords legacyJdbcTemplate.queryForList(sql, new java.util.Date(startTime), new java.util.Date(endTime)); int mismatchCount 0; for (MapString, 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 ?; ListMapString, 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(MapString, 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. 关键风险与应对防护 - **风险 1Binlog 积压导致的读写不一致** - *防护手段*设置 5 分钟对齐窗口切流前检查 MQ Lag 必须为 0否则拒绝切流。 - **风险 2分布式事务跨库一致性** - *防护手段*禁止跨库本地事务全面改用 RocketMQ 事务消息或 Seata AT 模式。 ## 4. 验收与回滚标准 - **切流完成标准**双写校验连续 48 小时 0 Mismatch 报告。 - **一键回滚预案**Nacos 动态开关在 3 秒内切回老库读取保留老库写双发。5. 服务拆分复盘总结分布式系统拆分不是一蹴而就的潇洒重构而是带伤换引擎的精密工程。通过 Canal 增量同步、双写补偿校验任务以及标准化 ADR 决策模板的落地团队才真正掌握了平滑拆分单体架构的技术主动权。