前阵子帮一个做电商数据分析的团队处理数据同步,他们在 MongoDB 里存了用户、订单、商品、库存等七八个集合,这些数据要求实时进 ClickHouse 做 OLAP 分析。最早是每个集合单独写一个同步脚本,MongoDB 侧一有新集合上线就要再开发一轮,麻烦不说,凌晨丢变更、断点续传、全量补数全靠手工。后来我把整条链路收敛到 Flink 1.20 上,用 Flink CDC 的 MongoDB connector 一次性订阅多个集合,通过路由规则映射到 ClickHouse 的对应表,一个 pipeline 跑完所有集合的同步,新集合上线只需要在配置里加一行路由。
这套方案对实时数仓入门、数据平台工程师、以及那些手头一堆 MongoDB 集合要往 ClickHouse 灌数据的人都适用。不需要专门维护 Kafka Connect 集群,也不用自己写消费 Change Streams 的 Java 任务,Flink 1.20 加上 Flink CDC 的 MongoDB connector 就能把全量快照和增量变更统一处理掉。下面我把整个方案的选型逻辑、配置细节、实测结果和踩坑记录完整展开,照着走基本能复现。
1. 为什么是 Flink:多集合同步的方案选型思考
1.1 同步选型的几条路线对比
多集合实时同步这件事,最开始摆在我面前的有好几条路,各有各的坑。最朴素的做法是自己写一个 Java 微服务,订阅 MongoDB 的 Change Streams,然后在内存里做数据转换,再批量写入 ClickHouse。这条路看起来灵活,但实际一深入就发现工作量全在细节里:每个集合要自己维护断点续传的 offset,MongoDB 重启之后 Change Stream 的 resume token 怎么存、丢了怎么办,全量数据补数时又要单独写一套扫描逻辑,更别说表结构变更、网络抖动重试这些边角问题。一套下来少说一个多月,还不算后续维护成本。
另一条常见路线是 Debezium + Kafka Connect 把数据推到 Kafka,再让 Flink 从 Kafka 消费写入 ClickHouse。这个方案生态成熟,但链路太长,中间多了一个 Kafka 集群要养活,配置、监控、容错都翻倍。如果团队里本来就没有 Kafka,为了同步数据硬上一个消息队列,性价比真的很低。而且 Kafka Connect 的多表路由配置也不轻松,每个 collection 依然要单独定义 connector。
用 Flink 就不太一样。Flink CDC 项目本身把"全量快照 + 增量变更"这个能力内建到了 connector 里,MongoDB 的 connector 可以直接从全量快照切到 Change Streams 的增量监听,中间由 Flink 的 checkpoint 机制兜底,任务挂掉能从最近一次一致状态恢复。更关键的是,一个 Flink 任务里可以同时订阅多个 MongoDB 集合,每个集合映射到 ClickHouse 的一张表,这就是"一套配置搞定多集合"的核心基础。对我这种要管一堆同步任务的人来说,一个任务取代七八个脚本,光是告警维护就轻松太多。
1.2 Flink CDC 的核心思路:全量快照加增量变更流
很多人第一次接触 Flink CDC 会以为它只是个"实时监听 binlog/oplog 的工具",其实不是。它真正的核心是"按表快照 + 增量订阅"一体化的设计。以 MongoDB 为例,MongoDB 官方从 3.6 开始支持 Change Streams,但它有个特点:Change Streams 只能读到开启监听之后的新变更,历史存量数据读不到。所以做同步必须分两步,先全量 dump 一遍存量数据,再切换到 Change Streams 监听增量。
Flink CDC 的 MongoDB connector 把这两步合并成了一个无缝的过程。任务启动后,connector 会先对目标集合执行全量扫描,扫描期间产生的增量变更也不会丢,Flink 会把这两个过程通过分布式快照机制关联起来。全量扫完之后自动切换到 Change Streams,继续读增量。整个过程对下游 ClickHouse 来说就是"数据一直在往里面流",不会出现存量全量没进、增量又不给读的空窗期。这对那些需要从零开始建立 ClickHouse 分析表的场景,省掉了非常多的编排工作。
2. 环境准备:把 Flink 1.20 的底座搭稳
2.1 版本匹配关系一定要先确认
Flink CDC 这个项目的版本和 Flink 主版本之间有明确的对应关系,这是最容易踩坑的地方。有人直接下了最新版 Flink CDC 去配合老版本 Flink,结果启动时报了一堆 NoSuchMethodError。我用的是 Flink 1.20,对应的 Flink CDC 需要选 3.2 及以上版本,这个组合在官方文档里有明确的支持矩阵,我当时实测下来也是稳定的。MongoDB 那边,Change Streams 功能要求数据库必须是副本集或分片集群,单节点实例不支持,这一点在选型阶段就要确认好,我见过不少人栽在这上面。
ClickHouse 侧的写入组件,官方 Flink 和 Flink CDC 没有内置标准 sink,所以我用的是社区广泛使用的 clickhouse-flink-connector,配合 Flink 的 JDBC 连接能力写入。这里提醒一句,clickhouse-flink-connector 有不同版本,注意它的 Flink 适配版本,最好选和你 Flink 大版本匹配的构建,否则会出现序列化器冲突。我这个场景里,把所有 jar 包放到 Flink 的 lib 目录后,重启 Flink 集群才生效,这点需要注意,不是放进去立刻就能用。
2.2 Flink 环境安装与依赖部署
环境我用的是 standalone 集群模式,两台机器组了一个小集群。Flink 1.20 解压后修改 conf/flink-conf.yaml 里的 jobmanager.memory.process.size 和 taskmanager.memory.process.size,我给 JobManager 留了 2GB,TaskManager 每台 4GB,因为要处理多个集合的快照扫描,内存给太小容易在 checkpoint 阶段频繁 Full GC。
依赖 jar 处理好坏直接决定你能不能少折腾一个晚上。核心需要准备这几个组件:flink-cdc-pipeline-connector-mongodb(或 Flink SQL 方式下的 mongodb-cdc connector)、clickhouse-flink-connector、以及 ClickHouse JDBC driver。把这些 jar 统一放到 $FLINK_HOME/lib 下,然后重启集群。如果用 flink-cdc 自带的提交脚本,那还需要把 flink-cdc-pipeline-connector-clickhouse 这类 pipeline connector 也放进去。我当时图省事,一开始只放了 MongoDB connector,结果 ClickHouse sink 一直识别不了,报 ClassNotFoundException,找了一圈才发现是 jar 没放全。
3. 一套配置的核心:路由规则与 YAML pipeline
3.1 YAML pipeline 整体结构
Flink CDC 从 3.x 开始提供了 YAML pipeline 的提交方式,这是我说的"一套配置"的真正载体。它把 source、sink、路由规则、并行度、checkpoint 全部写在一个 YAML 文件里,提交之后 Flink CDC 会自动把任务翻译成 Flink 作业去跑。相比以往用 Java/SQL 每条线写一段,YAML 方式的维护成本低非常多,尤其适合多集合同步的场景。
我这里给出一个实际跑通的配置骨架:
source: type: mongodb hosts: mongo01:27017,mongo02:27017 username: flink password: flink123 database: app_db collection: users,orders,products,inventory batch.size: 2048 poll.await.timeout.ms: 1000 heartbeat.interval.ms: 1000 sink: type: clickhouse jdbc-url: jdbc:clickhouse://clickhouse01:8123 username: default password: clickhouse123 table.prefix: ods_ batch.size: 2000 flush.interval-ms: 1000 pipeline: parallelism: 2 checkpoint.interval: 3000 route: - source-table: app_db.users sink-table: ods.users - source-table: app_db.orders sink-table: ods.orders - source-table: app_db.products sink-table: ods.products - source-table: app_db.inventory sink-table: ods.inventory这段配置里最核心的就是 route 列表。source 侧一次性声明了要订阅的多个集合,sink 侧通过 table.prefix 约定目标表前缀,route 里把 MongoDB 的每个集合一一对应到 ClickHouse 的表。新集合上线时,我只需要在 source 的 collection 里加一个名字,再在 route 里补一行,不需要写任何代码。
3.2 多集合映射与字段自动推导
如果对比传统方式,你会发现多集合同步的复杂度大部分来自"每个集合的表结构都不一样"。MongoDB 是 schemaless 的,不同集合的字段千差万别,有的字段在这个集合里没有、在那个集合里有,Flink CDC 的 MongoDB connector 会自动做 schema 推导,把 BSON 文档转成 Flink 的 RowData。
以 users 和 orders 两个集合举例。users 集合有 _id、name、email、created_at,orders 集合有 _id、order_no、user_id、total_amount、created_at。connector 在快照阶段扫描集合时,会动态识别每个集合的字段,生成对应的 Flink 表结构,然后通过 route 规则把数据分别写入 ClickHouse 的 ods.users 和 ods.orders 表。这比传统 JDBC 同步那种"一张表一个 SQL 映射"的方式灵活太多,因为 MongoDB 字段增减时,只要 ClickHouse 表里有对应列,数据就能正常落库;ClickHouse 没有的列,也可以在 sink 配置里做忽略或类型转换。
这里要补充一个实战细节:MongoDB 的 _id 默认是 ObjectId 类型,同步到 ClickHouse 时会被转成字符串。Flink CDC 内部默认把 _id 序列化成 24 位的十六进制字符串,ClickHouse 侧对应字段建议用 String 类型接收。如果你在 ClickHouse 里建表时把 _id 定义成了 FixedString,要注意长度必须是 24,否则写入时会报数据长度不符,我当时就因为这个浪费了半小时。
3.3 关键参数解释与调优逻辑
配置里的参数看起来多,真正决定任务稳定性的就几个。batch.size 控制 MongoDB connector 批量拉取记录的条数,我设成 2048 是经验值,太小会导致频繁的网络请求,太大在碰到大文档时容易内存抖动。poll.await.timeout.ms 是 Change Streams 轮询的空闲等待时间,设 1000 毫秒可以让延迟控制在秒级,又不至于把 CPU 空转在轮询上。heartbeat.interval.ms 这个参数容易被忽略,它的作用是让 connector 在没有业务数据时依然产生心跳记录,避免 Change Streams 的空闲连接被 MongoDB 服务端回收,尤其是 MongoDB 4.2 以上版本对空闲游标有清理逻辑,不配这个参数,任务长时间没数据时会突然断流。
parallelism 我设了 2,这是根据这里 Mongo 集合的数据量定的。多集合同步时并行度不是越大越好,每个并行子任务会分担一部分集合的读取和写入。如果集合多、数据分散,可以适当调高到 4 或 6,但要注意 ClickHouse 侧并发写入的压力,以及 checkpoint 时 state 的大小。checkpoint.interval 我设 3000 毫秒,兼顾了恢复粒度和写入性能,如果业务对延迟不敏感,可以放到 5 秒甚至 10 秒,减少 checkpoint 对吞吐的影响。
4. MongoDB 侧的准备与权限
4.1 Change Streams 的前提条件
标题里说"一套配置搞定",但 MongoDB 侧的前置条件如果不满足,Flink 这边配置写得再漂亮也白搭。Change Streams 要求 MongoDB 运行在副本集或者分片集群模式下,这是 Mongo 官方从 3.6 开始的设计约束。副本集不必是生产级的大集群,一个主节点加一个从节点也能开启 Change Streams,但不要用单实例开发模式的 mongod 去测试,否则 Flink 任务会抛出连接不支持 Change Streams 的错误。
另外,Change Streams 依赖 oplog,oplog 的大小决定了能回溯的变更窗口。如果 oplog 太小,Flink 任务在暂停恢复时,resume token 可能已经超出了 oplog 的范围,任务会报错要求重新初始化。我对接的这个电商场景,oplog 大小设的是 8GB,实测能支撑大约 12 小时的变更回溯。如果你的同步链路偶尔会停半天以上,建议把 oplog 调大,或者在恢复时接受重新做一次全量快照。
4.2 账号权限与连接配置
Flink CDC 读取 MongoDB 需要一个有权限的账号。这个账号不需要管理员权限,但至少要拥有目标数据库的 changeStream 和 read 权限,同时能读 local 库的 oplog(在部分版本和配置下需要)。我给达的授权语句类似这样:
db.getSiblingDB("admin").createUser({ user: "flink", pwd: "flink123", roles: [ { role: "read", db: "app_db" }, { role: "read", db: "local" } ] });然后是授予 Change Stream 权限。在 MongoDB 里,changeStream 权限是独立 action,可以通过自定义角色授予:
db.getSiblingDB("admin").createRole({ role: "flink_cdc_role", privileges: [ { resource: { db: "app_db", collection: "" }, actions: ["changeStream", "find", "listCollections"] } ], roles: [] });建好角色后,把这个角色授给 flink 用户。实际使用中,如果你发现任务能全量读数据但一到增量阶段就报错,十有八九是 changeStream 权限没配好。我这边最开始直接用了 read 角色,全量阶段一切正常,切增量后日志里报 not authorized on admin,一下就看出来是权限问题。
5. ClickHouse sink:写入链路与建表策略
5.1 ClickHouse 写入方式选择
ClickHouse 这边没有 Flink 官方 connector,现实中可以选 JDBC 驱动直写,也可以用社区封装的 clickhouse-flink-connector。我的选择是社区 connector,因为它在批处理、异步提交、失败重试上做了封装,比裸 JDBC 好用。connector 的核心写入逻辑是按批次攒数据,达到 batch.size 或 flush.interval 阈值后,一次性批量发送给 ClickHouse,这正好契合 ClickHouse 列式存储的批量写特性。
批量参数需要结合实际数据量调。batch.size 设 2000,flush.interval-ms 设 1000,意思是每秒或者每攒够 2000 条就刷一次。如果数据量大,可以把 batch.size 提到 5000 甚至 10000,ClickHouse 单次插入的吞吐会更好,但内存占用也会同步上升,TaskManager 给不够会出现 OOM。我实测下来,2000 条一批写入延迟已经能控制在 1 秒左右,对实时分析够用了。
5.2 字段类型映射对照
MongoDB 的 BSON 类型和 ClickHouse 的类型体系差别很大,好在 Flink CDC 已经做了大部分映射。实际建表时我整理了这样一张映射对照,方便团队新人直接参考:
| MongoDB BSON 类型 | Flink 类型 | ClickHouse 类型 |
|---|---|---|
| ObjectId | STRING | String |
| String | STRING | String |
| Int32 | INT | Int32 |
| Int64 | BIGINT | Int64 |
| Double | DOUBLE | Float64 |
| Boolean | BOOLEAN | Bool |
| Date | TIMESTAMP(3) | DateTime64(3) |
| Array | ARRAY | Array(T) |
| Embedded Document | ROW | String / JSON 类型(自定义处理) |
Embedded Document 是最容易出问题的,MongoDB 里嵌套文档在 ClickHouse 没有天然对应类型。我的做法是先把整个嵌套文档在 Flink 里转成 JSON 字符串,ClickHouse 侧用 String 存储,查询时再用 ClickHouse 的 JSON 函数解析。这样虽然牺牲了一点查询性能,但表结构稳定,不需要因为嵌套字段的变化去改表。
5.3 表引擎和主键选择
ClickHouse 建表时表引擎的选择直接影响同步语义。如果你要的是纯追加日志型数据(比如埋点、日志),用 MergeTree 就够了,ORDER BY 选一个业务时间字段或递增 ID 即可。但如果 MongoDB 集合里有数据更新,比如用户资料修改、订单状态流转,你要保证 ClickHouse 里同一个主键不会被插入多条脏数据,这时建议用 ReplacingMergeTree,并在建表时指定 version 字段。
实际操作里我会把 MongoDB 文档里的更新时间作为 version。建表示例:
CREATE TABLE ods.users ( _id String, name String, email String, updated_at DateTime64(3), created_at DateTime64(3) ) ENGINE = ReplacingMergeTree(updated_at) ORDER BY _id;这样同一 _id 的多次变更到达 ClickHouse 后,后台合并时会保留 updated_at 最新的那一条。注意,ReplacingMergeTree 的合并是后台发生的,不是写入时立刻生效,查询时如果想要强一致,可以用 FINAL 关键字或者按更新时间取 max。对实时数仓的 ODS 层来说,短期有重复行通常可以接受,最终一致性由合并保证。
6. 实操验证:从零到一跑通同步
6.1 启动任务与效果验证
配置准备好后,启动任务很简单。用 Flink CDC 的 shell 工具直接提交 YAML:
bin/flink-cdc.sh path/to/mongodb-to-clickhouse.yaml也可以用 Flink SQL 客户端方式,先把 MongoDB 和 ClickHouse 的 source/sink 表建好,再执行 INSERT INTO SELECT。SQL 方式适合那些需要精确控制字段映射、做数据清洗转换的场景,但多集合时每条线都要写一段 CREATE TABLE 和 INSERT,YAML 方式相比之下确实更贴合"一套配置"的主题。
任务启动后,我一般从两个角度验证。第一,最快的方式是看 ClickHouse 里表行数有没有涨:对 MongoDB 侧执行一次 insert 和一次 update,然后去目标表 count,应该能在几秒内看到变化。第二,看 Flink Web UI 的 task 状态和 checkpoints 状态,确认 checkpoint 一直显示 OK,这是任务长期稳定运行的底气。增量同步的核心验证是修改一条 MongoDB 文档,然后在 ClickHouse 里查对应主键,确认 updated_at 已经变成新值。
6.2 常见问题排查速查表
做这种多环节的实时同步,问题往往不出在核心逻辑,而是出在一些意想不到的边界条件。下面几个问题都是我实际遇到过、且社区里高频出现的:
| 现象 | 可能原因 | 排查方法 |
|---|---|---|
| 任务启动报错,提示 Change Streams 不可用 | MongoDB 是单节点,不是副本集 | 检查 rs.status(),确认是副本集 |
| 全量阶段正常,增量阶段报 not authorized | 账号缺 changeStream 权限 | 按上文角色授权,重启任务 |
| 写入 ClickHouse 报 DateTime 类型不匹配 | 时区或精度不一致 | ClickHouse 建表用 DateTime64(3),Flink 侧提前转换 |
| 同步延迟越来越大 | sink 批量参数太小或并行度不足 | 调大 batch.size,观察 CPU 是否打满 |
| 任务重启后报 resume token 失效 | oplog 太小,变更已超出回溯窗口 | 调大 oplog 容量,或接受重新快照 |
| ClickHouse 重启后偶发 failed to flush system log already exists | 系统表目录残留或进程未完全退出 | 先确认没有重复进程,必要时清理 system 相关日志表数据 |
这里特别说一下 ClickHouse 重启报错那个问题。ClickHouse 在启动时会尝试 flush 系统日志表,如果上一次进程没有干净退出,或者同时起了两个 server 进程,系统表目录里残留了元数据,就会报 "failed to flush system log already exists"。很多人一看到就去删 /var/lib/clickhouse 下的文件,非常危险。正确的排查顺序是:先确认没有残留的 clickhouse-server 进程,再查看 system.query_log 这类系统表目录是否有异常文件,确认没问题后重启。实测中绝大多数情况是重复进程导致的,杀掉旧进程重启就好。
7. 优化心得与扩展建议
最后再分享几个我在实际项目中反复调优得出的体会。第一个是批量参数要跟 checkpoint 配合着来,不要单独看某一项。batch.size 调大以后,一次 checkpoint 周期内攒的数据也会变多,如果 checkpoint 间隔太短,可能出现状态还没落盘下一批数据就又来了,造成背压。我的经验是 batch.size 增大一倍时,checkpoint 间隔也相应增加,保证两者匹配。
第二个是 MongoDB 集合出现 schema 变更时,不要慌。Flink CDC 对新增字段的处理策略是尽力而为,如果 ClickHouse 表里没有新增列,可以暂时忽略,等夜间维护窗口加上列再重启任务即可。但如果 MongoDB 侧删了字段,在 ClickHouse 里通常表现是写入默认值,不会导致任务失败。所以我的做法是 ClickHouse 表字段只增不减,用默认值兜底,这是成本最低的容错方式。
第三个建议是分库分表场景的扩展。如果你的 MongoDB 有多个库,每个库里都有相同结构的业务集合,比如分库后的 users,route 规则里可以用正则表达式匹配,像 source-table: shard_.*.users 这种写法,一套配置就能覆盖几十个集合。这个技能在业务分库规模较大的团队特别实用,上线新分库时完全不用动代码,只要确认 ClickHouse 目标表存在即可。
这套 Flink 1.20 加 MongoDB 多集合实时同步到 ClickHouse 的方案,我后来在好几个项目里复用,最直观的收益是同步任务的数量从十几个降到了两三个,监控告警范围大幅收缩。如果你正被一堆 MongoDB 集合的表同步折腾,照着这个思路配置,大概率能省下不少维护精力。