Skip to main content

概述

MongoDB Source Connector 使用 MongoDB Change Streams 读取数据库、集合或整个部署中的变更,并将事件写入 Kafka Topic,适合构建数据同步、流式处理、搜索索引和审计等数据链路。默认情况下,MongoDB 命名空间 database.collection 映射到同名 Kafka Topic,消息值包含完整的变更流事件。 Connector 可以只采集启动后产生的新变更,也可以先复制已有文档,再衔接后续变更。复制的文档会封装为形似插入操作的 Change Stream 事件,并经过与实时事件相同的 Topic、键和值映射流程。

前置条件

  • MongoDB 部署必须支持 Change Streams;连接账号需要对采集范围具备读取变更流所需权限,使用存量复制时还需要读取目标集合。
  • 使用变更前文档、扩展事件或 Stable API 等 MongoDB 能力时,服务端版本和数据库设置必须满足对应功能要求。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和支持 Change Streams 的 MongoDB 部署,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
替换 MongoDB 连接信息、数据库名和集合名。该配置从当前变更流位置开始采集,默认将事件写入 <database-name>.<collection-name> Topic;建议通过 Config Provider 或受保护的配置管理连接凭据。

配置

Connector 与任务

connector.class

指定 MongoDB Source Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 com.mongodb.kafka.connect.MongoSourceConnector
  • 必填:是

tasks.max

设置 Connect 允许创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须至少为 1;此 Connector 始终只创建一个 Source Task,设置为更大的值不会增加并行 Task。

MongoDB 连接与采集范围

connection.uri

MongoDB 连接字符串,用于指定主机、认证、Replica Set、TLS 和其他驱动选项。
  • 类型password
  • 默认值mongodb://localhost:27017,localhost:27018,localhost:27019
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效的 MongoDB URI。默认值通常只适合本地开发;不要在明文配置中保存凭据。

database

限定要监听的数据库。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:留空时监听整个 MongoDB 部署;设置数据库但不设置集合时监听该数据库中的集合。

collection

限定要监听的单个集合。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:不是正则表达式或列表;设置时必须同时设置 database

offset.partition.name

覆盖 Connector 用于标识 Source Offset 分区的名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:留空时根据连接主机、数据库和集合生成。修改此值会开始使用新的 Offset 身份,原有 Resume Token 不会自动迁移。

MongoDB Stable API

server.api.version

启用指定版本的 MongoDB Stable API。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:留空表示不启用;可用值为 1,要求 MongoDB 5.0 或更高版本。

server.api.deprecation.errors

启用 Stable API 弃用错误。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 server.api.version 非空时生效。

server.api.strict

启用 Stable API 严格模式。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 server.api.version 非空时生效。

TLS 证书

connection.ssl.truststore

指定 Connect Worker 本地的 Truststore 路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时使用 JVM 默认 KeyStore 类型加载;TLS 本身仍需由连接 URI 或部署设置启用。

connection.ssl.truststorePassword

指定 Truststore 密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 connection.ssl.truststore 非空时使用。

connection.ssl.keystore

指定 Connect Worker 本地的 Keystore 路径,用于提供客户端证书和私钥。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时使用 JVM 默认 KeyStore 类型加载。

connection.ssl.keystorePassword

指定 Keystore 及其中私钥的密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 connection.ssl.keystore 非空时使用;同一密码用于加载 KeyStore 和初始化私钥。

Change Stream

pipeline

为 Change Stream 设置聚合管道。
  • 类型string
  • 默认值[]
  • 重要级别:中
  • 有效值 / 注意事项:必须是 JSON 文档数组。管道需保留 Topic、Offset、键和值映射所需的 _idnsdocumentKeyfullDocument 字段;存量复制的高效筛选优先使用 startup.mode.copy.existing.pipeline

batch.size

设置 MongoDB Change Stream Cursor 的批次大小提示。
  • 类型int
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 00 表示沿用驱动和服务端行为,与 poll.max.batch.size 不同。

publish.full.document.only

