Skip to main content

概述

MongoDB Sink Connector 从 Kafka Topic 消费记录,将消息键和值转换为 BSON,并写入 MongoDB 数据库和集合。默认情况下,Connector 使用配置的数据库,并将 Kafka Topic 名称作为集合名称;也可以固定集合、按记录字段动态路由,或为不同 Topic 覆盖写入配置。 Connector 适合将业务事件、状态更新和时序数据从 Kafka 持续写入 MongoDB。普通写入路径可配置文档 _id、替换或更新策略、字段处理、批量写入和墓碑删除;对于受支持的 CDC 消息,还可以通过专用处理器把变更事件转换为 MongoDB 写操作。

前置条件

  • MongoDB 账号需具有目标数据库和集合所需的创建或写入权限;使用 TLS、客户端证书、Stable API、客户端字段级加密或自定义认证时,还需提前准备对应证书、密钥、密钥库和服务端能力。
  • 使用时序集合或 Stable API 时,MongoDB 版本需为 5.0 或更高;已有时序集合的时间字段、元数据字段和粒度应与 Connector 配置一致。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和 MongoDB,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
<mongodb-host> 替换为 Connect Worker 可访问的 MongoDB 地址。该配置把 mongodb-sink-input 中的 JSON 对象写入 inventory.eventsvalue.converter.schemas.enable 是 JSON Converter 的配置,位于 Connector 的 config 中。

配置

Connector 身份与任务

connector.class

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

tasks.max

Connector 可请求的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1;实际并行度还受 Topic 分区数和 Worker 分配影响。

Kafka Topic 选择

topics

指定要消费的 Kafka Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:与 topics.regex 必须且只能配置一个非空值。

topics.regex

使用 Java 正则表达式选择要消费的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效 Java 正则表达式;与 topics 必须且只能配置一个非空值。

MongoDB 连接与认证

connection.uri

MongoDB Driver 连接字符串,可包含主机、凭据和 URI 选项。
  • 类型password
  • 默认值mongodb://localhost:27017
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效 MongoDB Connection String;不要在文档或日志中暴露凭据。

mongo.custom.auth.mechanism.enable

启用自定义 MongoDB 凭据提供器。
  • 类型boolean-like string
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:该公开扩展未注册到 ConfigDef;仅可被 Boolean.parseBoolean 解析为 true 的值会启用提供器,启用后必须配置 mongo.custom.auth.mechanism.providerClass

mongo.custom.auth.mechanism.providerClass

指定自定义 CustomCredentialProvider 实现类。
  • 类型class-name string
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:启用自定义认证时必填;类必须位于插件 Classpath 中、实现 CustomCredentialProvider,并提供可访问的无参构造方法。

TLS

connection.ssl.truststore

MongoDB TLS 连接使用的 Truststore 文件。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时文件必须可读,并可使用 connection.ssl.truststorePassword 加载。

connection.ssl.truststorePassword

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

connection.ssl.keystore

客户端 TLS 凭据使用的 Keystore 文件。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时文件必须可读,并可使用 connection.ssl.keystorePassword 加载。

connection.ssl.keystorePassword

Keystore 及其中私钥的密码。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在 connection.ssl.keystore 非空时使用。

MongoDB Stable API

server.api.version

指定 MongoDB Stable API 版本。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:空值表示禁用;非空值必须是 Driver 支持的 ServerApiVersion,并要求 MongoDB 5.0 或更高版本。

server.api.deprecation.errors

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

server.api.strict

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

客户端字段级加密

csfle.enabled

为 Sink 写入启用 MongoDB 客户端字段级加密。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:设为 true 时必须配置 csfle.key.vault.namespacecsfle.local.master.key

csfle.key.vault.namespace

指定 database.collection 格式的密钥库命名空间。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项csfle.enabled=true 时必须为非空值。

csfle.local.master.key

指定 Base64 编码的 96 字节本地 KMS 主密钥。
  • 类型password
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项csfle.enabled=true 时必填,运行时会进行 Base64 解码。

