前阵子帮一个做电商数据分析的团队处理数据同步他们在 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。这条路看起来灵活但实际一深入就发现工作量全在细节里每个集合要自己维护断点续传的 offsetMongoDB 重启之后 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 留了 2GBTaskManager 每台 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 pipeline3.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_atorders 集合有 _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 依赖 oplogoplog 的大小决定了能回溯的变更窗口。如果 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 设 2000flush.interval-ms 设 1000意思是每秒或者每攒够 2000 条就刷一次。如果数据量大可以把 batch.size 提到 5000 甚至 10000ClickHouse 单次插入的吞吐会更好但内存占用也会同步上升TaskManager 给不够会出现 OOM。我实测下来2000 条一批写入延迟已经能控制在 1 秒左右对实时分析够用了。5.2 字段类型映射对照MongoDB 的 BSON 类型和 ClickHouse 的类型体系差别很大好在 Flink CDC 已经做了大部分映射。实际建表时我整理了这样一张映射对照方便团队新人直接参考MongoDB BSON 类型Flink 类型ClickHouse 类型ObjectIdSTRINGStringStringSTRINGStringInt32INTInt32Int64BIGINTInt64DoubleDOUBLEFloat64BooleanBOOLEANBoolDateTIMESTAMP(3)DateTime64(3)ArrayARRAYArray(T)Embedded DocumentROWString / 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 工具直接提交 YAMLbin/flink-cdc.sh path/to/mongodb-to-clickhouse.yaml也可以用 Flink SQL 客户端方式先把 MongoDB 和 ClickHouse 的 source/sink 表建好再执行 INSERT INTO SELECT。SQL 方式适合那些需要精确控制字段映射、做数据清洗转换的场景但多集合时每条线都要写一段 CREATE TABLE 和 INSERTYAML 方式相比之下确实更贴合一套配置的主题。任务启动后我一般从两个角度验证。第一最快的方式是看 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 有多个库每个库里都有相同结构的业务集合比如分库后的 usersroute 规则里可以用正则表达式匹配像 source-table: shard_.*.users 这种写法一套配置就能覆盖几十个集合。这个技能在业务分库规模较大的团队特别实用上线新分库时完全不用动代码只要确认 ClickHouse 目标表存在即可。这套 Flink 1.20 加 MongoDB 多集合实时同步到 ClickHouse 的方案我后来在好几个项目里复用最直观的收益是同步任务的数量从十几个降到了两三个监控告警范围大幅收缩。如果你正被一堆 MongoDB 集合的表同步折腾照着这个思路配置大概率能省下不少维护精力。