简介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 和 positionOracle 用 SCNSqlServer 用 LSNPostgreSQL 用 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-id1 log-binmysql-bin binlog_formatROW binlog_row_imageFULL expire_logs_days7改完重启 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 权限用于读 binlogREPLICATION 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 模式。先对数据库开 CDCUSE 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.confwal_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 到 20acks 按可靠性要求设 1 或 all分区数扩到 Writer 线程数的 2 倍。压测时用 kafka-producer-perf-test 先摸清单 broker 上限。5.3 Oracle 同步报 ORA-01291 找不到归档日志现象任务重启后报错日志里提示缺失某个归档日志序列。原因RMAN 备份策略删了归档日志但位点表里的 SCN 还指向那个日志。解决位点表增加日志序列号和归档路径字段启动前校验文件存在或者把归档保留时间设得比同步最大延迟长。后悔药是定期把位点表备份断链后能手动跳到最近的可用位点。5.4 Postgre 复制槽导致磁盘写满现象数据库磁盘使用率飙升pg_wal 目录巨大。原因同步任务停了但复制槽 active 为 falseWAL 无法回收。解决监控 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 作为同步目标端做轻量 ETL6.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 调大到 50000batch_size 降到 200用更多批次换平稳。监控 Channel 的队列深度持续接近 capacity 说明 Writer 是瓶颈要加并行或优化目标端索引。6.3 验证同步延迟的三个指标同步延迟不能只看任务状态要量化。第一个指标是位点延迟源端当前位点减去中间件已提交位点MySQL 用SHOW MASTER STATUS对比Oracle 用当前 SCN 对比。第二个指标是端到端延迟源端插入时间戳到目标端可见时间戳的差值在消息里带源端时间戳下游计算。第三个指标是 Channel 队列深度反映瞬时积压。三个指标一起看位点延迟大说明 Reader 慢端到端延迟大但位点延迟小说明 Writer 或网络慢队列深度高说明背压生效。我一般把这三个指标打到 Prometheus配告警阈值位点延迟超过 60 秒就查。这套方案值不值得做取决于你的同步链路数量。三条以内脚本能扛五条以上中间件的抽象收益就出来了。我自己的习惯是先用最小配置跑通一条 MySQL 到 Kafka压测到目标吞吐再逐个接 Oracle、SqlServer、Postgre。每接一个源端把踩过的坑写进配置模板下次直接复用。希望帮到你。本文还有配套的精品资源点击获取