csfle.schema.map

指定控制自动字段加密的 JSON Schema Map。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时必须能解析为从命名空间到 Schema 的 BSON 文档。

csfle.bypass.query.analysis

跳过自动 CS-FLE 查询分析。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:仅在启用 CS-FLE 时相关;设为 false 需要 mongocryptdcrypt_shared 等自动加密支持。

目标命名空间

database

默认目标 MongoDB 数据库。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须为非空字符串;使用字段路径映射器时可按记录替换。
  • 必填:是

collection

默认目标 MongoDB 集合。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:使用默认命名空间映射器时,空值表示使用 Kafka Topic 名称。

namespace.mapper

指定选择目标数据库和集合的 NamespaceMapper 实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.sink.namespace.mapping.DefaultNamespaceMapper
  • 重要级别:高
  • 有效值 / 注意事项:必须是实现 NamespaceMapper 且具有公开无参构造方法的完整类名;字段路径配置要求使用 com.mongodb.kafka.connect.sink.namespace.mapping.FieldPathNamespaceMapper

namespace.mapper.key.database.field

从键文档字段路径读取目标数据库名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:要求使用 FieldPathNamespaceMapper,且不能与 namespace.mapper.value.database.field 同时配置。

namespace.mapper.key.collection.field

从键文档字段路径读取目标集合名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:要求使用 FieldPathNamespaceMapper,且不能与 namespace.mapper.value.collection.field 同时配置。

namespace.mapper.value.database.field

从值文档字段路径读取目标数据库名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:要求使用 FieldPathNamespaceMapper,且不能与 namespace.mapper.key.database.field 同时配置。

namespace.mapper.value.collection.field

从值文档字段路径读取目标集合名称。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:要求使用 FieldPathNamespaceMapper,且不能与 namespace.mapper.key.collection.field 同时配置。

namespace.mapper.error.if.invalid

控制命名空间字段缺失或不是字符串时是否令记录失败。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:用于 FieldPathNamespaceMapperfalse 回退到默认数据库或集合,true 抛出 DataException

topic.override.<topic>.<property>

为一个 Kafka Topic 覆盖 Topic 级 Sink 配置。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:必须使用实际 Topic 和已注册属性组成具体键,不能直接使用模板键;覆盖值会合并到全局配置之上,不能覆盖 connection.uritopics。使用 topics.regex 时,覆盖项中的 Topic 必须匹配该正则表达式。

记录转换

key.converter

覆盖 Connector 级消息键 Converter。
  • 类型class
  • 默认值:无固定默认值
  • 重要级别:低
  • 有效值 / 注意事项:未配置时继承 Worker 设置;必须是可实例化的 Converter,其子配置使用 key.converter. 前缀。

value.converter

覆盖 Connector 级消息值 Converter。
  • 类型class
  • 默认值:无固定默认值
  • 重要级别:低
  • 有效值 / 注意事项:未配置时继承 Worker 设置;必须是可实例化的 Converter,其子配置使用 value.converter. 前缀。

文档标识

document.id.strategy

指定 DocumentIdAdder 生成 MongoDB _idIdStrategy 实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.sink.processor.id.strategy.BsonOidStrategy
  • 重要级别:高
  • 有效值 / 注意事项:使用实现 IdStrategy 的完整类名;空字符串虽可通过 ConfigDef 校验,但初始化时会失败。配置 CDC 处理器时忽略。

document.id.strategy.overwrite.existing

允许 DocumentIdAdder 覆盖值文档中已有的 _id
  • 类型boolean
  • 默认值false
  • 重要级别:高
  • 有效值 / 注意事项:仅在 DocumentIdAdder 运行时相关;配置 CDC 处理器时忽略。

document.id.strategy.uuid.format

指定 UuidStrategy 生成 UUID 的 BSON 表示形式。
  • 类型string
  • 默认值string
  • 重要级别:高
  • 有效值 / 注意事项:不区分大小写,可选 stringbinary;仅与 UuidStrategy 配合使用。

