简介:DBSyncer(简称dbs)是一款开源的数据同步中间件,面向需要跨库、跨源数据流转的开发者与运维人员,解决MySQL、Oracle、SqlServer、PostgreSQL、Elasticsearch、Kafka、File、SQL等多种异构数据源之间的全量与增量同步问题。资源包共737个文件,以472个Java源码为核心,辅以html、css、js等前端页面资源,以及xml、sql、json、sh、bat等配置与脚本文件,另有少量字体、图片等静态素材,压缩包约2.07MB,结构完整,便于二次开发与本地部署调试。该中间件支持上传插件自定义同步转换业务,并提供全量与增量数据统计图、应用性能预警等监控能力,读者可借此研究同步任务的调度机制、插件扩展方式与监控指标实现,也可作为数据库管理监控与数据同步场景的工程参考。目前已有747人学习下载,适合具备一定Java与数据库基础、希望深入理解数据同步中间件设计的中高级开发者。
1. 数据同步中间件选型:为什么异构数据源统一同步是个真问题
凌晨两点被电话叫醒,MySQL 到 Kafka 的同步链路断了,下游实时看板全部停更。这种场景做过数据集成的工程师都不陌生。业务系统里同时跑着 MySQL、Oracle、SqlServer、PostgreSQL,日志侧要进 Kafka,文件侧要落对象存储或 FTP,分析侧还要灌到 ClickHouse 或 Doris。每加一条链路就写一套定时脚本,每换一个源端就重写一遍连接逻辑,维护成本指数级上升。一款开源的数据同步中间件要解决的正是这个问题:用统一的抽象层把异构数据源的读取、转换、写入标准化,让 MySQL、Oracle、SqlServer、Postgre、File、Kafka、SQL 这些场景共用一套配置和运行时。它适合正在被多源同步折磨的数据平台工程师、中间件开发者,以及需要快速搭建同步链路的团队。读完你能判断这个方向值不值得投入,也能照着把最小链路跑通。
2. 拆解同步中间件的核心抽象:Reader、Writer、Channel 怎么分工
2.1 为什么异构同步不能靠脚本堆砌
脚本堆砌的问题不在于能不能跑,而在于不可观测、不可复用、不可扩展。一条 MySQL 到 Kafka 的脚本里,连接管理、字段映射、断点续传、错误重试、限流全揉在一起,换一个源端就要复制粘贴再改。数据同步中间件的常见做法是引入三层抽象:Reader 负责从源端拉数据,Channel 负责缓冲和流控,Writer 负责写入目标端。三者通过统一的数据记录模型通信,记录里包含表名、操作类型、字段列表、时间戳。这样 MySQL Reader 和 Oracle Reader 对上层暴露的接口一致,Kafka Writer 和 File Writer 也一致。选型时要重点看这个抽象层是否干净:Reader 是否支持增量位点、Writer 是否支持批量提交、Channel 是否支持背压。如果中间件把源端特有逻辑泄漏到 Writer 层,扩展新数据源时就会很痛苦。
2.2 用配置描述一条同步链路的最小结构
一条同步链路在中间件里通常用一个 JSON 或 YAML 描述。下面是一个最小示例,把 MySQL 的增量数据同步到 Kafka:
{ "job": { "name": "mysql_to_kafka_orders", "reader": { "type": "mysql", "connection": { "host": "127.0.0.1", "port": 3306, "database": "shop", "username": "sync_user", "password": "sync_pass" }, "table": "orders", "mode": "incremental", "position": { "type": "binlog", "start_file": "mysql-bin.000001", "start_pos": 4 }, "columns": ["id", "user_id", "amount", "status", "created_at"] }, "channel": { "type": "memory", "capacity": 10000, "batch_size": 500 }, "writer": { "type": "kafka", "connection": { "bootstrap_servers": "127.0.0.1:9092", "topic": "shop_orders" }, "format": "json", "acks": "1" } } }这段配置里,reader.type 决定用哪个源端插件,mode 为 incremental 时走 binlog 增量,position 记录起始位点。channel.capacity 是内存队列容量,batch_size 是攒批大小,直接影响吞吐和延迟。writer.acks 设为 1 表示 Kafka leader 写入即返回,追求吞吐可以设 1,追求可靠设 all。参数怎么调后面章节会展开,这里先建立结构认知:中间件的配置就是围绕 Reader、Channel、Writer 三段填参数。
2.3 增量同步的位点管理是可靠性的命门
全量同步简单,增量同步才是生产环境的常态。增量同步的核心是位点管理:记录上次同步到哪里,重启后从位点继续。MySQL 用 binlog file 和 position,Oracle 用 SCN,SqlServer 用 LSN,PostgreSQL 用 LSN 或 slot。中间件需要把位点持久化,常见做法是存到本地文件、ZooKeeper 或源端的一张元数据表。位点提交时机很关键:如果先提交位点再写目标端,中间崩溃会丢数据;如果先写目标端再提交位点,中间崩溃会重复。多数中间件选择至少一次语义,配合目标端幂等写入来去重。选型时要确认位点存储是否支持高可用,单机文件位点在容器漂移后会失效。我一般会把位点存到源端同库的元数据表,跟着源端备份走,恢复时不会丢。
3. 从零跑通 MySQL 到 Kafka 的增量同步链路
3.1 环境准备与源端 MySQL 的 binlog 配置
先确认 MySQL 开了 binlog 且格式为 ROW。登录 MySQL 执行:
SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format'; SHOW VARIABLES LIKE 'binlog_row_image';如果 log_bin 是 OFF,需要在 my.cnf 里加:
[mysqld] server-id=1 log-bin=mysql-bin binlog_format=ROW binlog_row_image=FULL expire_logs_days=7改完重启 MySQL。binlog_format 必须是 ROW,因为中间件解析的是行级变更,STATEMENT 格式拿不到变更前后的完整字段。binlog_row_image 设为 FULL 保证 update 语句能拿到所有列的新值。expire_logs_days 控制 binlog 保留天数,设太短会导致位点过期后无法续传,设太长会占磁盘,一般 7 天起步。然后创建同步账号并授权:
CREATE USER 'sync_user'@'%' IDENTIFIED BY 'sync_pass'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'sync_user'@'%'; FLUSH PRIVILEGES;REPLICATION SLAVE 权限用于读 binlog,REPLICATION CLIENT 用于查位点。只给 SELECT 是不够的,这是新手常踩的坑。
3.2 目标端 Kafka 的 topic 与分区规划
Kafka 侧先建 topic。分区数决定并行消费能力,一般按目标端写入并行度来定:
kafka-topics.sh --create \ --bootstrap-server 127.0.0.1:9092 \ --topic shop_orders \ --partitions 6 \ --replication-factor 2分区数建议是 Writer 并行线程数的整数倍,6 个分区配 3 个写线程比较顺。replication-factor 生产环境至少 2,单副本挂一台 broker 就丢数据。消息格式用 JSON 时,建议在消息头里带上源端表名和操作类型,方便下游按表路由。如果下游是 Flink 或 Spark,可以直接消费 JSON 解析;如果下游是另一个数据库,中间件通常还支持 Avro 或 Protobuf 格式,序列化开销更小但需要 schema registry。我一般先用 JSON 跑通,压测后再换二进制格式。
3.3 启动同步任务并验证数据一致性
配置文件和依赖就绪后,启动任务:
./bin/sync-launcher.sh --config conf/mysql_to_kafka_orders.json --job mysql_to_kafka_orders启动后看日志确认三件事:Reader 是否成功连上 MySQL 并拿到起始位点,Channel 是否有数据流入,Writer 是否成功发送到 Kafka。验证数据用 Kafka 控制台消费者:
kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 \ --topic shop_orders --from-beginning --max-messages 10然后在 MySQL 里插一条测试数据:
INSERT INTO orders (user_id, amount, status, created_at) VALUES (1001, 99.50, 'paid', NOW());再消费一次,看是否出现对应 JSON。一致性验证要做双向:源端插、改、删各一条,确认目标端消息的操作类型分别是 INSERT、UPDATE、DELETE。如果 UPDATE 消息里只有变更列没有全列,检查 binlog_row_image 是否为 FULL。如果 DELETE 消息只有主键,那是正常的,下游按主键删即可。
4. 多源适配的坑:Oracle、SqlServer、Postgre 各自踩过什么
4.1 Oracle 的 SCN 与 LogMiner 配置要点
Oracle 增量同步比 MySQL 麻烦,因为需要开归档日志和补充日志。先确认归档模式:
SELECT log_mode FROM v$database; SELECT supplemental_log_data_min, supplemental_log_data_pk, supplemental_log_data_all FROM v$database;log_mode 必须是 ARCHIVELOG,否则 LogMiner 读不到变更。补充日志要开最小补充日志和主键补充日志:
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA; ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (PRIMARY KEY) COLUMNS;不开补充日志,UPDATE 和 DELETE 的日志里可能缺主键,中间件无法定位要改哪一行。SCN 位点管理上,Oracle 的 SCN 增长很快,位点表要定期清理,否则元数据表膨胀。另外 Oracle 的 LogMiner 查询对 UNDO 表空间有压力,同步任务并发高时要把 LogMiner 会话数调大,常见参数是加_logminer_max_parallelism,具体值按 CPU 核数来。踩过的坑是:归档日志被 RMAN 删了但位点还在,重启后报 ORA-01291 找不到日志,解决方法是位点表里记录日志序列号,启动前校验日志是否存在。
4.2 SqlServer 的 CDC 开启与 LSN 位点
SqlServer 走 CDC 模式。先对数据库开 CDC:
USE shop; EXEC sys.sp_cdc_enable_db; EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'orders', @role_name = NULL, @supports_net_changes = 1;开完查 CDC 实例:
SELECT * FROM cdc.change_tables;LSN 位点存在 cdc.lsn_time_mapping 里,中间件读 cdc.dbo_orders_CT 捕获表。坑在于 CDC 的清理任务默认只保留 3 天,位点超过 3 天没推进,捕获表数据被清掉就断链了。要改清理策略:
EXEC sys.sp_cdc_change_job @job_type = 'cleanup', @retention = 10080;retention 单位是分钟,10080 是 7 天。另外 SqlServer 的 CDC 对 DDL 不友好,源表加列后捕获表不会自动加列,需要重新开 CDC 并重置位点,这个操作会丢中间数据,生产环境要停业务窗口做。
4.3 Postgre 的逻辑复制槽与 WAL 保留
Postgre 用逻辑复制槽。先改 postgresql.conf:
wal_level = logical max_replication_slots = 10 max_wal_senders = 10重启后创建复制槽:
SELECT * FROM pg_create_logical_replication_slot('sync_slot', 'pgoutput');复制槽会阻止 WAL 被回收,如果同步任务停了但槽还在,WAL 会一直堆积直到磁盘满。这是 Postgre 同步最危险的坑。监控上要盯pg_replication_slots的active和restart_lsn,任务停了要手动删槽:
SELECT pg_drop_replication_slot('sync_slot');另外 Postgre 的逻辑复制默认不复制 DDL,源表加列后中间件需要重新拉 schema,否则新列写不到目标端。常见做法是配置 schema 自动刷新间隔,或者监听 DDL 事件手动触发。
5. 避坑与排查:同步链路断了的 5 个血泪现场
5.1 位点提交了但数据没到目标端
现象:任务重启后从新位点开始,但目标端缺了一段数据。原因:中间件先提交位点再异步写目标端,写失败时位点已推进。解决:改成先写目标端再提交位点,或者目标端做幂等去重配合至少一次语义。检查配置里位点提交模式,如果是 async 且没有重试队列,改成 sync 提交。
5.2 Kafka 消息延迟高但吞吐上不去
现象:消费端 lag 持续增长,Producer 端吞吐只有几 MB/s。原因:batch_size 太小、linger.ms 为 0、acks 为 all 且分区数不足。解决:batch_size 调到 16384 以上,linger.ms 设 5 到 20,acks 按可靠性要求设 1 或 all,分区数扩到 Writer 线程数的 2 倍。压测时用 kafka-producer-perf-test 先摸清单 broker 上限。
5.3 Oracle 同步报 ORA-01291 找不到归档日志
现象:任务重启后报错,日志里提示缺失某个归档日志序列。原因:RMAN 备份策略删了归档日志,但位点表里的 SCN 还指向那个日志。解决:位点表增加日志序列号和归档路径字段,启动前校验文件存在;或者把归档保留时间设得比同步最大延迟长。后悔药是定期把位点表备份,断链后能手动跳到最近的可用位点。
5.4 Postgre 复制槽导致磁盘写满
现象:数据库磁盘使用率飙升,pg_wal 目录巨大。原因:同步任务停了但复制槽 active 为 false,WAL 无法回收。解决:监控 pg_replication_slots,任务停止超过阈值就告警;确认不再需要该槽后手动删除。预防措施是给复制槽设 max_slot_wal_keep_size,超过就自动失效,但会断链,要权衡。
5.5 字段类型映射导致精度丢失
现象:MySQL 的 decimal(18,4) 同步到 Kafka 后变成科学计数法,下游解析出错。原因:中间件默认用 double 序列化 decimal,精度丢失。解决:配置里指定 decimal 用字符串格式输出,或者用 Avro 的 decimal 逻辑类型。Oracle 的 NUMBER 同理,不指定精度时默认映射成 double。字段映射表要在任务上线前逐列核对,尤其是金额、时间戳、大整数。
6. 进阶技巧:用 SQL 作为同步目标端做轻量 ETL
6.1 SQL Writer 的批量 upsert 与冲突处理
很多场景目标端不是 Kafka 而是另一个数据库,中间件的 SQL Writer 支持把变更写成 insert、update、delete 或 upsert。以 MySQL 目标端为例,upsert 语法:
INSERT INTO orders_sync (id, user_id, amount, status, created_at) VALUES (?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE user_id = VALUES(user_id), amount = VALUES(amount), status = VALUES(status), created_at = VALUES(created_at);批量提交时把多条 values 拼在一起,batch_size 设 500 到 1000 比较稳。冲突处理策略有三种:ignore 跳过、update 覆盖、error 报错。增量同步一般用 update 覆盖,保证最终一致。如果源端有物理删除,目标端要么物理删要么软删,软删需要加 deleted 标记列,中间件配置里指定 delete 转 update。
6.2 用 Channel 做流控和背压保护目标端
Channel 不只是缓冲,还能做背压。当 Writer 写入变慢,Channel 队列满,Reader 自动降速。配置里 capacity 和 batch_size 的比值决定背压灵敏度。capacity 10000、batch_size 500 时,队列能攒 20 批,Writer 短暂抖动不会影响 Reader。如果目标端是 Oracle 这种写入慢的库,capacity 调大到 50000,batch_size 降到 200,用更多批次换平稳。监控 Channel 的队列深度,持续接近 capacity 说明 Writer 是瓶颈,要加并行或优化目标端索引。
6.3 验证同步延迟的三个指标
同步延迟不能只看任务状态,要量化。第一个指标是位点延迟:源端当前位点减去中间件已提交位点,MySQL 用SHOW MASTER STATUS对比,Oracle 用当前 SCN 对比。第二个指标是端到端延迟:源端插入时间戳到目标端可见时间戳的差值,在消息里带源端时间戳,下游计算。第三个指标是 Channel 队列深度,反映瞬时积压。三个指标一起看:位点延迟大说明 Reader 慢,端到端延迟大但位点延迟小说明 Writer 或网络慢,队列深度高说明背压生效。我一般把这三个指标打到 Prometheus,配告警阈值,位点延迟超过 60 秒就查。
这套方案值不值得做,取决于你的同步链路数量。三条以内脚本能扛,五条以上中间件的抽象收益就出来了。我自己的习惯是先用最小配置跑通一条 MySQL 到 Kafka,压测到目标吞吐,再逐个接 Oracle、SqlServer、Postgre。每接一个源端,把踩过的坑写进配置模板,下次直接复用。希望帮到你。
本文还有配套的精品资源,点击获取