概述
Debezium MongoDB CDC Source Connector 从支持变更流的 MongoDB 部署读取文档变化,并将文档的新增、更新和删除事件发布到 Kafka Topic。Topic 名称由逻辑前缀、数据库和集合组成;初始快照会先发布匹配集合中的现有文档,随后继续发布变更流事件。快照记录使用 READ 操作,新增、更新和删除分别对应 CREATE、UPDATE 和 DELETE 操作;删除事件默认还会产生 tombstone 记录。 它适用于将 MongoDB 中的业务数据实时送入数据湖、搜索、分析或下游服务,也适用于已有数据初始化后持续同步的 CDC 链路。源端恢复依赖 Kafka Connect 持久化的 Offset 以及 MongoDB 仍保留对应的变更流历史;消费者应按稳定键或来源位置处理可能的重复事件。前置条件
- MongoDB 必须部署为支持 Change Streams 的副本集或其他受支持的集群部署;独立模式 MongoDB 不能作为该 Connector 的变更流来源。
- 用于连接 MongoDB 的账号必须具备读取目标数据库和集合的权限;如果启用信号或增量快照,还必须具备读取信号集合的权限。写入信号由负责发出信号的外部客户端完成,该客户端需要相应的写入权限。
- 使用 TLS 时,Kafka Connect Worker 必须能够读取相应的信任库、密钥库及其密码,并且 MongoDB 服务端证书的主机名必须能够通过校验。
授权许可
使用 Apache License 2.0。快速开始
准备可访问的 Kafka Connect 集群、Kafka 集群和 MongoDB 变更流源端,确认 Worker 能访问 MongoDB 和 Connector 插件目录,并准备一个有权读取源端的账号。创建和管理 Connector 时,参见 AutoMQ 的管理 Connector。<database-host>、<database-port>、<replica-set-name>、<mongodb-user>、<mongodb-password> 和 <logical-topic-prefix> 替换为实际值。默认 snapshot.mode=initial 时,如果没有可用的既有 Offset,Connector 会先对匹配到的源数据执行初始快照,再继续读取变更流;默认的 MongoDB 连接认证源是 admin。Converter 的具体格式应与消费者约定一致。
配置
连接与认证
topic.prefix
为 Connector 生成的 MongoDB Topic 提供唯一的逻辑前缀。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:只能包含字母、数字、连字符、点和下划线;同一 Kafka 集群中的不同 MongoDB Source 实例应使用不同前缀。
mongodb.connection.string
MongoDB 连接字符串。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:必须是有效的 MongoDB ConnectionString;Connector 校验配置时会建立实时连接进行验证。副本集连接应在连接字符串中提供副本集信息。
mongodb.user
MongoDB 用户名。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:高
- 必填:否
- 有效值 / 注意事项:当凭证已嵌入连接字符串,或使用自定义认证提供方时可以省略。
mongodb.password
MongoDB 密码。
- 类型:
PASSWORD - 默认值:无固定默认值
- 重要级别:高
- 必填:否
- 有效值 / 注意事项:可与
mongodb.user配合使用;不要将真实密码写入日志、文档示例或版本库。
mongodb.authsource
保存 MongoDB 用户凭证的认证数据库。
- 类型:
STRING - 默认值:
admin - 重要级别:中
- 有效值 / 注意事项:当用户不在
admin数据库中创建时,设置为实际认证数据库。
mongodb.ssl.enabled
是否为 MongoDB 连接启用 TLS/SSL。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:启用后应同时准备服务端证书校验所需的信任配置。
mongodb.ssl.invalid.hostname.allowed
是否允许 MongoDB 服务端证书的主机名校验失败。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true会关闭主机名保护并带来中间人攻击风险,仅适用于受控诊断环境。
mongodb.ssl.keystore
用于双向 TLS 认证的客户端密钥库路径。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 必填:否
- 有效值 / 注意事项:仅在 MongoDB 要求客户端证书时设置。
mongodb.ssl.keystore.password
客户端密钥库密码。
- 类型:
PASSWORD - 默认值:无固定默认值
- 重要级别:中
- 必填:否
- 有效值 / 注意事项:密钥库受密码保护时必须设置。
mongodb.ssl.keystore.type
客户端密钥库类型。
- 类型:
STRING - 默认值:
PKCS12 - 重要级别:中
- 有效值 / 注意事项:仅在设置
mongodb.ssl.keystore时生效,并应与密钥库实际格式一致。
mongodb.ssl.truststore
用于验证 MongoDB 服务端证书的信任库路径。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 必填:否
- 有效值 / 注意事项:在 Worker 不使用默认信任库,或需要专用 CA 时设置。
mongodb.ssl.truststore.password
信任库密码。
- 类型:
PASSWORD - 默认值:无固定默认值
- 重要级别:中
- 必填:否
- 有效值 / 注意事项:信任库需要解锁或验证密码时设置。
mongodb.ssl.truststore.type
信任库类型。
- 类型:
STRING - 默认值:
PKCS12 - 重要级别:中
- 有效值 / 注意事项:仅在设置
mongodb.ssl.truststore时生效,并应与信任库实际格式一致。
mongodb.connect.timeout.ms
建立 MongoDB 连接的超时时间,单位为毫秒。
- 类型:
INT - 默认值:
10000 - 重要级别:低
- 有效值 / 注意事项:应根据网络延迟和部署位置调整;过小可能导致启动或重连失败。
mongodb.server.selection.timeout.ms
MongoDB 驱动选择可用服务端的超时时间,单位为毫秒。
- 类型:
INT - 默认值:
30000 - 重要级别:低
- 有效值 / 注意事项:服务端故障转移或跨网络部署时可适当增大。
mongodb.socket.timeout.ms
MongoDB Socket 操作超时时间,单位为毫秒。
- 类型:
INT - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:
0表示不设置 Socket 超时。
mongodb.heartbeat.frequency.ms
MongoDB 驱动集群监控心跳频率,单位为毫秒。
- 类型:
INT - 默认值:
10000 - 重要级别:低
- 有效值 / 注意事项:影响驱动发现集群状态变化的节奏。
mongodb.poll.interval.ms
轮询 MongoDB 副本集成员变化的间隔,单位为毫秒。
- 类型:
LONG - 默认值:
30000 - 重要级别:中
- 有效值 / 注意事项:必须为正数;需要更快发现拓扑变化时可缩短,但会增加轮询开销。
connection.validation.timeout.ms
Connector 配置校验阶段验证 MongoDB 连接的超时时间,单位为毫秒。
- 类型:
LONG - 默认值:
60000 - 重要级别:低
- 有效值 / 注意事项:必须为正数;它只控制连接验证,不替代运行阶段的连接超时配置。
源范围与字段
database.include.list
指定要捕获的数据库范围。
- 类型:
LIST - 默认值:无固定默认值
- 重要级别:高
- 有效值 / 注意事项:默认按正则表达式匹配;不能与
database.exclude.list同时设置。使用名称匹配模式时按字面值匹配。
database.exclude.list
指定不捕获的数据库范围。
- 类型:
LIST - 默认值:无固定默认值
- 重要级别:高
- 有效值 / 注意事项:默认按正则表达式匹配;不能与
database.include.list同时设置。
collection.include.list
指定要捕获的集合范围。
- 类型:
LIST - 默认值:无固定默认值
- 重要级别:高
- 有效值 / 注意事项:使用完整的
database.collection名称或模式;不能与collection.exclude.list同时设置。
collection.exclude.list
指定不捕获的集合范围。
- 类型:
LIST - 默认值:无固定默认值
- 重要级别:高
- 有效值 / 注意事项:使用完整的
database.collection名称或模式;不能与collection.include.list同时设置。
field.exclude.list
从事件中排除指定字段。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:使用逗号分隔的
database.collection.field或嵌套字段名称;数据库和集合支持*。
field.renames
重命名事件中的字段。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:使用逗号分隔的
旧字段:新字段映射,并使用完整字段限定名;嵌套字段名称必须符合字段映射语法。
快照与初始化
snapshot.mode
选择初始快照和后续变更流的启动方式。
- 类型:
STRING - 默认值:
initial - 重要级别:低
- 有效值 / 注意事项:可选
always、initial、never、no_data、initial_only、when_needed、configuration_based或custom。never已弃用,使用no_data替代。initial_only完成快照后不会继续持续捕获变更;no_data和never要求已有源状态及解码所需的结构前提。
snapshot.collection.filter.overrides
为指定集合声明快照读取过滤条件。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:使用逗号分隔的
database.collection名称;每个名称都必须配套设置同名的snapshot.collection.filter.overrides.<db>.<collection>属性。
snapshot.delay.ms
开始快照前等待的时间,单位为毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:必须为非负数;适合需要让源端或下游准备完成后再初始化的场景。
streaming.delay.ms
快照完成后开始读取变更流前等待的时间,单位为毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:必须为非负数;只影响快照到持续捕获之间的等待。
snapshot.fetch.size
快照读取时使用的批量大小。
- 类型:
INT - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:必须为非负整数;应结合文档大小、MongoDB 负载和 Worker 内存调整。
snapshot.max.threads
快照阶段用于并行读取集合的最大线程数。
- 类型:
INT - 默认值:
1 - 重要级别:中
- 有效值 / 注意事项:必须为正数;实际线程数不超过集合数量。它增加的是单个 Connector 内部的快照并行度,不增加 Kafka Connect Task 数量。
snapshot.include.collection.list
指定需要执行快照的集合范围。
- 类型:
LIST - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:使用正则表达式列表;它用于限制快照集合,不等同于持续变更流的捕获范围。
incremental.snapshot.chunk.size
增量快照每个分块读取的文档数量。
- 类型:
INT - 默认值:
1024 - 重要级别:中
- 有效值 / 注意事项:必须为非负数;较小分块降低单次读取压力,但会增加增量快照管理开销。
incremental.snapshot.watermarking.strategy
选择增量快照的水位标记策略。
- 类型:
STRING - 默认值:
INSERT_INSERT - 重要级别:低
- 有效值 / 注意事项:可选
INSERT_INSERT或INSERT_DELETE;必须启用信号能力才能使用信号驱动的增量快照。
缓冲与发布
max.batch.size
单次从内部队列取出并发布的最大记录数。
- 类型:
INT - 默认值:
2048 - 重要级别:中
- 有效值 / 注意事项:必须为正数,且应小于
max.queue.size。
max.queue.size
内部事件队列允许缓存的最大记录数。
- 类型:
INT - 默认值:
8192 - 重要级别:中
- 有效值 / 注意事项:必须大于
max.batch.size;增大后可以吸收短时波动,但会增加内存占用和故障前缓存量。
max.queue.size.in.bytes
内部事件队列允许缓存的最大字节数。
- 类型:
LONG - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:必须为非负数;
0表示不启用按字节限制。MongoDB 文档大小差异较大时可用它约束内存。
poll.interval.ms
Connector 从内部队列轮询记录的间隔,单位为毫秒。
- 类型:
LONG - 默认值:
500 - 重要级别:中
- 有效值 / 注意事项:必须为正数;较小值通常降低发布等待,但会增加轮询频率。
query.fetch.size
快照或查询读取时使用的驱动抓取大小。
- 类型:
INT - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:必须为非负整数;
0使用驱动或数据库默认行为。
事件与心跳
tombstones.on.delete
是否在删除事件后发送 tombstone 记录。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:
true发送 DELETE 和 tombstone;false只发送 DELETE。下游压缩 Topic 或状态存储应按其删除约定选择。
skipped.operations
指定不发布的操作类型。
- 类型:
LIST - 默认值:
[t] - 重要级别:低
- 有效值 / 注意事项:可选操作码
c、u、d、t或none;默认跳过 truncate 操作。
heartbeat.interval.ms
发送心跳事件的间隔,单位为毫秒。
- 类型:
INT - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:必须为非负数;
0禁用心跳事件。低流量源可启用心跳以提高进度和 Offset 可见性。
heartbeat.topics.prefix
心跳 Topic 的名称前缀。
- 类型:
STRING - 默认值:
__debezium-heartbeat - 重要级别:低
- 有效值 / 注意事项:只有启用心跳事件时才需要特别关注;应与 Topic 命名和权限策略一致。
extended.headers.enabled
是否在源记录中写入 Debezium 扩展上下文 Header。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:关闭后下游将无法依赖这些扩展 Header 获取相应上下文。
信号与增量快照
signal.data.collection
存放 Connector 信号的 MongoDB 数据库和集合。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:中
- 有效值 / 注意事项:格式为
database.collection;未设置时禁用信号处理。使用增量快照前应创建该集合并授予相应权限。
signal.poll.interval.ms
轮询信号集合的间隔,单位为毫秒。
- 类型:
LONG - 默认值:
5000 - 重要级别:中
- 有效值 / 注意事项:必须为正数;越小越快发现信号,但会增加源端读取频率。
signal.enabled.channels
启用的信号通道列表。
- 类型:
LIST - 默认值:
[source] - 重要级别:中
- 有效值 / 注意事项:默认启用
source通道;只有配置了对应信号源和权限时才应增加其他通道。
Topic、Schema 与扩展
topic.naming.strategy
生成 Kafka Topic 名称的策略类。
- 类型:
CLASS - 默认值:
io.debezium.schema.DefaultTopicNamingStrategy - 重要级别:中
- 有效值 / 注意事项:类必须实现
TopicNamingStrategy;更换策略会改变 Topic 路由,可能影响下游订阅和历史数据兼容性。
schema.name.adjustment.mode
调整 Schema 名称以适配特定序列化格式的模式。
- 类型:
STRING - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:可选
none、avro或avro_unicode;选择前确认下游 Schema Registry 或 Converter 的命名要求。
sourceinfo.struct.maker
创建 MongoDB Source 信息结构的实现类。
- 类型:
CLASS - 默认值:
io.debezium.connector.mongodb.MongoDbSourceInfoStructMaker - 重要级别:低
- 有效值 / 注意事项:自定义类必须提供兼容的
SourceInfoStructMaker实现;除非有明确的事件元数据扩展需求,否则保留默认值。
converters
注册自定义 Converter 及其配置前缀。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:低
- 有效值 / 注意事项:使用逗号分隔的自定义 Converter 名称,并为每个 Converter 提供对应的类型和选项;自定义类必须在 Worker 插件路径中可见。
post.processors
注册事件后处理器及其配置前缀。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:低
- 有效值 / 注意事项:使用逗号分隔的后处理器名称,并提供各自的类型和选项;后处理器会改变发布事件的内容或元数据。
错误处理
event.processing.failure.handling.mode
处理损坏或无法解析事件的方式。
- 类型:
STRING - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:可选
fail、warn或ignore。warn和ignore可能跳过无法处理的事件,应配合错误日志和业务完整性检查使用。
errors.max.retries
可重试源端错误的最大重试次数。
- 类型:
INT - 默认值:
-1 - 重要级别:低
- 有效值 / 注意事项:
-1表示不限次数,0表示禁用重试,正数表示有限重试;重试不是端到端无重复交付保证。
errors.retry.timeout
Kafka Connect Worker 对失败记录或任务进行重试的总时长,单位为毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:
0表示禁用该重试窗口,-1表示无限重试;它属于 Worker 错误策略,应与任务失败告警一起规划。
errors.tolerance
Kafka Connect Worker 对错误的容忍级别。
- 类型:
STRING - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可选
none或all;all会跳过问题记录,不能代替错误分析和补偿流程。
Kafka Connect 框架
connector.class
指定 Kafka Connect 要实例化的 Connector 类。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:使用
io.debezium.connector.mongodb.MongoDbConnector。
tasks.max
Connector 可使用的最大 Kafka Connect Task 数量。
- 类型:
INT - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1;该 Connector 为 MongoDB 复制返回一个 Task 配置,增大此值不会为每个副本集创建独立的 Kafka Connect Task。副本集复制和快照并行度由 Connector 内部线程处理。
key.converter
序列化源记录 Key 的 Converter 类。
- 类型:
CLASS - 默认值:无固定默认值
- 重要级别:低
- 有效值 / 注意事项:必须实现 Kafka Connect
Converter;未在 Connector 级别设置时使用 Worker 默认值。
value.converter
序列化 MongoDB CDC 事件 Value 的 Converter 类。
- 类型:
CLASS - 默认值:无固定默认值
- 重要级别:高
- 有效值 / 注意事项:必须实现 Kafka Connect
Converter;应与下游消费者的反序列化方式一致。
最佳实践
首次接入时限定初始化范围并继续增量同步
适用业务场景:首次接入已有 MongoDB 数据,需要先建立指定数据库和集合的存量基线,再持续接收后续变化,避免把无关集合的存量数据一并发布。 配置示例:snapshot.mode=initial 会发布匹配集合的 READ 记录并继续捕获变更;后续新增集合不会自动补发其已有存量,若需要补发应使用单独的快照规划。collection.include.list 与 collection.exclude.list 不能同时配置。
已有数据基线时只接入后续变化
适用业务场景:下游已经通过备份、批处理或其他方式建立了数据基线,只需要从 Connector 启动位置开始接收后续 MongoDB 变化,不希望再次发布现有文档。 配置示例:no_data 不发送现有文档的快照记录,但仍要求下游具备解析后续事件所需的结构和基线。切换到该模式前应明确下游基线对应的时间点,并保留稳定的 topic.prefix 和 Worker Offset;不要用它绕过未完成的初始快照或补回已经过期的 MongoDB 变更历史。
低流量源持续运行时提高进度可见性
适用业务场景:业务集合平时写入较少,但仍需要确认 Connector 持续读取、维护源端进度,并在下游监控中区分“没有变化”和“没有运行”。 配置示例:监控
监控内容
监控 Kafka Connect Worker、Connector 和 Task 的健康状态及重启次数,检查源端到 Kafka 的吞吐、端到端延迟、积压、Offset 提交时间和提交失败;同时关注 MongoDB 连接错误、重试、任务错误日志、事件处理失败、Kafka 发布失败以及 Worker JVM 的堆内存、GC、线程和 CPU。启用errors.tolerance=all 或错误处理相关 DLQ 配置时,还应监控被跳过记录和 DLQ 流量,避免把错误容忍误判为正常同步。
导入 Grafana 大盘
下载 AutoMQ Connect Grafana Dashboard,在 Grafana 中选择已经采集 Kafka Connect 指标的数据源,确认集群、Worker、Connector 和 Task 标签与指标采集配置一致,然后通过 Dashboards 的导入功能上传 JSON 并映射数据源。限制条件
- 独立模式 MongoDB 不提供该 Connector 所需的 Change Streams,因此不能作为本 Connector 的变更流来源。
tasks.max增大不会按副本集拆分为多个 Kafka Connect Task;副本集复制使用 Connector 内部线程,快照线程也不会改变 Task 数量。initial_only只执行快照,不会在快照结束后继续持续捕获变更。no_data和已弃用的never不会发布现有文档快照,并要求已有源状态及下游解码所需的结构前提。- 增量快照的窗口去重用于处理快照与并发变化的碰撞,不构成端到端 exactly-once 保证,也不代表单一全局时间点快照。
- 变更流管道不能使用 MongoDB 索引;使用正则范围过滤时还会增加内部字段变换开销。
- 连接重启或 Offset 恢复时,如果 MongoDB 已清理保存的变更流历史,Connector 不能从旧位置无间隙地重放所有中间事件。
- 内部缓冲区不是持久化检查点;增大队列只会改变内存、背压和延迟行为,不改变交付保证。
- 在 Kafka 发布完成但 Worker Offset 尚未持久化的故障窗口内,重启后可能重新发布已经发送过的事件;下游应能处理重复记录。
常见问题
为什么 Connector 连接 MongoDB 失败?
检查 MongoDB 是否为支持 Change Streams 的副本集或集群部署,连接字符串是否包含正确的主机、端口和副本集信息,认证数据库是否由mongodb.authsource 正确指定,以及账号是否具备目标数据库读取权限。启用 TLS 时,再检查信任库、证书主机名和 Worker 文件权限;不要通过 mongodb.ssl.invalid.hostname.allowed=true 绕过生产环境证书校验。
为什么启动后没有收到现有文档?
确认snapshot.mode 没有设置为 no_data、never 或其他不执行数据快照的模式,并检查 database.include.list、collection.include.list、排除列表和 snapshot.include.collection.list 是否排除了目标集合。若下游已经保存过同一 topic.prefix 的 Offset,重新启动不会因为重启自动重新发送完整存量;需要按数据基线和快照策略重新规划。
为什么只收到部分集合的事件?
检查数据库和集合包含列表、排除列表及其正则表达式,确认没有同时配置同一层级的 include 和 exclude。还要检查snapshot.include.collection.list 是否只限制了快照,避免把快照范围和持续变更流范围混为一谈。
为什么删除事件后又出现一条空值记录?
这是tombstones.on.delete=true 的默认行为:Connector 先发布 DELETE,再发布同一 Key 的 tombstone,便于 Kafka 日志压缩清理旧值。如果下游不接受 tombstone,可设置 tombstones.on.delete=false,但必须同步评估下游删除和状态清理逻辑。
为什么重启后出现重复事件?
Kafka Connect Offset 的持久化可能晚于事件发布;如果故障发生在发布和 Offset 刷新之间,恢复时可能从较早位置重新发布事件。确认topic.prefix、Worker Offset 存储和 MongoDB 变更历史保持不变,并让下游使用稳定文档 Key、来源位置或业务幂等逻辑处理重复。Connector 不承诺端到端 exactly-once。
为什么启用增量快照后没有执行?
确认已设置signal.data.collection,该集合存在且账号具备读取信号的权限,signal.enabled.channels 包含 source,并且增量快照配置使用了支持的水位策略。还要确认信号文档使用正确的数据库和集合名称,以及相关集合具备 Connector 能读取的文档键。