document.id.strategy.partial.key.projection.type

指定 PartialKeyStrategy 的专用字段投影模式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值回退到 key.projection.type;可用 noneallowlistblocklist,仍兼容已弃用的 whitelistblacklist,应分别改用 allowlistblocklist

document.id.strategy.partial.key.projection.list

指定 PartialKeyStrategy 使用的字段列表。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值回退到 key.projection.list;仅与 PartialKeyStrategy 配合使用。

document.id.strategy.partial.value.projection.type

指定 PartialValueStrategy 的专用字段投影模式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值回退到 value.projection.type;可用 noneallowlistblocklist,仍兼容已弃用的 whitelistblacklist,应分别改用 allowlistblocklist

document.id.strategy.partial.value.projection.list

指定 PartialValueStrategy 使用的字段列表。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值回退到 value.projection.list;仅与 PartialValueStrategy 配合使用。

写入与删除

writemodel.strategy

指定为非空记录构建 MongoDB WriteModel 的实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.sink.writemodel.strategy.DefaultWriteModelStrategy
  • 重要级别:低
  • 有效值 / 注意事项:必须实现 WriteModelStrategy;默认执行替换写入,时序集合改用插入。配置 CDC 处理器时忽略。

delete.on.null.values

将值为 null 的墓碑记录转换为基于消息键的 MongoDB 删除。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:使用默认删除策略时,document.id.strategy 必须为 FullKeyStrategyPartialKeyStrategyProvidedInKeyStrategy;配置 CDC 处理器时忽略。

delete.writemodel.strategy

指定为墓碑记录构建 MongoDB 删除模型的实现类。
  • 类型string
  • 默认值com.mongodb.kafka.connect.sink.writemodel.strategy.DeleteOneDefaultStrategy
  • 重要级别:低
  • 有效值 / 注意事项:仅在 delete.on.null.values=true 时使用;自定义类必须实现 WriteModelStrategy

max.batch.size

限制每个 MongoDB 批量写入的最大记录数。
  • 类型int
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:至少为 00 表示不按此配置拆分同一 Topic 和命名空间的记录组。该值限制记录数,不限制 BSON 字节数。

bulk.write.ordered

控制 MongoDB Bulk Write 是否按有序模式执行。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项:有序模式在首个写错误后停止当前批次后续写入;无序模式可继续当前批次中的其他写入,但不保证记录执行顺序。

rate.limiting.timeout

Connector 触发限速时休眠的毫秒数。
  • 类型int
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:至少为 0;仅当本配置与 rate.limiting.every.n 均大于 0 时启用限速。

rate.limiting.every.n

指定每处理多少个批次触发一次限速休眠。
  • 类型int
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:至少为 0;仅当本配置与 rate.limiting.timeout 均大于 0 时启用限速。

后处理与字段变换

post.processor.chain

指定普通写入路径中依次执行的 PostProcessor 类。
  • 类型list
  • 默认值com.mongodb.kafka.connect.sink.processor.DocumentIdAdder
  • 重要级别:低
  • 有效值 / 注意事项:每项必须是完整实现类名;如果省略 DocumentIdAdder,Connector 仍会自动将其放到链首。配置 CDC 处理器时忽略。

key.projection.type

指定键字段投影模式。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项:可用 noneallowlistblocklistwhitelistblacklist 已弃用,应分别改用 allowlistblocklist。配置 CDC 处理器时忽略。

key.projection.list

指定键字段投影使用的逗号分隔字段路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:按 key.projection.type 解释;模式为 none 或配置 CDC 处理器时忽略。

value.projection.type

指定值字段投影模式。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项:可用 noneallowlistblocklistwhitelistblacklist 已弃用,应分别改用 allowlistblocklist。配置 CDC 处理器时忽略。

value.projection.list

指定值字段投影使用的逗号分隔字段路径。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:按 value.projection.type 解释;模式为 none 或配置 CDC 处理器时忽略。