仅将变更事件的 fullDocument 作为消息值发布。
  • 类型boolean
  • 默认值false
  • 重要级别:高
  • 有效值 / 注意事项:设为 true 时强制使用 updateLookup;没有 fullDocument 的事件会被过滤,除非同时启用删除 Tombstone。

publish.full.document.only.tombstone.on.delete

在仅发布完整文档模式下,为缺少文档型 fullDocument 的事件发送空值 Tombstone;删除事件是最常见的情形。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 publish.full.document.only=true 时生效。实现按是否存在文档型 fullDocument 判断,不只检查事件的 operationType 是否为 delete

change.stream.document.key.as.key

控制非 Schema 输出是否使用 Change Stream 的 documentKey 作为 Kafka 消息键。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:设为 false 或事件没有 documentKey 时,键使用 Resume Token 文档;output.format.key=schema 时忽略此项。

change.stream.full.document.before.change

设置 Change Stream 返回变更前文档的方式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:可用值为 defaultoffwhenAvailablerequired 或空字符串。需要 MongoDB 6.0 或更高版本并为集合启用 Pre-image;不影响存量复制记录。

change.stream.full.document

设置 Change Stream 返回完整文档的方式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:可用值为 defaultupdateLookupwhenAvailablerequired 或空字符串;publish.full.document.only=true 时强制使用 updateLookup

change.stream.show.expanded.events

控制是否请求扩展 Change Stream 事件。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:扩展 DDL 事件需要 MongoDB 6.0 或更高版本,消除歧义的更新路径需要 MongoDB 6.1 或更高版本。

collation

为 Change Stream 设置 Collation。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:留空表示不设置;否则必须是有效的 MongoDB Collation JSON 文档。该设置不应用于存量复制聚合。

poll.max.batch.size

限制单次 poll 返回的最大记录数。
  • 类型int
  • 默认值1000
  • 重要级别:低
  • 有效值 / 注意事项:必须至少为 1;只限制记录数量,不限制单批总字节数。

poll.await.time.ms

设置 Change Stream 等待新事件的最长时间。
  • 类型long
  • 默认值5000
  • 重要级别:低
  • 有效值 / 注意事项:必须至少为 1,单位为毫秒;该值不是 Kafka Connect 的 Poll Timeout。

Topic 映射

topic.mapper

指定将 MongoDB 命名空间映射为 Kafka Topic 的实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.source.topic.mapping.DefaultTopicMapper
  • 重要级别:高
  • 有效值 / 注意事项:必须是可加载且实现 TopicMapper 的完整 Java 类名。自定义实现可以读取自身定义的附加属性。

topic.separator

指定默认 Topic 名各部分之间的分隔符。
  • 类型string
  • 默认值.
  • 重要级别:低
  • 有效值 / 注意事项:用于连接前缀、数据库、集合和后缀;Connector 不校验最终 Topic 名是否合法。

topic.prefix

为默认映射产生的所有 Topic 添加前缀。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:非空前缀后会追加 topic.separator

topic.suffix

为默认映射产生的所有 Topic 添加后缀。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:非空后缀前会添加 topic.separator

topic.namespace.map

