一个基于 RuoYi-Vue 二次开发的轻量 MySQL 同步平台:把「配 DataX JSON / 手写 Canal 消费端 / SSH 隧道跑 SQL」这三件苦差事,收敛成浏览器里的点选操作。

本文不讲空话:架构设计、核心代码实现与性能实测数字,全部来自项目当前代码与运行日志。

目录


一、为什么要造这个轮子?

团队过去三年在数据同步上反复踩坑:

场景常见做法真实痛点
业务库拆分(订单库 → 订单中台库)DataX 写 JSON字段映射手写 joiner,运维同学上手成本高,改一个列名要动配置
实时订阅 binlog部署 Canal + 自写消费端要自己处理位点、多线程 ACK、异常回滚,错一次就丢/重复数据
上线前对账两边各 count(*) 比总数只知道"差 37 行",不知道是哪 37 行、哪个字段,修不了
临时数据修复登录跳板机手写 INSERT ... SELECT无审计、无日志,出事没人知道谁改的
DBA 离职交接脚本散落在各台机器新人接手无从下手

于是有了 DataMove,目标很明确:把"同步"变成"配置 + 一键启动",把"对账"变成"看得见差异 + 一键补齐"。


二、能力总览

┌─────────────────────────────────────────────────────────────────────┐
│                          DataMove 能力地图                           │
├──────────────┬──────────────────────────────────────────────────────┤
│ 数据源        │ MySQL 5.7 / 8.x,多源并存,密码 AES 加密存储            │
│ 全量同步      │ 按主键 ID / 按时间字段,断点续传 + 幂等覆盖              │
│ 分片并行      │ FULL+ID 模式按主键区间切 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.13(RuoYi-Vue 4.8.1)
同步原生 JDBC 分批 + Canal Client 1.1.4(增量)
校验源/目标双游标流式归并(setFetchSize(Integer.MIN_VALUE),内存只驻留一行)
加密AES(AES/CBC/PKCS5Padding + Base64)
调度Spring @Scheduled + 自建线程池
告警钉钉 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();                                // ❌ 失败不推进,下次从断点重来
        // 记日志 + 告警,不静默吞异常
    }
}

配套的两个硬要求:

  1. 幂等写入:用 INSERT ... ON DUPLICATE KEY UPDATE,而不是 REPLACE INTOREPLACE 内部是先 DELETE 再 INSERT,会改自增主键、触发级联删除;ON DUPLICATE KEY UPDATE 只更新既有列,重跑不脏数据。
  2. 流式读取:MySQL JDBC 默认会把结果集全部拉进内存,必须显式 setFetchSize(batchSize),否则百万行必 OOM。

4.3 分片并行:把 31 分钟压到几分钟

FULL + ID 模式下,任务配置 shard_count > 1 时会走分片:

  1. 查 SELECT MIN(id), MAX(id)
  2. RangeSplitter.split(min, max, shardCount) 把主键区间等宽切成 N 段(纯函数,好测);
  3. 每个分片一个线程 + 一对独立的源/目标连接,各自 PreparedStatement 与 setFetchSize,区间互不重叠;
  4. 进度按 synchronized(progressLock) 聚合后回写 sync_task_progress,全局游标取各分片最大值。

与断点续传的互斥规则(很重要,否则会重复同步):已有断点(progress.total_rows > 0)时自动回退单线程续传;空表、区间划分失败、只能切出 1 段时也回退。回退时会打一行日志"分片并行回退单线程续传",避免出现"看着是分片任务,实际在重复搬数据"的灵异现象。

4.4 batch_size 和 shard_count 是两回事

维度批次 batch_size分片 shard_count
切分方向时间维度"分次"空间维度"分区"
工作方式单线程按 LIMIT 游标循环,一批 commit 完再读下一批整表按主键区间切 N 块,N 个线程同时各管一块
目的控内存、缩短事务、失败只重试一小批提升并行吞吐(耗时近线性 / N)
推荐值1000 ~ 100001 ~ 16(参考目标库写入能力)
串行 (shard_count=1):
  线程1: [批次1]→[批次2]→[批次3]→...→[批次N]        总时间 T