field.renamer.mapping

指定精确字段路径重命名规则的 JSON 数组。
  • 类型string
  • 默认值[]
  • 重要级别:低
  • 有效值 / 注意事项:每个启用的对象必须包含 oldNamenewName,并在 post.processor.chain 中配置相应重命名处理器。

field.renamer.regexp

指定正则表达式字段重命名规则的 JSON 数组。
  • 类型string
  • 默认值[]
  • 重要级别:低
  • 有效值 / 注意事项:对象必须是有效的正则重命名设置,并在 post.processor.chain 中配置相应处理器;不要使用未定义的 field.renamer.regex

field.value.transformer

指定应用于选定值字段的自定义 FieldValueTransformer
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:非空时必须是具有公开无参构造方法的完整实现类名,Connector 会在普通写入路径自动追加对应后处理器。

field.value.transformer.fields

指定由 field.value.transformer 处理的逗号分隔字段名。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:仅在配置 Transformer 时有意义;按字段名匹配,也会匹配嵌套文档和文档数组中的同名字段,而不是完整路径。

CDC 处理

change.data.capture.handler

指定解释 CDC Envelope 并构建 MongoDB WriteModel 的 CdcHandler 实现类。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:非空时切换到 CDC 路径,并忽略后处理链、字段重命名、键值投影、普通写入模型、文档 ID 配置和 delete.on.null.values

change.data.capture.handler.remove.null.values

从 CDC 替换文档和 Update 的 $set 内容中移除 null 字段。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅在配置 CDC Handler 时生效;保留 $unset,对删除操作无影响。

时序集合

timeseries.timefield

指定 MongoDB 时序集合的顶层 BSON 日期时间字段。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:非空时启用时序模式,要求 MongoDB 5.0 或更高版本;每条记录都需包含该字段,且值为 BSON 日期时间或可被自动转换。

timeseries.metafield

指定时序集合的可选顶层元数据字段。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:非空时要求配置 timeseries.timefield;不能与时间字段或 _id 相同,MongoDB 不允许数组元数据。

timeseries.expire.after.seconds

指定新建时序集合的数据保留秒数。
  • 类型long
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:至少为 0;非零值要求配置 timeseries.timefield0 表示创建集合时不设置 expireAfterSeconds

timeseries.granularity

指定新建时序集合的预期时间间隔粒度。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:可选 secondsminuteshours;非空值要求配置 timeseries.timefield

timeseries.timefield.auto.convert

将 Epoch 毫秒数或字符串时间字段转换为 BSON 日期时间。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅在时序模式中添加转换处理器;转换失败时保留原值,后续写入仍可能失败。

timeseries.timefield.auto.convert.date.format

指定解析字符串时间字段的 DateTimeFormatter Pattern。
  • 类型string
  • 默认值yyyy-MM-dd[['T'][ ]][HH:mm:ss[[.][SSSSSS][SSS]][ ]VV[ ]'['VV']'][HH:mm:ss[[.][SSSSSS][SSS]][ ]X][HH:mm:ss[[.][SSSSSS][SSS]]]
  • 重要级别:低
  • 有效值 / 注意事项:必须能被 DateTimeFormatter.ofPattern 接受,仅在自动转换字符串时间值时使用。

timeseries.timefield.auto.convert.locale.language.tag

指定时间字段 Formatter 使用的 Locale Language Tag。
  • 类型string
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值使用 Locale.ROOT;非空值必须能被 Locale.forLanguageTag 接受。

错误处理与重试

errors.tolerance

配置 Kafka Connect 框架容错,并作为 Connector 内部容错的基础值。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:Kafka Connect 3.9.1 仅接受 noneall;不要在此处使用 Connector ConfigDef 接受但框架会拒绝的 data。启用通用错误继续处理和 DLQ 时使用 all

mongo.errors.tolerance

