用 Flink MongoDB CDC 把文档库接入实时数仓
一、依赖与运行环境准备
1.1 Maven 依赖(工程开发)
如果你是用 Java / Scala 自己写 Flink 程序,需要在工程里引入:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-mongodb-cdc</artifactId>
<version>3.5.0</version>
</dependency>
1.2 Flink SQL Client / Standalone 集群
如果你更多是用 Flink SQL 做任务,可以下载对应的 SQL JAR:
flink-sql-connector-mongodb-cdc-3.5.0.jar
放到:
<FLINK_HOME>/lib/
重启 Flink 集群后,就可以在 SQL 里直接:
WITH ('connector' = 'mongodb-cdc', ...)
二、MongoDB 端前置条件:版本、集群与权限
MongoDB CDC 的实现依赖 Change Streams(3.6 引入),不是直接扒 oplog,所以对 Mongo 端的要求比较明确。
2.1 版本与部署要求
必须满足:
-
MongoDB 版本 ≥ 3.6
- Change Streams 功能从 3.6 起提供。
-
部署方式:Replica Set 或 Sharded Cluster
- Standalone 不支持 Change Streams。
-
存储引擎:WiredTiger
- 现代 Mongo 默认就是它。
-
复制协议版本:pv1
- 4.0 之后只支持 pv1,新集群默认就是 pv1。
2.2 权限:为 Flink 准备一个“监听用户”
MongoDB CDC 内部整合了官方的 Kafka Connector,需要用户具备 changeStream + read 等权限。
可以参考下面这段简单授权脚本:
use admin;
db.createRole(
{
role: "flinkrole",
privileges: [{
// 所有非系统库、非系统集合
resource: { db: "", collection: "" },
actions: [
"splitVector",
"listDatabases",
"listCollections",
"collStats",
"find",
"changeStream"
]
}],
roles: [
// 针对分片集群快照拆分时,需要读 config 库
{ role: 'read', db: 'config' }
]
}
);
db.createUser(
{
user: 'flinkuser',
pwd: 'flinkpw',
roles: [
{ role: 'flinkrole', db: 'admin' }
]
}
);
后续如果用到了数据库正则匹配(
databaseList用 regex),还需要readAnyDatabase角色。
三、在 Flink SQL 中定义 MongoDB CDC 表
先上一个完整示例,感受一下语法:
-- 注册 MongoDB 中的 products 集合为一张动态表
CREATE TABLE products (
_id STRING, -- 必须声明
name STRING,
weight DECIMAL(10,3),
tags ARRAY<STRING>, -- 数组
price ROW<amount DECIMAL(10,2), -- 嵌套文档
currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>, -- 嵌套文档数组
PRIMARY KEY(_id) NOT ENFORCED
) WITH (
'connector' = 'mongodb-cdc',
'hosts' = 'localhost:27017,localhost:27018,localhost:27019',
'username' = 'flinkuser',
'password' = 'flinkpw',
'database' = 'inventory',
'collection'= 'products'
);
-- 像查表一样查 Mongo 集合:先读快照,再持续读变更
SELECT * FROM products;
这里有几个关键点:
-
_id必须显式声明,并作为主键- Mongo 的变更事件里,Delete 只包含
_id和分片键; - 如果你把其他字段设成主键,Delete 事件就没法对齐下游数据;
- 所以 Flink CDC 明确要求主键只能是
_id。
- Mongo 的变更事件里,Delete 只包含
-
Flink SQL 支持丰富的复杂类型:
ARRAY<STRING>ROW<...>ARRAY<ROW<...>>等嵌套结构
-
因为 Change Stream 默认没有 Update Before,所以变更流语义是 Upsert:
- 插入、更新后、删除都可以映射到 Upsert 流;
- 配合下游 Upsert Sink(如 Upsert Kafka、Upsert JDBC 等)就能保证幂等。
四、核心 Connector 参数拆解
本文只挑最关键的参数解释,方便你直接抄配置。
4.1 基础连接参数
| 参数名 | 必填 | 默认 | 说明 |
|---|---|---|---|
| connector | 是 | - | 固定为 mongodb-cdc |
| scheme | 否 | mongodb | 协议,mongodb 或 mongodb+srv |
| hosts | 是 | - | Mongo 节点列表,逗号分隔:host1:27017,host2:27018 |
| username | 否 | - | Mongo 认证用户名 |
| password | 否 | - | Mongo 认证密码 |
| database | 否 | - | 要监控的数据库,支持正则;不填则表示所有数据库 |
| collection | 否 | - | 要监控的集合,支持正则;不填则表示所有集合 |
| connection.options | 否 | - | 附加连接参数,如 replicaSet=myset&connectTimeoutMS=300000 |
实践中通常会同时限定
database和collection,减少无关扫描。
4.2 启动模式:scan.startup.mode
支持三种模式:
'scan.startup.mode' = 'initial' | 'latest-offset' | 'timestamp'
语义如下:
-
initial(默认)-
首次启动:
- 对目标集合做一次全量快照;
-
然后:
- 从当前最新的 Change Stream resumeToken 继续消费。
-
-
latest-offset- 不做快照;
- 直接从“当前时间点”开始消费变更,只关心之后发生的变化。
-
timestamp-
跳过快照,从指定的时间戳开始消费 Change Stream。
-
示例:
'scan.startup.mode' = 'timestamp', 'scan.startup.timestamp-millis' = '1667232000000'
-
在 DataStream API 中对应的是 StartupOptions:
MongoDBSource.<String>builder()
.startupOptions(StartupOptions.latest()) // 从最新位置开始
.startupOptions(StartupOptions.timestamp(1667232000000L))// 从指定时间点开始
.build();
4.3 快照相关参数
当 scan.startup.mode = 'initial' 时,你可以通过以下参数控制快照行为:
-
initial.snapshotting.queue.size(默认 16000)- 拷贝数据时内部缓冲队列大小;
-
initial.snapshotting.max.threads(默认=处理器核数)- 首次快照拷贝使用的线程数;
-
initial.snapshotting.pipeline- 一段 JSON 数组,代表 Mongo Aggregation Pipeline,用于在快照阶段过滤数据。
例如:只快照 closed=false 的文档:
'initial.snapshotting.pipeline'
= '[ { "$match": { "closed": "false" } } ]'
注意:
- 只在
scan.startup.mode = 'initial'时生效; - 只用于 Debezium 模式,不能配合增量快照使用;
- 是标准 Mongo 聚合语法的 JSON 字符串。
4.4 增量快照与并行快照
从 Flink CDC 2.3.0 开始,Mongo 也支持 增量快照 模式:
-
开关:
'scan.incremental.snapshot.enabled' = 'true' -
要求:
- MongoDB 版本 ≥ 4.0。
相关参数:
scan.incremental.snapshot.chunk.size.mb(默认 64)
每个快照 chunk 的大小(MB);scan.incremental.snapshot.chunk.samples(默认 20)
基于 sample 分片策略时,每个 chunk 的采样数;scan.incremental.close-idle-reader.enabled
快照完成后是否关闭空闲 Reader;scan.cursor.no-timeout(默认 true)
禁用 10 分钟空闲超时,特别适合长时间快照任务;scan.incremental.snapshot.unbounded-chunk-first.enabled
优先调度大块数据,降低 OOM 风险;scan.incremental.snapshot.backfill.skip(慎用)
跳过快照期间的变更回填,会退化为“至少一次+可能重放”。
实际建议:
- 集合数据量大、更新频率高 ⇒ 强烈建议开启增量快照;
- Demo 或数据量很小 ⇒ 可以关闭,但注意快照阶段 checkpoint 超时问题。
4.5 心跳:heartbeat.interval.ms
'heartbeat.interval.ms' = '10000' -- 比如 10s
当集合变更不频繁时:
- Change Stream 的 resumeToken 长期不推进;
- 如果这期间你做了 checkpoint / savepoint,之后恢复时可能发现 token 已经过期。
因此官方建议:
对更新不频繁的集合,一定要把
heartbeat.interval.ms设为大于 0 的值,保证 resumeToken 持续往前走。
五、Exactly-Once 与 Full Changelog
5.1 Exactly-Once 处理语义
MongoDB CDC 是一个标准的 Flink Source:
- 首先读取初始快照;
- 然后持续消费 Change Stream 变更事件;
- 结合 Flink Checkpoint,对 offset / resumeToken 做状态管理。
只要下游 Sink 支持幂等或事务,就可以实现端到端的 Exactly-Once 处理语义。
5.2 Full Changelog(MongoDB 6.0+)
MongoDB 6.0 起,Change Stream 支持输出 pre-image 和 post-image:
- pre-image:文档被更新、替换、删除之前的版本;
- post-image:文档被插入、替换、更新之后的版本。
Flink MongoDB CDC 利用这个能力,可以直接产出完整的变更流:
- Insert ⇒
+I - UpdateBefore ⇒
-U - UpdateAfter ⇒
+U - Delete ⇒
-D
这样可以省掉下游一个 ChangelogNormalize 节点,拓扑更简单,延迟更低。
要启用 Full Changelog,需要:
-
MongoDB 版本 ≥ 6.0;
-
数据库级开启 preAndPostImages:
db.runCommand({ setClusterParameter: { changeStreamOptions: { preAndPostImages: { expireAfterSeconds: 'off' // 可替换为自定义过期时间 } } } }) -
集合级开启 changeStreamPreAndPostImages:
db.runCommand({ collMod: "<< collection name >>", changeStreamPreAndPostImages: { enabled: true } }) -
Flink 端开启
scan.full-changelog:-
DataStream:
MongoDBSource.<String>builder() .scanFullChangelog(true) ... .build(); -
SQL:
CREATE TABLE mongodb_source (...) WITH ( 'connector' = 'mongodb-cdc', 'scan.full-changelog' = 'true', ... );
-
六、DataStream API 使用示例
6.1 基础 SourceFunction 用法
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.cdc.connectors.mongodb.MongoDBSource;
public class MongoDBSourceExample {
public static void main(String[] args) throws Exception {
SourceFunction<String> sourceFunction = MongoDBSource.<String>builder()
.hosts("localhost:27017")
.username("flink")
.password("flinkpw")
.databaseList("inventory") // 支持正则
.collectionList("inventory.products", "inventory.orders") // 支持正则
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(sourceFunction)
.print()
.setParallelism(1); // 为了保证输出顺序,sink 一般设为 1
env.execute();
}
}
6.2 增量快照并行 Source(2.3.0+)
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.cdc.connectors.mongodb.source.MongoDBSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;
public class MongoDBIncrementalSourceExample {
public static void main(String[] args) throws Exception {
MongoDBSource<String> mongoSource =
MongoDBSource.<String>builder()
.hosts("localhost:27017")
.databaseList("inventory") // 正则也可
.collectionList("inventory.products", "inventory.orders")
.username("flink")
.password("flinkpw")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 启用 checkpoint
env.enableCheckpointing(3000);
env.fromSource(
mongoSource,
WatermarkStrategy.noWatermarks(),
"MongoDBIncrementalSource")
.setParallelism(2)
.print()
.setParallelism(1);
env.execute("Print MongoDB Snapshot + Change Stream");
}
}
如果
databaseList使用了正则,需要为用户授予readAnyDatabase角色。
七、元数据列与 Source 指标
7.1 元数据列(Metadata)
MongoDB CDC 可以把一些元信息作为只读列暴露出来,典型的有:
database_name:数据库名;collection_name:集合名;op_ts:操作发生时间;row_kind:变更类型,+I/-D/-U/+U。
完整示例:
CREATE TABLE products (
db_name STRING METADATA FROM 'database_name' VIRTUAL,
collection_name STRING METADATA FROM 'collection_name' VIRTUAL,
operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,
operation STRING METADATA FROM 'row_kind' VIRTUAL,
_id STRING,
name STRING,
weight DECIMAL(10,3),
tags ARRAY<STRING>,
price ROW<amount DECIMAL(10,2), currency STRING>,
suppliers ARRAY<ROW<name STRING, address STRING>>,
PRIMARY KEY(_id) NOT ENFORCED
) WITH (
'connector' = 'mongodb-cdc',
'hosts' = 'localhost:27017,localhost:27018,localhost:27019',
'username' = 'flinkuser',
'password' = 'flinkpw',
'database' = 'inventory',
'collection' = 'products'
);
注意:
row_kind在有回撤(Update Before / Delete)时会参与比对,一些复杂 SQL 场景下可能需要小心处理,比方说只在简单同步任务中使用它来做审计。
7.2 Source Metrics 监控指标
Connector 会用分组 namespace.schema.table 暴露一组指标:
(Mongo 中 schema 为空字符串,所以最终类似 inventory.products)
isSnapshotting:当前是否处于快照阶段;isStreamReading:当前是否在消费增量变更;numTablesSnapshotted/numTablesRemaining:
已完成快照/尚未快照的集合数量;numSnapshotSplitsProcessed/Remaining/Finished:
快照拆分块处理进度;snapshotStartTime/snapshotEndTime:
快照开始/结束时间。
配合 Prometheus + Grafana,可以很方便地做一套 CDC 状态看板。
八、BSON → Flink SQL 类型映射
最后再看一下 BSON 类型和 Flink SQL 类型的映射,方便你设计 Schema:
| BSON 类型 | Flink SQL 类型 |
|---|---|
| TINYINT / SMALLINT / Int | INT / SMALLINT(按位宽映射) |
| Long | BIGINT |
| FLOAT / Double | DOUBLE |
| Decimal128 | DECIMAL(p, s) |
| Boolean | BOOLEAN |
| Date / Timestamp(日期) | DATE / TIME / TIMESTAMP(3/0) / TIMESTAMP_LTZ |
| String / ObjectId / UUID / Regex… | STRING |
| BinData | BYTES |
| Object | ROW |
| Array | ARRAY |
| DBPointer | ROW<$ref STRING, $id STRING> |
| GeoJSON Point | ROW<type STRING, coordinates ARRAY> |
| GeoJSON Line | ROW<type STRING, coordinates ARRAY<ARRAY>> |
实战里要特别注意:
- 金额类字段:建议统一成
DECIMAL(18,2)或DECIMAL(20,4)这类固定精度; - 时间类字段:跨时区场景推荐统一用
TIMESTAMP_LTZ(3)。
九、实战小结与配置建议
如果你准备在生产中落地 Flink MongoDB CDC,可以参考这样一套 checklist:
-
先小规模 PoC:
- 起一个测试 replica set;
- 给 Flink 配一个
flinkuser; - 用 Flink SQL 注册 1~2 个集合先跑通。
-
根据数据量与业务需求选择启动模式:
- 需要完整历史 + 后续变更 ⇒
scan.startup.mode='initial'; - 只关心从现在开始 ⇒
scan.startup.mode='latest-offset'。
- 需要完整历史 + 后续变更 ⇒
-
大集合 & 高频更新 ⇒ 建议开启增量快照:
scan.incremental.snapshot.enabled = true;- 调整
chunk.size.mb、chunk.samples等参数。
-
慢变更集合一定要配置心跳:
heartbeat.interval.ms > 0,避免 resumeToken 过期。
-
Mongo ≥ 6.0 时优先考虑 Full Changelog:
- 开启 preAndPostImages +
scan.full-changelog=true; - 下游直接拿到完整的
+I/-D/-U/+U流,算子拓扑更简洁。
- 开启 preAndPostImages +
更多推荐
所有评论(0)