分片并行 (shard_count=4):
  线程1: [批次1]→[批次5]→...   ┐
  线程2: [批次2]→[批次6]→...   ├─ 同时进行           总时间 ≈ T/4
  线程3: [批次3]→[批次7]→...   ┘
  线程4: [批次4]→[批次8]→...

4.5 暂停 / 停止怎么做

用 SyncContext 持有 AtomicBoolean pauseFlag / stopFlag,主循环(以及每个分片线程内循环)在批次边界检查标志位并优雅收口;暂停会先把当前批次写完再停,因此不会产生半批脏数据。任务级单实例保护用 static Set<Long> RUNNING_TASK + synchronized start(),同一任务不会有两个 worker。


五、核心设计 2:Canal 增量同步的三层过滤与 ACK 顺序

5.1 主流程

connect() → subscribe(源库\.任务表) → while(true) { getWithoutAck(1000) }
                                              ↓
                                    apply(entry) → 目标库 JDBC 执行
                                              ↓
                                   connector.ack(batchId)   // ✅ 成功才 ACK
                                              ↓
                                   upsert sync_canal_position

5.2 三层事件过滤,让无关事件根本不进客户端

手段效果
库/表过滤Canal 服务端 subscribe 表达式从 .*\..* 收紧为 源库\.任务表(库表名做正则转义)不相关库表的 binlog 事件根本不进客户端,省网络与解析开销
DML 类型过滤任务配置 binlog_dml_types(INSERT / UPDATE / DELETE 按需勾选),未勾选事件计数丢弃归档库只收 INSERT、审计库不要 DELETE 等场景
字段过滤任务配置 ignore_fields,拼 SQL 时自动剔除这些列目标库自带的 create_time / tenant_id 不被源库覆盖

几个实现细节值得说明:

  • key 列永远保留:即使把主键填进 ignore_fields 也会被跳过,否则 UPDATE / DELETE 就没有 WHERE 条件了。
  • 忽略字段按映射后的目标列名匹配:配了字段映射也不用改这里。
  • 忽略字段同时被校验引擎读取update_timeON 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_ONLY + setFetchSize(Integer.MIN_VALUE)百万行级别内存里只驻留一行
进度回填频率每 5000 行抽屉里 2s 轮询展示已比对行数 / 三类差异数
差异明细留存上限2000 条超出置 truncated=1 并明确提示"仅保留前 N 条",但统计数字始终全量准确
修复批大小200 条 / 批每条回填修复状态(已修复 / 失败 / 跳过)
侧效应不改 sync_task.status校验是只读旁路,避免"任务已完成却显示运行中"的状态错乱

两条设计红线:

  1. 不会删除目标库任何数据。 EXTRA 行只在明细里标出来、修复状态直接置为 SKIPPED。删除是不可逆的破坏性动作,不应该藏在"一键"按钮背后。
  2. 修复只改不一致的列。 依据 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(带 BOM,Excel 打开中文不乱码)。
  • 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),一条失败消息同时投递:

  1. 任务上配置的钉钉机器人 dingtalk_webhook
  2. 任务上配置的 alert_email(SMTP 服务器为全局配置)。

任一渠道未配置自动跳过,互不影响。触发场景覆盖:启动失败(源/目标库连接失败)、全量同步异常、增量批次异常、增量连接异常、表结构同步失败、校验发现差异、校验失败、修复失败。

工程细节:告警走独立单线程池(队列 200,DiscardOldestPolicy 丢弃最旧消息),异步发送不阻塞同步线程;任何告警异常只记日志,不影响任务状态流转 —— 告警是附属功能,不能反过来拖垮主链路。


十、配置项与表结构速查

sync_task 中与同步行为相关的字段:

