一、依赖与运行环境准备

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 版本与部署要求

必须满足:

  1. MongoDB 版本 ≥ 3.6

    • Change Streams 功能从 3.6 起提供。
  2. 部署方式:Replica Set 或 Sharded Cluster

    • Standalone 不支持 Change Streams。
  3. 存储引擎:WiredTiger

    • 现代 Mongo 默认就是它。
  4. 复制协议版本: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;

这里有几个关键点:

  1. _id 必须显式声明,并作为主键

    • Mongo 的变更事件里,Delete 只包含 _id 和分片键;
    • 如果你把其他字段设成主键,Delete 事件就没法对齐下游数据;
    • 所以 Flink CDC 明确要求主键只能是 _id
  2. Flink SQL 支持丰富的复杂类型:

    • ARRAY<STRING>
    • ROW<...>
    • ARRAY<ROW<...>> 等嵌套结构
  3. 因为 Change Stream 默认没有 Update Before,所以变更流语义是 Upsert

    • 插入、更新后、删除都可以映射到 Upsert 流;
    • 配合下游 Upsert Sink(如 Upsert Kafka、Upsert JDBC 等)就能保证幂等。

四、核心 Connector 参数拆解

本文只挑最关键的参数解释,方便你直接抄配置。

4.1 基础连接参数

参数名必填默认说明
connector-固定为 mongodb-cdc
schememongodb协议,mongodbmongodb+srv
hosts-Mongo 节点列表,逗号分隔:host1:27017,host2:27018
username-Mongo 认证用户名
password-Mongo 认证密码
database-要监控的数据库,支持正则;不填则表示所有数据库
collection-要监控的集合,支持正则;不填则表示所有集合
connection.options-附加连接参数,如 replicaSet=myset&connectTimeoutMS=300000

实践中通常会同时限定 databasecollection,减少无关扫描。

4.2 启动模式:scan.startup.mode

支持三种模式:

'scan.startup.mode' = 'initial' | 'latest-offset' | 'timestamp'

语义如下:

  1. initial(默认)

    • 首次启动:

      • 对目标集合做一次全量快照
    • 然后:

      • 从当前最新的 Change Stream resumeToken 继续消费。
  2. latest-offset

    • 不做快照;
    • 直接从“当前时间点”开始消费变更,只关心之后发生的变化。
  3. 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-imagepost-image

  • pre-image:文档被更新、替换、删除之前的版本;
  • post-image:文档被插入、替换、更新之后的版本。

Flink MongoDB CDC 利用这个能力,可以直接产出完整的变更流:

  • Insert ⇒ +I
  • UpdateBefore ⇒ -U
  • UpdateAfter ⇒ +U
  • Delete ⇒ -D

这样可以省掉下游一个 ChangelogNormalize 节点,拓扑更简单,延迟更低。

要启用 Full Changelog,需要:

  1. MongoDB 版本 ≥ 6.0;

  2. 数据库级开启 preAndPostImages:

    db.runCommand({
      setClusterParameter: {
        changeStreamOptions: {
          preAndPostImages: {
            expireAfterSeconds: 'off' // 可替换为自定义过期时间
          }
        }
      }
    })
    
  3. 集合级开启 changeStreamPreAndPostImages:

    db.runCommand({
      collMod: "<< collection name >>",
      changeStreamPreAndPostImages: { enabled: true }
    })
    
  4. 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 / IntINT / SMALLINT(按位宽映射)
LongBIGINT
FLOAT / DoubleDOUBLE
Decimal128DECIMAL(p, s)
BooleanBOOLEAN
Date / Timestamp(日期)DATE / TIME / TIMESTAMP(3/0) / TIMESTAMP_LTZ
String / ObjectId / UUID / Regex…STRING
BinDataBYTES
ObjectROW
ArrayARRAY
DBPointerROW<$ref STRING, $id STRING>
GeoJSON PointROW<type STRING, coordinates ARRAY>
GeoJSON LineROW<type STRING, coordinates ARRAY<ARRAY>>

实战里要特别注意:

  • 金额类字段:建议统一成 DECIMAL(18,2)DECIMAL(20,4) 这类固定精度;
  • 时间类字段:跨时区场景推荐统一用 TIMESTAMP_LTZ(3)

九、实战小结与配置建议

如果你准备在生产中落地 Flink MongoDB CDC,可以参考这样一套 checklist:

  1. 先小规模 PoC:

    • 起一个测试 replica set;
    • 给 Flink 配一个 flinkuser
    • 用 Flink SQL 注册 1~2 个集合先跑通。
  2. 根据数据量与业务需求选择启动模式:

    • 需要完整历史 + 后续变更 ⇒ scan.startup.mode='initial'
    • 只关心从现在开始 ⇒ scan.startup.mode='latest-offset'
  3. 大集合 & 高频更新 ⇒ 建议开启增量快照:

    • scan.incremental.snapshot.enabled = true
    • 调整 chunk.size.mbchunk.samples 等参数。
  4. 慢变更集合一定要配置心跳:

    • heartbeat.interval.ms > 0,避免 resumeToken 过期。
  5. Mongo ≥ 6.0 时优先考虑 Full Changelog:

    • 开启 preAndPostImages + scan.full-changelog=true
    • 下游直接拿到完整的 +I/-D/-U/+U 流,算子拓扑更简洁。

更多推荐