独立覆盖 Connector 内部容错模式。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 nonedataall;只需容忍 Connector 数据错误时,在 Kafka Connect 3.9.1 上应将 data 配置在此项。

errors.log.enable

启用 Kafka Connect 框架和 Connector 的错误日志。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:设为 true 时也记录被容忍的 Connector 错误;未容忍错误无论此值如何都会记录。

mongo.errors.log.enable

独立覆盖 Connector 内部错误日志设置。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:显式配置后覆盖 errors.log.enable;未容忍的 Connector 错误仍会记录。

errors.retry.timeout

配置 Kafka Connect 框架中可重试处理阶段的总重试时长,单位毫秒。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 禁用框架重试,-1 表示无限重试;该配置不保证重试 MongoDB Bulk Write,也不会恢复已移除的 max.num.retriesretries.defer.timeout 行为。MongoDB Driver 的可重试写入由 connection.uri 等 Driver 选项控制。

errors.retry.delay.max.ms

配置连续框架重试之间的最大延迟,单位毫秒。
  • 类型long
  • 默认值60000
  • 重要级别:中
  • 有效值 / 注意事项:仅在 errors.retry.timeout 允许重试时相关;达到上限后 Kafka Connect 会加入随机抖动。

死信队列

errors.deadletterqueue.topic.name

指定接收错误 Sink 记录的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:空值禁用 DLQ;该 Topic 不能出现在 topics 中或匹配 topics.regex。通常同时配置 errors.tolerance=all 和相应 Connector 容错模式。

errors.deadletterqueue.topic.replication.factor

指定 Kafka Connect 自动创建 DLQ Topic 时使用的副本因子。
  • 类型short
  • 默认值3
  • 重要级别:中
  • 有效值 / 注意事项:仅在 DLQ Topic 尚不存在时使用,值必须受 Kafka 集群支持。

errors.deadletterqueue.context.headers.enable

为 DLQ 记录添加 __connect.errors.* 上下文 Header。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在配置 DLQ 时相关;Header 可能包含运行元数据,应按数据治理要求处理。

最佳实践

使用稳定记录身份控制恢复后的写入结果

适用业务场景:Connector 已开始持续写入,业务希望 Worker 重启或 Offset 重放时,同一 Kafka 记录仍定位到同一 MongoDB 文档,而不是因为默认随机 ObjectId 再创建一份文档。此策略适合按 Kafka 记录身份落库、不要求多个不同事件合并到同一业务实体的场景。 在快速开始配置基础上使用以下完整 Connector 配置:
KafkaMetaDataStrategy 根据 Topic、Partition 和 Offset 生成稳定 _id,配合默认替换并 Upsert 的写入模型,可使同一记录的重放定位到同一文档。它不提供端到端 Exactly-Once:MongoDB 写入与 Kafka Offset 提交之间没有原子事务,触发器、校验和其他外部副作用仍可能再次执行。

按分区和写入压力扩展吞吐

适用业务场景:基础链路已运行,需要利用多个 Kafka 分区提高并行度,并限制单次 MongoDB Bulk Write 的记录数,以平衡吞吐、请求大小、延迟和失败影响范围。 在快速开始配置基础上使用以下完整 Connector 配置,并按实际 Topic 分区数和 MongoDB 容量调整示例值:
tasks.max 只是允许的最大 Task 数,实际并行度受 Topic 分区分配限制;多个 Task 可能并发写同一集合。max.batch.size 按记录数拆分批次而不是按 BSON 字节数;增大该值通常能减少请求次数,但也可能增加单批延迟和错误影响范围。不同 Task、批次和命名空间之间不提供全局顺序或原子性。

将可容忍错误发送到死信队列

适用业务场景:链路进入稳定运维阶段,个别格式错误或写入失败的记录不应立即停止整个 Task,同时运维人员需要保留错误记录和上下文以便排查、修复和重放。 在快速开始配置基础上使用以下完整 Connector 配置,并提前创建独立的 DLQ Topic:
DLQ Topic 不能同时作为输入 Topic。被容忍的记录会从 MongoDB 写入中排除,并允许当前处理继续;这不是自动修复或重试机制。持续监控 DLQ,确认其中不含敏感上下文,并建立修复后重放流程;Write Concern 错误的写入结果可能不确定,不能仅凭 DLQ 记录推断目标端一定未写入。