使用 JSON 对象将 MongoDB 命名空间映射到指定 Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:支持完整命名空间、数据库、按声明顺序匹配的 /正则/* 映射;映射结果仍会应用前缀和后缀。模板可使用 dbsepcollsep_collcoll_sepsep_coll_sep 变量;需自行避免非法名称和 Topic 冲突。

输出格式与 Schema

output.format.key

设置消息键的输出格式。
  • 类型string
  • 默认值json
  • 重要级别:高
  • 有效值 / 注意事项:可用值为 jsonbsonschema,不区分大小写;应与 key.converter 兼容。

output.format.value

设置消息值的输出格式。
  • 类型string
  • 默认值json
  • 重要级别:高
  • 有效值 / 注意事项:可用值为 jsonbsonschema,不区分大小写;应与 value.converter 兼容。

output.schema.infer.value

为每条消息值自动推断 Connect Schema。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 output.format.value=schema 时生效。文档结构变化可能使同一 Topic 出现多个 Schema。

output.schema.key

定义 Schema 格式消息键使用的 Avro Schema。
  • 类型string
  • 默认值{ "type": "record", "name": "keySchema", "fields" : [{"name": "_id", "type": "string"}]}
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效且与完整 Change Stream 文档转换兼容的 Connector Avro Schema;仅在 output.format.key=schema 时读取。

output.schema.value

定义 Schema 格式消息值使用的 Avro Schema。
  • 类型string
  • 默认值:内置 ChangeStream Avro Schema
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效且与消息值兼容的 Connector Avro Schema;output.schema.infer.value=true 或输出格式不是 schema 时忽略。

output.json.formatter

指定 BSON 到 JSON 的格式化实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.source.json.formatter.DefaultJson
  • 重要级别:中
  • 有效值 / 注意事项:必须是可加载且实现 JsonWriterSettingsProvider 的完整 Java 类名;JSON 输出及 Schema 转换会使用,原始 BSON 输出会忽略。

key.converter

覆盖 Worker 级消息键 Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项null 表示继承 Worker 配置;所选 Converter 必须兼容 output.format.key 产生的 String、BSON Bytes 或 Schemaful Connect Data。

value.converter

覆盖 Worker 级消息值 Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项null 表示继承 Worker 配置;所选 Converter 必须兼容 output.format.value 产生的 String、BSON Bytes 或 Schemaful Connect Data。

启动位置与存量复制

startup.mode

设置没有可用 Source Offset 时的启动方式。
  • 类型string
  • 默认值:空字符串,实际按 latest 处理
  • 重要级别:中
  • 有效值 / 注意事项:可用值为 latesttimestampcopy_existing 或空字符串。已有 Resume Token 优先于该设置;显式非空值优先于弃用的 copy.existing

startup.mode.timestamp.start.at.operation.time

设置 Timestamp 启动模式的起始操作时间。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 startup.mode=timestamp 且没有可用 Offset 时应用;支持十进制 Epoch 秒、秒精度 ISO-8601 Instant 或 Canonical Extended JSON BSON Timestamp。

startup.mode.copy.existing.max.threads

设置存量复制读取线程数。
  • 类型int
  • 默认值Runtime.getRuntime().availableProcessors()(运行环境可用处理器数量)
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 1,仅在 startup.mode=copy_existing 时使用;实际线程数不会超过选中的命名空间数量。显式设置时覆盖弃用项 copy.existing.max.threads

startup.mode.copy.existing.queue.size

设置存量复制线程与 Source Task 之间的内存队列容量。
  • 类型int
  • 默认值16000
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 1,仅在 startup.mode=copy_existing 时使用;显式设置时覆盖弃用项 copy.existing.queue.size

startup.mode.copy.existing.pipeline

设置只用于读取存量文档的聚合管道。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:必须为空或 JSON 文档数组,仅在 startup.mode=copy_existing 时使用;它在合成变更事件和通用 pipeline 之前执行。显式设置时覆盖弃用项 copy.existing.pipeline

startup.mode.copy.existing.namespace.regex

使用正则表达式筛选启动时发现的存量复制命名空间。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:匹配完整的 database.collection,使用 Java 正则的子串匹配语义,仅在 startup.mode=copy_existing 时使用。显式设置时覆盖弃用项 copy.existing.namespace.regex

startup.mode.copy.existing.allow.disk.use

控制存量复制聚合是否允许 MongoDB 使用磁盘。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:仅在 startup.mode=copy_existing 时使用;显式设置时覆盖弃用项 copy.existing.allow.disk.use

copy.existing

使用旧式开关启用存量复制。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在 startup.mode 为空时生效。
  • 已弃用:是
  • 替代项startup.mode=copy_existing

copy.existing.max.threads

设置旧式存量复制线程数。
  • 类型int
  • 默认值Runtime.getRuntime().availableProcessors()(运行环境可用处理器数量)
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 1;只有未显式设置替代项时才作为回退值。
  • 已弃用:是
  • 替代项startup.mode.copy.existing.max.threads

copy.existing.queue.size

设置旧式存量复制内存队列容量。
  • 类型int
  • 默认值16000
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 1;只有未显式设置替代项时才作为回退值。
  • 已弃用:是
  • 替代项startup.mode.copy.existing.queue.size

copy.existing.pipeline

设置旧式存量复制聚合管道。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:必须为空或 JSON 文档数组;只有未显式设置替代项时才作为回退值。
  • 已弃用:是
  • 替代项startup.mode.copy.existing.pipeline

copy.existing.namespace.regex

设置旧式存量复制命名空间正则表达式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:必须是有效 Java 正则;只有未显式设置替代项时才作为回退值。
  • 已弃用:是
  • 替代项startup.mode.copy.existing.namespace.regex

copy.existing.allow.disk.use

设置旧式存量复制聚合是否允许使用磁盘。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:只有未显式设置替代项时才作为回退值。
  • 已弃用:是
  • 替代项startup.mode.copy.existing.allow.disk.use

错误处理与心跳

errors.tolerance

设置 Kafka Connect 框架及 MongoDB Connector 的错误容忍模式。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可用值为 noneall,不区分大小写。all 可跳过转换失败并允许部分 Change Stream 恢复分支无 Resume Token 重建;mongo.errors.tolerance 仅覆盖 Connector 自身对该值的读取。

mongo.errors.tolerance

仅覆盖 MongoDB Connector 内部的错误容忍模式。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可用值为 noneall;只要显式设置,就优先于 Connector 对 errors.tolerance 的读取,但不会替代 Kafka Connect 框架的同名设置。

errors.log.enable

控制 Kafka Connect 框架及 MongoDB Connector 是否记录被容忍的错误。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:容忍模式为 none 时 Connector 仍会记录转换错误;mongo.errors.log.enable 可仅覆盖 Connector 自身行为。

mongo.errors.log.enable

仅覆盖 MongoDB Connector 是否记录被容忍的转换错误。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:显式设置时优先于 Connector 对 errors.log.enable 的读取,不改变 Kafka Connect 框架日志设置。

errors.deadletterqueue.topic.name

设置 MongoDB 转换失败记录写入的 DLQ Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 Connector 的有效容忍模式为 all 时使用;Connector 不创建或校验 Topic。mongo.errors.deadletterqueue.topic.name 可覆盖 Connector 对该值的读取。

mongo.errors.deadletterqueue.topic.name

仅覆盖 MongoDB Connector 转换失败记录使用的 DLQ Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:要求 Connector 的有效错误容忍模式为 all;仅影响 Connector 转换错误处理。

heartbeat.interval.ms

设置无业务记录时生成心跳记录的最小间隔。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 0,单位为毫秒;0 表示禁用。只有出现新的 Post-batch Resume Token 且当前 Poll 没有业务记录时才会生成心跳。

heartbeat.topic.name

指定心跳记录写入的 Kafka Topic。
  • 类型string
  • 默认值__mongodb_heartbeats
  • 重要级别:中
  • 有效值 / 注意事项:必须为非空字符串,仅在 heartbeat.interval.ms 大于 0 且生成心跳时使用;Topic 需预先满足命名、权限和可用性要求。

自定义凭据

mongo.custom.auth.mechanism.enable

启用自定义 MongoDB 凭据提供器。
  • 类型string
  • 默认值:无
  • 重要级别:未声明
  • 有效值 / 注意事项:仅字符串值 true(不区分大小写)会启用;该属性未注册到 ConfigDef,因此不会出现在 ConfigDef 元数据中,也不会接受 ConfigDef 类型校验。

mongo.custom.auth.mechanism.providerClass

指定自定义 MongoDB 凭据提供器实现类。
  • 类型string
  • 默认值:无
  • 重要级别:未声明
  • 有效值 / 注意事项:启用自定义认证时必填。类必须可加载、实现 CustomCredentialProvider、具有可访问的无参构造方法,并在所有可运行 Task 的 Worker 上提供实现及依赖。

最佳实践

首次接入已有集合并持续采集增量

适用业务场景:目标集合已经包含存量文档,需要先建立 Kafka 数据基线,再持续接收复制期间及之后产生的新变更。该模式适合首次接入,不适合作为每次重启都重复执行的初始化流程。 在快速开始配置基础上使用以下完整 Connector 配置:
关键说明:Connector 启动时复制所选范围内已存在的文档,并衔接 Change Stream。复制中断后重启会重新执行整轮复制,消费者应使用稳定业务键或其他方式处理重复;复制不是跨多个集合的一致性快照。

为业务命名空间设置稳定的 Topic 名称

适用业务场景:MongoDB 的数据库或集合名称不适合作为下游 Topic 契约,或者迁移过程中需要保持现有 Topic 名称不变。显式映射可将源端命名与 Kafka 消费接口解耦。 使用以下 Connector 配置将一个命名空间路由到固定 Topic:
关键说明:映射在前缀和后缀处理之前完成。调整映射不会改变 Source Offset 身份,但必须提前确认目标 Topic 的命名、权限、分区策略以及是否会与其他映射结果冲突。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 每个 Connector 实例只创建一个 Source Task;增大 tasks.max 不能水平扩展 Change Stream 消费。
  • copy_existing 中途重启会从头重新复制,已写入 Kafka 的存量记录可能重复;该过程也不是跨集合的一致性快照。
  • 多命名空间复制可能交错,Topic 映射可以拆分数据流,且 Connector 不指定 Kafka Partition,因此不保证跨命名空间、Topic 或 Partition 的全局顺序。
  • publish.full.document.only=true 会过滤没有 fullDocument 的事件,除非同时启用 Tombstone 选项生成空值记录。
  • 错误容忍设为 all 且未配置 DLQ 时,转换失败记录会被丢弃,后续 Offset 可能永久越过该事件。
  • Resume Token 无效或历史已丢失时,容忍分支可能无 Token 重建 Change Stream,无法补回已经离开 Oplog 的数据区间。
  • Connector 不提供 MongoDB 到 Kafka 的端到端 Exactly-once、无重复或无丢失保证;Offset 提交窗口、存量复制和错误容忍均可能产生重复或数据缺口。

常见问题

Connector 运行后没有收到新消息

检查 MongoDB 部署是否支持 Change Streams、连接账号是否具有所选范围的权限,以及 databasecollection 是否指向实际写入的位置。确认默认 Topic <database>.<collection>topic.namespace.map 的目标 Topic 正确;若使用 pipeline,还要确认它没有过滤全部事件或移除 _idns 等必需字段。

更新事件中没有最新完整文档

默认 Change Stream 事件不一定包含更新后的完整文档。需要完整文档时设置 change.stream.full.document=updateLookup;如果只希望发布完整文档,可设置 publish.full.document.only=true,但需同时评估没有 fullDocument 的事件会被过滤,以及删除事件是否需要 Tombstone。

重启后从意外位置开始采集

Connector 优先使用已提交的 Resume Token,只有没有可用 Offset 时才应用 startup.mode。检查是否修改了连接主机、databasecollectionoffset.partition.name,这些变化可能产生新的 Offset 身份;恢复旧配置后确认原 Offset 仍然存在,再重启 Connector。

存量复制后出现重复记录

存量复制期间的文档会作为插入事件发布,复制中断重启会重新执行整轮复制,同时发生的更新还可能作为后续 Change Stream 事件出现。消费者应根据稳定的文档键实现幂等处理,并在重启前判断是否需要继续使用原 Offset 或重新建立数据基线。

转换错误没有进入 DLQ

确认有效的 Connector 错误容忍模式为 all,DLQ Topic 已存在且 Worker 有写权限,并检查是否设置了 mongo.errors.* 覆盖项。Topic 映射为空、仅发布完整文档时缺少 fullDocument 等过滤行为不属于转换失败,不会写入 DLQ。