概述
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。<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、键和值映射所需的
_id、ns、documentKey或fullDocument字段;存量复制的高效筛选优先使用startup.mode.copy.existing.pipeline。
batch.size
设置 MongoDB Change Stream Cursor 的批次大小提示。
- 类型:
int - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:必须至少为
0;0表示沿用驱动和服务端行为,与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 - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:可用值为
default、off、whenAvailable、required或空字符串。需要 MongoDB 6.0 或更高版本并为集合启用 Pre-image;不影响存量复制记录。
change.stream.full.document
设置 Change Stream 返回完整文档的方式。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:可用值为
default、updateLookup、whenAvailable、required或空字符串;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 - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:支持完整命名空间、数据库、按声明顺序匹配的
/正则/和*映射;映射结果仍会应用前缀和后缀。模板可使用db、sep、coll、sep_coll、coll_sep和sep_coll_sep变量;需自行避免非法名称和 Topic 冲突。
输出格式与 Schema
output.format.key
设置消息键的输出格式。
- 类型:
string - 默认值:
json - 重要级别:高
- 有效值 / 注意事项:可用值为
json、bson、schema,不区分大小写;应与key.converter兼容。
output.format.value
设置消息值的输出格式。
- 类型:
string - 默认值:
json - 重要级别:高
- 有效值 / 注意事项:可用值为
json、bson、schema,不区分大小写;应与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 - 默认值:内置
ChangeStreamAvro 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处理 - 重要级别:中
- 有效值 / 注意事项:可用值为
latest、timestamp、copy_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 - 重要级别:中
- 有效值 / 注意事项:可用值为
none、all,不区分大小写。all可跳过转换失败并允许部分 Change Stream 恢复分支无 Resume Token 重建;mongo.errors.tolerance仅覆盖 Connector 自身对该值的读取。
mongo.errors.tolerance
仅覆盖 MongoDB Connector 内部的错误容忍模式。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可用值为
none、all;只要显式设置,就优先于 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 配置:为业务命名空间设置稳定的 Topic 名称
适用业务场景:MongoDB 的数据库或集合名称不适合作为下游 Topic 契约,或者迁移过程中需要保持现有 Topic 名称不变。显式映射可将源端命名与 Kafka 消费接口解耦。 使用以下 Connector 配置将一个命名空间路由到固定 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、连接账号是否具有所选范围的权限,以及database 和 collection 是否指向实际写入的位置。确认默认 Topic <database>.<collection> 或 topic.namespace.map 的目标 Topic 正确;若使用 pipeline,还要确认它没有过滤全部事件或移除 _id、ns 等必需字段。
更新事件中没有最新完整文档
默认 Change Stream 事件不一定包含更新后的完整文档。需要完整文档时设置change.stream.full.document=updateLookup;如果只希望发布完整文档,可设置 publish.full.document.only=true,但需同时评估没有 fullDocument 的事件会被过滤,以及删除事件是否需要 Tombstone。
重启后从意外位置开始采集
Connector 优先使用已提交的 Resume Token,只有没有可用 Offset 时才应用startup.mode。检查是否修改了连接主机、database、collection 或 offset.partition.name,这些变化可能产生新的 Offset 身份;恢复旧配置后确认原 Offset 仍然存在,再重启 Connector。
存量复制后出现重复记录
存量复制期间的文档会作为插入事件发布,复制中断重启会重新执行整轮复制,同时发生的更新还可能作为后续 Change Stream 事件出现。消费者应根据稳定的文档键实现幂等处理,并在重启前判断是否需要继续使用原 Offset 或重新建立数据基线。转换错误没有进入 DLQ
确认有效的 Connector 错误容忍模式为all,DLQ Topic 已存在且 Worker 有写权限,并检查是否设置了 mongo.errors.* 覆盖项。Topic 映射为空、仅发布完整文档时缺少 fullDocument 等过滤行为不属于转换失败,不会写入 DLQ。