ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

开源一个可视化 MySQL 数据同步工具 DataMove:全量 + Canal 增量 + 数据对账一键修复 - 二

开源一个可视化 MySQL 数据同步工具 DataMove:全量 + Canal 增量 + 数据对账一键修复 - 二 在线体验http://8.140.200.84/官网http://8.140.200.84/website/一个基于 RuoYi-Vue 二次开发的轻量 MySQL 同步平台把「配 DataX JSON / 手写 Canal 消费端 / SSH 隧道跑 SQL」这三件苦差事收敛成浏览器里的点选操作。本文不讲空话架构设计、核心代码实现与性能实测数字全部来自项目当前代码与运行日志。目录一、为什么要造这个轮子二、能力总览三、技术栈与整体架构四、核心设计 1全量同步的断点续传 幂等 分片并行五、核心设计 2Canal 增量同步的三层过滤与 ACK 顺序六、核心设计 3数据校验与一键修复双游标归并七、核心设计 4在线 SQL 工作台八、可观测性实时指标 / 运行历史 / 日志 / 审计九、告警钉钉 邮件双通道十、配置项与表结构速查十一、性能实测10 万 / 100 万行全量同步十二、10 分钟跑起来十三、写在最后一、为什么要造这个轮子团队过去三年在数据同步上反复踩坑场景常见做法真实痛点业务库拆分订单库 → 订单中台库DataX 写 JSON字段映射手写 joiner运维同学上手成本高改一个列名要动配置实时订阅 binlog部署 Canal 自写消费端要自己处理位点、多线程 ACK、异常回滚错一次就丢/重复数据上线前对账两边各count(*)比总数只知道差 37 行不知道是哪 37 行、哪个字段修不了临时数据修复登录跳板机手写INSERT ... SELECT无审计、无日志出事没人知道谁改的DBA 离职交接脚本散落在各台机器新人接手无从下手于是有了 DataMove目标很明确把同步变成配置 一键启动把对账变成看得见差异 一键补齐。二、能力总览┌─────────────────────────────────────────────────────────────────────┐ │ DataMove 能力地图 │ ├──────────────┬──────────────────────────────────────────────────────┤ │ 数据源 │ MySQL 5.7 / 8.x多源并存密码 AES 加密存储 │ │ 全量同步 │ 按主键 ID / 按时间字段断点续传 幂等覆盖 │ │ 分片并行 │ FULLID 模式按主键区间切 N 片多线程并行读写 │ │ 增量同步 │ Canal Client 订阅 binlog ROW库表/DML 类型/字段三层过滤 │ │ 表结构同步 │ DDL 任务目标表不存在则按源表 SHOW CREATE TABLE 建表 │ │ 字段映射 │ kettle 风格拖拽连线源列 → 目标列 │ │ 数据校验 │ 双游标流式归并定位到行 字段级差异 │ │ 一键修复 │ 补缺失 INSERT 修不一致 UPDATE**不删目标库数据** │ │ 数据中心 │ 在线分页浏览 / 增删改数据 表结构与索引管理 │ │ SQL 工作台 │ CodeMirror 编辑器补全/高亮多语句EXPLAIN 执行计划 │ │ 可观测性 │ 实时指标大盘 运行历史 批次日志 字段级审计日志 │ │ 告警 │ 钉钉 Webhook 邮件异步不阻塞同步线程 │ └──────────────┴──────────────────────────────────────────────────────┘对比 DataX / Canal-adapter / 云 DTS 这类方案DataMove 的差异化在零代码 全中文 自带对账与审计适合中小团队开箱即用。三、技术栈与整体架构层级技术后端Spring Boot 2.7.18 MyBatis-Plus 3.5.5 Spring Security JWT前端Vue 2.7 Element UI 2.13RuoYi-Vue 4.8.1同步原生 JDBC 分批 Canal Client 1.1.4增量校验源/目标双游标流式归并setFetchSize(Integer.MIN_VALUE)内存只驻留一行加密AESAES/CBC/PKCS5Padding Base64调度SpringScheduled 自建线程池告警钉钉 Webhook SMTP 邮件数据库MySQL 8.0前端 (Vue 2 Element UI)后端 (SpringBoot 2.7)存储基础设施HTTPJDBC 分批TCP 11111订阅 ROW binlog双游标归并 / 回放修复数据源 / 任务 / 日志SQL 工作台 / 数据中心Controller 层任务服务FullSyncEngine全量引擎RangeSplitter分片切分DdlSyncEngine表结构引擎CanalSyncEngine增量引擎DataVerifyEngine校验 一键修复SqlConsoleSQL 引擎TaskMetrics实时指标SyncLogService异步日志AuditLogService审计AlertUtils钉钉 邮件MySQL 8.0元数据 业务Canal Server订阅 ROW binlog前端只负责展示与配置所有同步动作都落在后端引擎三个引擎全量/增量/校验共用同一套指标、日志与告警设施。四、核心设计 1全量同步的断点续传 幂等 分片并行4.1 两种游标模式模式适用场景WHERE 片段按主键 ID有自增主键的业务表WHERE id #{lastId} ORDER BY id ASC LIMIT #{batchSize}按时间字段无自增主键 / 日志型表WHERE time #{lastTime} OR (time #{lastTime} AND id #{lastIdInBatch}) ORDER BY time ASC, id ASC LIMIT ...按时间模式的关键在同秒数据先用time #{lastTime} AND id #{lastIdInBatch}把同一秒内剩余行捞完再推进lastTime避免漏数据 / 重复拉取。小细节时间断点必须用rs.getTimestamp()读取不能用getObject()—— 否则时区/driver 差异会让断点值漂移。4.2 断点推进的唯一原则整批数据成功写入并 commit 之后才更新断点。while (!ctx.isStopped() 有下一批) { rows selectBatch(lastId); // PreparedStatement setFetchSize(batchSize) 流式读 try { batchInsertOrUpdate(rows); // INSERT ... ON DUPLICATE KEY UPDATE updateProgress(rows.lastId()); // ✅ 提交成功后才推进断点 } catch (BatchUpdateException e) { rollback(); // ❌ 失败不推进下次从断点重来 // 记日志 告警不静默吞异常 } }配套的两个硬要求幂等写入用INSERT ... ON DUPLICATE KEY UPDATE而不是REPLACE INTO。REPLACE内部是先DELETE再INSERT会改自增主键、触发级联删除ON DUPLICATE KEY UPDATE只更新既有列重跑不脏数据。流式读取MySQL JDBC 默认会把结果集全部拉进内存必须显式setFetchSize(batchSize)否则百万行必 OOM。4.3 分片并行把 31 分钟压到几分钟FULL ID模式下任务配置shard_count 1时会走分片查SELECT MIN(id), MAX(id)RangeSplitter.split(min, max, shardCount)把主键区间等宽切成 N 段纯函数好测每个分片一个线程 一对独立的源/目标连接各自PreparedStatement与setFetchSize区间互不重叠进度按synchronized(progressLock)聚合后回写sync_task_progress全局游标取各分片最大值。与断点续传的互斥规则很重要否则会重复同步已有断点progress.total_rows 0时自动回退单线程续传空表、区间划分失败、只能切出 1 段时也回退。回退时会打一行日志分片并行回退单线程续传避免出现看着是分片任务实际在重复搬数据的灵异现象。4.4batch_size和shard_count是两回事维度批次batch_size分片shard_count切分方向时间维度分次空间维度分区工作方式单线程按LIMIT游标循环一批 commit 完再读下一批整表按主键区间切 N 块N 个线程同时各管一块目的控内存、缩短事务、失败只重试一小批提升并行吞吐耗时近线性 / N推荐值1000 ~ 100001 ~ 16参考目标库写入能力串行 (shard_count1): 线程1: [批次1]→[批次2]→[批次3]→...→[批次N] 总时间 T 分片并行 (shard_count4): 线程1: [批次1]→[批次5]→... ┐ 线程2: [批次2]→[批次6]→... ├─ 同时进行 总时间 ≈ T/4 线程3: [批次3]→[批次7]→... ┘ 线程4: [批次4]→[批次8]→...4.5 暂停 / 停止怎么做用SyncContext持有AtomicBoolean pauseFlag / stopFlag主循环以及每个分片线程内循环在批次边界检查标志位并优雅收口暂停会先把当前批次写完再停因此不会产生半批脏数据。任务级单实例保护用static SetLong RUNNING_TASKsynchronized start()同一任务不会有两个 worker。五、核心设计 2Canal 增量同步的三层过滤与 ACK 顺序5.1 主流程connect() → subscribe(源库\.任务表) → while(true) { getWithoutAck(1000) } ↓ apply(entry) → 目标库 JDBC 执行 ↓ connector.ack(batchId) // ✅ 成功才 ACK ↓ upsert sync_canal_position5.2 三层事件过滤让无关事件根本不进客户端层手段效果库/表过滤Canal 服务端subscribe表达式从.*\..*收紧为源库\.任务表库表名做正则转义不相关库表的 binlog 事件根本不进客户端省网络与解析开销DML 类型过滤任务配置binlog_dml_typesINSERT / UPDATE / DELETE 按需勾选未勾选事件计数丢弃归档库只收 INSERT、审计库不要 DELETE 等场景字段过滤任务配置ignore_fields拼 SQL 时自动剔除这些列目标库自带的create_time/tenant_id不被源库覆盖几个实现细节值得说明key 列永远保留即使把主键填进ignore_fields也会被跳过否则 UPDATE / DELETE 就没有 WHERE 条件了。忽略字段按映射后的目标列名匹配配了字段映射也不用改这里。忽略字段同时被校验引擎读取update_time、ON UPDATE CURRENT_TIMESTAMP这类数据库自动维护的列天然与源库不同不忽略会刷出满屏假差异。不勾 不过滤老任务零感知。5.3 位点必须在 ACK 之后推进try { apply(events); // 先写目标库 connector.ack(batchId); // ✅ ACK 成功 upsertPosition(journal, pos); // ✅ 再落位点 } catch (Exception e) { connector.rollback(batchId); // ❌ 不 ACK、不推进位点下次重放 alert(task, 增量批次异常, ...); }顺序反了会怎样如果先写位点、后 ACK程序在这两步之间崩溃重启后位点已经前移这批事件就永久丢了。反过来先 ACK 后写位点最坏情况是重复执行一批而写入是幂等的可接受。在分布式同步里可能重复永远优于可能丢。5.4 反向回环与 UPDATE / DELETE 定位回环防护源库写入如果与目标库同实例目标库的写入也会产生 binlog 被 Canal 捕获。引擎在处理每条 entry 时校验schema.equalsIgnoreCase(源库名)并跳过非任务表事件从源头掐掉自己同步自己。UPDATE 精确定位使用before镜像里的 key 列拼WHERE命中 0 行时降级为 upsert命中 1 行打 WARN。DELETE 定位收集全部key 列CanalRowKeys拼 WHERE避免只按单列删除误伤。六、核心设计 3数据校验与一键修复双游标归并同类工具一般只做count(*)比对只能说差 37 行。DataMove 的做法是源库和目标库各开一个流式游标按主键升序做双指针归并源游标: 1 5 9 12 ←→ 目标游标: 1 5 7 12 ↓ 源 key 目标 key → MISSING 源有目标无 → 一键补 INSERT 源 key 目标 key → EXTRA 目标有源无 → 只报告不删 key 相等 → 比字段值不同则 MISMATCH只 UPDATE 差异列相比分段 checksum归并不需要额外建索引、不需要全表排序而且天然能定位到行和字段相比count(*)它才是一键修复的前提。工程实现上的几个约束项取值原因游标读取方式TYPE_FORWARD_ONLYsetFetchSize(Integer.MIN_VALUE)百万行级别内存里只驻留一行进度回填频率每 5000 行抽屉里 2s 轮询展示已比对行数 / 三类差异数差异明细留存上限2000 条超出置truncated1并明确提示仅保留前 N 条但统计数字始终全量准确修复批大小200 条 / 批每条回填修复状态已修复 / 失败 / 跳过侧效应不改sync_task.status校验是只读旁路避免任务已完成却显示运行中的状态错乱两条设计红线不会删除目标库任何数据。EXTRA行只在明细里标出来、修复状态直接置为SKIPPED。删除是不可逆的破坏性动作不应该藏在一键按钮背后。修复只改不一致的列。依据diff_fields生成UPDATE目标库其它列的原值保持不动。另外校验与修复共用ignore_fields且RowDiffUtils.normalize()会处理BigDecimal尾零1.0→1、byte[]转b64:前缀等看着不同其实相同的值避免假差异。七、核心设计 4在线 SQL 工作台同步之外日常运维最常用的就是随便连个库跑条 SQL。DataMove 内置了一个 Web 版 SQL 工作台CodeMirror 5编辑器体验SQL 语法高亮text/x-mysql、括号匹配、当前行高亮补全源包含 SQL 关键字 表名 字段名Ctrl 空格手动唤起、输入标识符自动联想、输入.直接列出该表字段。表结构助手右侧面板展示表信息引擎 / 字符集 / 列 / 索引 / DDL点列名即插入编辑器光标处。多语句执行按分号拆分会跳过引号、反引号、--/#//* */注释里的分号每条语句独立展示结果单次最多 20 条。安全护栏单语句setQueryTimeout(30s)防卡死查询结果上限 1000 行EXPLAIN ANALYZE需要显式二次确认因为它会真的执行SQL。执行计划可视化type ALL / index标橙、Using filesort/Using temporary标红、rows ≥ 10000标橙一眼看出慢在哪。结果导出一键导出 Excel列宽按中文宽度自适应或 CSV带 BOMExcel 打开中文不乱码。SQL 收藏收藏常用 SQL带标题、标签、团队共享开关与使用次数统计。全部落痕每次执行都写sync_sql_log操作人 / IP / SQL / 耗时可在「SQL 执行日志」里检索与导出。八、可观测性实时指标 / 运行历史 / 日志 / 审计8.1 任务大盘实时进程内维护TaskMetrics前端有运行中任务时 3s 自动刷新实时速率10s 滑动窗口计算行/秒不是总耗时平均能立刻反映掉速ETA按源表总行数与当前速率估算剩余时间瓶颈库对比源库读取耗时与目标库写入耗时比值超过1.3判定瓶颈侧增量任务对比等待 binlog 事件与应用变更分片监控每个分片的区间、游标位置、实时速率、状态同步中 / 已完成 / 失败暂停 / 继续 / 停止按钮直接在大盘上操作。8.2 运行历史表sync_task_run实时监控看现在运行历史看过去。每次「启动」任务写一条记录概览卡运行次数 / 成功 / 失败 / 运行中 / 同步行数 / 失败行数 / 平均耗时 / 平均速率近 7 / 14 / 30 天趋势图柱 运行次数成功、失败堆叠线 同步行数多维筛选 关键字检索 CSV 导出口径与页面筛选一致运行中每 5s 回填进度服务重启也能看到跑到哪了完成 / 失败 / 停止 / 暂停都会各自收口一条历史不留悬挂的运行中。8.3 批次日志表sync_task_log每个批次一条含批次号、分片号、同步模式、位点起止、本批行数、累计行数、单批耗时、行/秒、异常信息。增量任务的位点是binlog 文件:offset管理员一眼能定位到具体 binlog 位置。8.4 字段级审计日志表sync_audit_log不是笼统的任务被修改过而是谁 / 什么时候 / 改了哪个任务的哪个字段old → new/ 从哪个 IP 和 UA 改的。同一次请求修改多个字段共享一个revision_id详情页可一键展开同一时刻发生了什么任务名、操作人写库时快照任务改名后历史依然读得懂写入是旁路审计落库失败只log.warn绝不搞挂主流程 —— 合规功能不能反过来影响业务UI 上行底色按操作类型区分旧值红色删除线、新值绿色高亮。九、告警钉钉 邮件双通道统一入口AlertUtils.alert(task, subject, content)一条失败消息同时投递任务上配置的钉钉机器人dingtalk_webhook任务上配置的alert_emailSMTP 服务器为全局配置。任一渠道未配置自动跳过互不影响。触发场景覆盖启动失败源/目标库连接失败、全量同步异常、增量批次异常、增量连接异常、表结构同步失败、校验发现差异、校验失败、修复失败。工程细节告警走独立单线程池队列 200DiscardOldestPolicy丢弃最旧消息异步发送不阻塞同步线程任何告警异常只记日志不影响任务状态流转—— 告警是附属功能不能反过来拖垮主链路。十、配置项与表结构速查sync_task中与同步行为相关的字段字段含义task_typeFULL/INCR/DDLsync_modeID/TIME/BINLOG/DDLbatch_size批次大小默认 1000shard_count分片数默认 1 1 且 FULLID 才启用id_field/time_field/start_id/start_time游标字段与起点overwrite_flag覆盖式全量先 TRUNCATE 目标表再拉全表binlog_dml_types增量 DML 类型过滤INSERT,UPDATE,DELETE子集ignore_fields忽略字段按映射后目标列名匹配key 列永不被忽略dingtalk_webhook/alert_email告警通道canal_host/canal_port/canal_destination增量任务的 Canal 连接数据表一览表用途sync_datasource数据源配置密码 AES 加密sync_task任务主表sync_task_field_mapping字段映射源列 → 目标列前端 SVG 拖拽连线sync_task_progress断点进度last_sync_max_id/last_sync_timesync_task_log批次明细日志sync_task_run运行历史sync_task_verify校验运行记录进度 三类差异统计 修复结果sync_task_diff差异明细类型 / 主键 / 字段差异 / 修复状态sync_canal_positionCanal 增量位点sync_sql_log/sync_sql_favoriteSQL 执行日志 / SQL 收藏sync_audit_log字段级审计日志初始化脚本sql/datamove.sql已包含全部结构另有 9 个幂等升级脚本供老库平滑升级。十一、性能实测10 万 / 100 万行全量同步下列数字全部取自sync_task_log批次明细与sync_task_progress断点记录的真实落库数据未做估算或美化。11.1 测试环境项值MySQL8.0.41Docker 单实例innodb_buffer_pool_size128 MBinnodb_flush_log_at_trx_commit1每事务刷盘源/目标库同一实例不同库127.0.0.1:3306任务FULL · ID 模式 ·overwrite_flag1测试表12 个字段 2 个二级索引110 万行约 363 MB造数方式INSERT ... SELECT数字表笛卡尔积源库与目标库是同一个 MySQL 实例读写在同一个 buffer pool 里循环因此这是单机上限数据真实跨机同步还要再扣掉网络 RTT 与带宽开销。11.2 总耗时用例源表行数batch_size批次数同步总耗时吞吐单行耗时失败批次A10 万行100,02310,0001218.0 s5,545 行/秒0.180 ms0B100 万行1,100,023100,00013204.1 s5,391 行/秒0.186 ms0批次明细节选批次用例 A区间 → 本批行数 / 耗时用例 B区间 → 本批行数 / 耗时1185 → 10,184 / 10,000 行 / 1,993 ms185 → 100,184 / 100,000 行 / 17,412 ms540,185 → 50,184 / 10,000 行 / 1,813 ms524,465 → 624,464 / 100,000 行 / 17,779 ms980,185 → 90,184 / 10,000 行 / 1,732 ms1,048,745 → 1,148,744 / 100,000 行 /23,794 ms末批-1100,185 → 100,207 / 23 行 / 8 ms1,410,885 → 1,410,907 / 23 行 / 9 ms末批100,207 → 100,207 / 0 行 / 4 ms1,410,907 → 1,410,907 / 0 行 / 5 ms合计12 批 / 100,023 行 / 18,038 ms13 批 / 1,100,023 行 / 204,081 ms11.3 结论线性度好百万行没有劣化数据量放大 11 倍单行耗时仅从 0.180 ms 涨到 0.186 ms3.3%吞吐只降 2.8%。说明setFetchSize流式读取生效全程无 OOM、无 GC 抖动。batch_size不是吞吐的决定因素1 万 vs 10 万吞吐几乎不变瓶颈在单行写目标表2 个二级索引维护 每批 commit 刷盘。既然对吞吐不敏感就按失败回滚代价选 ——推荐 1000 ~ 10000。单批耗时存在约 30% 的离群点用例 B 第 9 批 23,794 ms比中位数高 30%属 InnoDB checkpoint 刷脏页 二级索引页分裂造成的周期性写放大抖动不是引擎逻辑问题。结果零差异同步后逐行比对两侧均 1,100,023 行字段不一致 0 行、目标缺失 0 行。线性外推按 5,391 行/秒1000 万行单线程约需 1855 秒约 31 分钟—— 这正是分片并行存在的意义切 8 片理论上可压进 5 分钟量级。十二、10 分钟跑起来1) 环境JDK 1.8、Maven 3.6、MySQL 5.7 / 8.x增量同步另需 Canal Server。2) 初始化数据库mysql -uroot -p sql/datamove.sql老库升级幂等脚本按日期顺序执行mysql -uroot -p datamove sql/upgrade_20260921_alert_email.sql mysql -uroot -p datamove sql/upgrade_20260921_field_mapping.sql # ... 其余脚本见 sql/ 目录3) 启动后端mvn clean spring-boot:run # 或打包mvn package -DskipTests后端 APIhttp://localhost:8080/Swaggerhttp://localhost:8080/swagger-ui/index.html4) 启动前端cd ruoyi-ui npm install npm run dev打开http://localhost:80默认账号admin / admin123。5) 三步跑通第一个任务数据源添加源库与目标库 → 点「测试连接」任务新建任务全量 / 增量 / 表结构填批次大小、分片数、忽略字段、告警 Webhook → 点「启动」看效果任务大盘看实时速率与瓶颈同步日志看批次明细跑完用「数据校验」对账并一键补齐差异。增量任务前置条件源库开启 ROW 模式SET GLOBAL binlog_format ROW;然后在 Canal 的instance.properties中指向源库创建增量任务时填 Canal Host / Port默认 11111/ Destination。十三、写在最后DataMove 的定位很清晰给中小企业 / 个人 / 外包团队用的零代码 MySQL 同步工具—— 配好就能跑跑完能对账出错有告警动过有审计。核心同步能力全量 / 增量 / DDL / 分片 / 字段映射 / 校验修复 / 大盘全部开源采用 Apache-2.0 协议个人与公司均可免费商用企业级特性脱敏、多租户细粒度权限、集群监控、SLA 支持为商业版范畴详见仓库LICENSING.md。如果这个项目对你有帮助欢迎 Star 鼓励Giteehttps://gitee.com/qingtian2023/datamove文档README.mdQUICKSTART.md Swagger UI反馈欢迎 Issue / PR作者DataMove 团队
返回列表