字段含义
task_typeFULL / INCR / DDL
sync_modeID / TIME / BINLOG / DDL
batch_size批次大小,默认 1000
shard_count分片数,默认 1(> 1 且 FULL+ID 才启用)
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_time
sync_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.41(Docker 单实例)
innodb_buffer_pool_size128 MB
innodb_flush_log_at_trx_commit1(每事务刷盘)
源/目标库同一实例不同库,127.0.0.1:3306
任务FULL · ID 模式 · overwrite_flag=1
测试表12 个字段 + 2 个二级索引,110 万行约 363 MB
造数方式INSERT ... SELECT 数字表笛卡尔积

源库与目标库是同一个 MySQL 实例,读写在同一个 buffer pool 里循环,因此这是"单机上限"数据;真实跨机同步还要再扣掉网络 RTT 与带宽开销。

11.2 总耗时

用例源表行数batch_size批次数同步总耗时吞吐单行耗时失败批次
A:10 万行100,02310,0001218.0 s5,545 行/秒0.180 ms0
B:100 万行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 ms
540,185 → 50,184 / 10,000 行 / 1,813 ms524,465 → 624,464 / 100,000 行 / 17,779 ms
980,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 ms

11.3 结论

  1. 线性度好,百万行没有劣化:数据量放大 11 倍,单行耗时仅从 0.180 ms 涨到 0.186 ms(+3.3%),吞吐只降 2.8%。说明 setFetchSize 流式读取生效,全程无 OOM、无 GC 抖动。
  2. batch_size 不是吞吐的决定因素:1 万 vs 10 万,吞吐几乎不变,瓶颈在单行写目标表(2 个二级索引维护 + 每批 commit 刷盘)。既然对吞吐不敏感,就按"失败回滚代价"选 —— 推荐 1000 ~ 10000
  3. 单批耗时存在约 30% 的离群点:用例 B 第 9 批 23,794 ms,比中位数高 30%,属 InnoDB checkpoint 刷脏页 + 二级索引页分裂造成的周期性写放大抖动,不是引擎逻辑问题。
  4. 结果零差异:同步后逐行比对,两侧均 1,100,023 行,字段不一致 0 行、目标缺失 0 行。
  5. 线性外推:按 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
  • 后端 API:http://localhost:8080/
  • Swagger:http://localhost:8080/swagger-ui/index.html

4) 启动前端

cd ruoyi-ui
npm install
npm run dev

打开 http://localhost:80,默认账号 admin / admin123

5) 三步跑通第一个任务

  1. 数据源:添加源库与目标库 → 点「测试连接」;
  2. 任务:新建任务(全量 / 增量 / 表结构),填批次大小、分片数、忽略字段、告警 Webhook → 点「启动」;
  3. 看效果:任务大盘看实时速率与瓶颈,同步日志看批次明细,跑完用「数据校验」对账并一键补齐差异。

增量任务前置条件(源库开启 ROW 模式):

SET GLOBAL binlog_format = 'ROW';

然后在 Canal 的 instance.properties 中指向源库,创建增量任务时填 Canal Host / Port(默认 11111)/ Destination。


十三、写在最后

DataMove 的定位很清晰:

给中小企业 / 个人 / 外包团队用的"零代码 MySQL 同步工具" —— 配好就能跑,跑完能对账,出错有告警,动过有审计。

核心同步能力(全量 / 增量 / DDL / 分片 / 字段映射 / 校验修复 / 大盘)全部开源,采用 Apache-2.0 协议,个人与公司均可免费商用;企业级特性(脱敏、多租户细粒度权限、集群监控、SLA 支持)为商业版范畴,详见仓库 LICENSING.md

如果这个项目对你有帮助,欢迎 Star 鼓励:

  • Gitee:https://gitee.com/qingtian2023/datamove
  • 文档:README.md + QUICKSTART.md + Swagger UI
  • 反馈:欢迎 Issue / PR

作者:DataMove 团队

Logo

一站式 AI 云服务平台

更多推荐