监控

监控内容

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

导入 Grafana 大盘

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

限制条件

  • MongoDB 写入与 Kafka Offset 提交之间没有原子事务;Worker 在写入成功但 Offset 持久化前停止时,恢复后可能再次处理相同记录,因此不能据此声明 Exactly-Once、无重复或无丢失。
  • 写入不跨批次、命名空间或 Task 保持原子性或全局顺序;提高 tasks.max 可能使多个 Task 并发写入同一集合。
  • 启用 change.data.capture.handler 后会绕过普通后处理链、文档 ID 策略、普通 WriteModel 策略和墓碑删除配置,不能混用这些路径来组合行为。
  • InsertOne 没有 Upsert 行为;稳定 _id 的重放可能产生重复键错误,随机 _id 的重放可能创建额外文档。
  • 被容忍的错误记录不会由 Connector 自动重试;没有可用的 Errant Record Reporter 或未配置 DLQ 时,这些记录可能只被跳过。
  • 时序模式和 Stable API 要求 MongoDB 5.0 或更高版本,且时序集合的时间字段必须是 BSON 日期时间或可成功转换的值。

常见问题

Connector 启动时提示必须配置 Topic

topicstopics.regex 都为空或同时非空会导致校验失败。仅保留一种选择方式:明确列出 Topic 时使用 topics,需要动态匹配时使用有效的 Java 正则表达式配置 topics.regex;同时确认 DLQ Topic 不在输入范围内。

记录写入了错误的数据库或集合

默认映射器使用 database,并在 collection 为空时使用 Kafka Topic 名称。检查是否配置了 topic.override.<topic>.<property>FieldPathNamespaceMapper;字段路由值必须是字符串,且键和值不能同时为同一命名空间组件提供字段。需要在字段无效时直接失败,可设置 namespace.mapper.error.if.invalid=true

重启后 MongoDB 中出现额外文档

默认 BsonOidStrategy 每次处理都会生成新的 ObjectId,Offset 重放时可能产生新文档。根据业务身份选择稳定 ID 策略,例如按 Kafka Topic、Partition 和 Offset 生成 ID,或从稳定的消息键和值中提取 ID,并结合替换或更新 WriteModel 评估重放结果;稳定 ID 只能控制目标文档定位,不代表 Exactly-Once。

值为 null 的记录没有删除文档

默认 delete.on.null.values=false。启用后,默认删除策略还要求使用 FullKeyStrategyPartialKeyStrategyProvidedInKeyStrategy,并且消息键必须能形成有效删除条件。配置 CDC Handler 时,此设置会被忽略,应按对应 CDC Handler 的删除事件格式检查消息。

配置 errors.tolerance=data 时校验失败

MongoDB Connector 的 ConfigDef 接受 data,但 Kafka Connect 3.9.1 的框架配置仅接受 noneall。需要只容忍 Connector 内部的数据处理错误时,将 mongo.errors.tolerance=data 配置在 Connector 专用项中,并根据是否需要框架继续处理或 DLQ 决定将 errors.tolerance 保持为 none 或设置为 allmongo.errors.tolerance 不控制 Key/Value Converter 阶段的错误;此类错误仍由 Kafka Connect 框架的 errors.tolerance 处理。

已配置容错但 DLQ 中没有错误记录

检查 errors.deadletterqueue.topic.name 是否非空、DLQ Topic 是否与输入 Topic 分离、errors.tolerance 是否允许框架继续处理,以及 Connector 内部容错是否由 mongo.errors.tolerance 覆盖。还应确认运行环境提供可用的 Errant Record Reporter;否则 Connector 可能使用空 Reporter,错误记录不会写入 DLQ。