MongoDB 文档同步 SDK / 工具(Java 8+)
主路径:MongoDB → MongoDB(自建 / 云上均支持)。
同构兼容库亦可:Amazon DocumentDB、阿里云 DDS 等 MongoDB 协议兼容的文档库。
Kafka Sink:sink.type=kafka,把变更写成 mongo-kafka 兼容的 Change Stream 消息,供下游消费或再经 mongo-kafka Sink 落库。
面向迁移、灾备、多活与跨架构搬迁:一套 API / 一条命令,完成 全量 + 增量、DDL 跟随、多库表过滤 与 数据校验。可嵌入 Java 业务进程,也可脚本启动。
./bin/mongosync.sh -f conf/mongo-sync.properties # 启动同步
./bin/mongosync.sh --config conf/mongo-sync.properties --shutdown
./bin/verify.sh -f conf/mongo-verify.properties # 数据比对QQ 交流群:983986505(使用问题、需求反馈、经验交流欢迎加群)
| 诉求 | mongo-sync 怎么做 |
|---|---|
| Mongo → Mongo 主路径 | 全量∥增量、DDL、分桶有序写,覆盖迁库 / 灾备主场景 |
| Mongo → Kafka | sink.type=kafka,消息格式对齐 mongo-kafka Source(Change Stream JSON/BSON) |
| 同构上云(DocumentDB / DDS) | 标准 Mongo 驱动写入协议兼容库,便于迁云 |
| 迁库 / 扩容不停服 | 全量∥增量并行(FULL_AND_INCREMENTAL),UPSERT 兜底窗口重复 |
| 跨架构互传 | 自动识别 standalone / 副本集 / 分片,匹配读任务(Oplog / ChangeStream) |
| 分片集群增量 | mongos 拉全量 + ChangeStream@mongos(MongoDB 3.6+;不再提供多分片 OPLOG) |
| 大表全量加速 | 按 _id 切段多任务并行读 |
| 结构一起走 | 启动预建集合 / 索引;运行中 DDL(删表、改名、建删索引)可落地 |
| 写序与吞吐 | _id 分桶 + LMAX Disruptor 背压;唯一索引自动有序写 |
| 可嵌入 / 可脚本 | SDK(MongoSyncClient)或 mongosync.sh 配置文件启动 |
| 迁完可验 | verify.sh:COUNT / ID / FULL 三种比对 |
Sink 不感知 捕获协议——无论 Oplog 还是 ChangeStream,统一变成 TransferEvent / DdlEvent 再写入 Sink。
二者是姊妹产品,控制面(start / pauseIncremental / canCommit / commit、全量∥增量、分桶有序)对齐,数据面互不替代:
| mongo-sync(本仓) | rds-sync | |
|---|---|---|
| 源 | MongoDB(Oplog / ChangeStream) | MySQL / Oracle / PostgreSQL |
| Sink | MongoDB、DocumentDB / DDS、Kafka(Change Stream) | MySQL JDBC、Kafka(行级 envelope) |
| 事件契约 | TransferEvent / DdlEvent |
RowChange / DdlEvent |
| 典型场景 | 迁库、灾备、分片升级、投递 Kafka | 异构关系库搬迁、投递 Kafka、国产化转型 |
需要 MySQL / Oracle / PostgreSQL 时请走 rds-sync,不要在本仓找 JDBC 源。
- 四种同步模式:仅全量、全量∥持续增量、全量后追平再停、仅增量
- 双Sink 形态:MongoDB(默认)/ Kafka(mongo-kafka Change Stream 消息)
- 双捕获通道:ChangeStream(推荐 / MongoDB 7.0+);Oplog 3.2–6.0(V1/V2/V3 解析)
- 架构自适应:
capture.mode=AUTO按源端拓扑匹配读计划(禁止在 mongos / standalone 上误拉 Oplog) - 多库表:白/黑名单、
ns变换(MongoMultiSyncClient) - 元数据:
bootstrapCollection/bootstrapIndexes可分别开关;支持跳过 TTL 索引 - 位点:可选文件持久化(
offset.store.dir)+ 周期心跳日志 - 迁移状态机:
MigrationProgress/canCommit/commit;canCommit要求全量完成、pipeline 排空,且增量滞后 ≤commit.max.lag.ms(默认 10000) - 捕获窗口告警:全量∥增量期间监控锚定位点相对 oplog 最早条目的余量(
window.warn.seconds,默认 3600);逼近阈值打WINDOW WARN - 独立增量 pause:
pauseIncremental/resumeIncremental(全量可继续;HTTP:/api/v1/pauseIncremental) - 校验:
VerifyMain支持单表 / 多表白名单
- 能确保数据总量一致
- 能确保数据信息一致
- 能确保异构系统数据同步一致
- 能确保数据索引一致
- 能确保数据结构一致
- 全量数据复制
- 实时数据同步
- 增量数据同步
- 自定义同步范围
- 复合数据同步方案
- 100% 传输带宽利用率
- 可控 CPU 利用率
- 内存使用率可配置
- 支持多表并传
- 体积小巧
- 断点续传
- 支持多版本 MongoDB 同步
| 全量 | Oplog | ChangeStream | |
|---|---|---|---|
| standalone | ✅ | ❌ | ❌(仅 FULL) |
| 副本集 | ✅ | ✅ | ✅ |
| mongos | ✅ | ❌(改写为各 shard) | ✅ |
| 某 shard | — | ✅ | — |
cd mongo-sync
chmod +x bin/*.sh
# 编辑配置:源/Sink URI、库表或 namespace.white
cp doc/examples/mongo-sync.example.properties my-sync.properties
./bin/mongosync.sh -f my-sync.properties
# Ctrl+C 优雅停止
./bin/verify.sh -f doc/examples/mongo-verify.example.propertiesMongoSyncClient sync = MongoSyncClient.create(MongoSyncClient.builder()
.sourceUri("mongodb://src/?replicaSet=rs0")
.sinkUri("mongodb://sink")
.mapCollection("demo", "orders")
.captureMode(CaptureMode.AUTO)
.syncMode(SyncMode.FULL_AND_INCREMENTAL)
.offsetStoreDir("./data/offsets")
.writeErrorHandler((bucket, event, err) -> {
// 生产务必处理写失败
}));
sync.start();多库表:
MongoMultiSyncClient multi = MongoMultiSyncClient.create(MongoMultiSyncConfig.builder()
.sourceUri("mongodb://src/?replicaSet=rs0")
.sinkUri("mongodb://sink")
.namespaceWhite("demo;app.orders")
.syncMode(SyncMode.FULL_AND_INCREMENTAL)
.offsetStoreDir("./data/offsets")
.writeErrorHandler((bucket, event, err) -> { }));
multi.start();Kafka Sink(sink.uri 为 bootstrap servers):
MongoSyncClient.create(MongoSyncClient.builder()
.sourceUri("mongodb://src/?replicaSet=rs0")
.sinkUri("127.0.0.1:9092")
.sinkType(SinkType.KAFKA)
.mapCollection("demo", "orders")
.kafkaTopicPrefix("mongo")
.syncMode(SyncMode.FULL_AND_INCREMENTAL)
.offsetStoreDir("./data/offsets")
.writeErrorHandler((bucket, event, err) -> { }));支持多种数据同步方案,包括全量数据复制、实时数据同步、增量数据同步、自定义同步范围以及复合数据同步方案,可根据业务需求灵活选择。
提供高效数据校验功能,能确保数据总量一致、数据信息一致、异构系统数据同步一致、数据索引一致以及数据结构一致,全方位保障数据准确性。
采用高速数据同步机制,实现 100% 传输带宽利用率,支持可控 CPU 利用率,内存使用率可配置,并支持多表并传,确保同步过程高效稳定。同时支持断点续传功能,避免网络中断导致的数据丢失。
用姊妹产品 rds-sync。mongo-sync 只做文档库(MongoDB / 协议兼容库 / Kafka Sink);关系库的全量切分、JDBC 写入和 binlog/redo/WAL 增量在 rds-sync。两边的同步模式与迁移控制 API 故意对齐,便于同一套运维习惯。
| 文档 | 说明 |
|---|---|
| bin/README.md | 脚本入口详解 |
| doc/ARCHITECTURE.md | 架构、能力清单与已知限制 |
| doc/examples/mongo-sync.example.properties | 同步配置示例 |
| doc/examples/mongo-sync-kafka.example.properties | Kafka Sink配置示例 |
| doc/examples/mongo-verify.example.properties | 校验配置示例 |
| doc/oplog/ | 各版本 Oplog 样例 |
| rds-sync | 关系库同步(MySQL / Oracle / PostgreSQL) |
- MongoDB → MongoDB:迁库、扩容、跨机房 / 多活(主推)
- MongoDB → Kafka:变更投递到 Kafka,格式对齐 mongo-kafka Source
- 同构上云:自建 Mongo → DocumentDB / DDS 等协议兼容库
- 灾备与只读副本:持续增量同步到备端
- 架构升级:副本集 ↔ 分片、跨版本(捕获通道随版本自动收紧)
- 业务内嵌同步:以 SDK 嵌入现有 Java 服务,统一事件模型
- 关系库搬迁:不在本仓;见 rds-sync
生产请配置
offset.store.dir、注入SyncWriteErrorHandler,切换前用verify.sh抽检。DocumentDB / DDS 与社区版在部分算子、DDL 上可能有差异,迁云前务必验证。更细限制见 架构说明。
- QQ 交流群:983986505