概述
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.events;value.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.namespace和csfle.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需要mongocryptd或crypt_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 - 重要级别:中
- 有效值 / 注意事项:用于
FieldPathNamespaceMapper;false回退到默认数据库或集合,true抛出DataException。
topic.override.<topic>.<property>
为一个 Kafka Topic 覆盖 Topic 级 Sink 配置。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:必须使用实际 Topic 和已注册属性组成具体键,不能直接使用模板键;覆盖值会合并到全局配置之上,不能覆盖
connection.uri和topics。使用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 _id 的 IdStrategy 实现类。
- 类型:
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 - 重要级别:高
- 有效值 / 注意事项:不区分大小写,可选
string或binary;仅与UuidStrategy配合使用。
document.id.strategy.partial.key.projection.type
指定 PartialKeyStrategy 的专用字段投影模式。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:空值回退到
key.projection.type;可用none、allowlist、blocklist,仍兼容已弃用的whitelist和blacklist,应分别改用allowlist和blocklist。
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;可用none、allowlist、blocklist,仍兼容已弃用的whitelist和blacklist,应分别改用allowlist和blocklist。
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必须为FullKeyStrategy、PartialKeyStrategy或ProvidedInKeyStrategy;配置 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 - 重要级别:中
- 有效值 / 注意事项:至少为
0;0表示不按此配置拆分同一 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 - 重要级别:低
- 有效值 / 注意事项:可用
none、allowlist、blocklist;whitelist和blacklist已弃用,应分别改用allowlist和blocklist。配置 CDC 处理器时忽略。
key.projection.list
指定键字段投影使用的逗号分隔字段路径。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:按
key.projection.type解释;模式为none或配置 CDC 处理器时忽略。
value.projection.type
指定值字段投影模式。
- 类型:
string - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:可用
none、allowlist、blocklist;whitelist和blacklist已弃用,应分别改用allowlist和blocklist。配置 CDC 处理器时忽略。
value.projection.list
指定值字段投影使用的逗号分隔字段路径。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:按
value.projection.type解释;模式为none或配置 CDC 处理器时忽略。
field.renamer.mapping
指定精确字段路径重命名规则的 JSON 数组。
- 类型:
string - 默认值:
[] - 重要级别:低
- 有效值 / 注意事项:每个启用的对象必须包含
oldName和newName,并在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.timefield,0表示创建集合时不设置expireAfterSeconds。
timeseries.granularity
指定新建时序集合的预期时间间隔粒度。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:可选
seconds、minutes或hours;非空值要求配置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 仅接受
none或all;不要在此处使用 Connector ConfigDef 接受但框架会拒绝的data。启用通用错误继续处理和 DLQ 时使用all。
mongo.errors.tolerance
独立覆盖 Connector 内部容错模式。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可选
none、data或all;只需容忍 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.retries或retries.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:监控
监控内容
关注 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
topics 和 topics.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。启用后,默认删除策略还要求使用 FullKeyStrategy、PartialKeyStrategy 或 ProvidedInKeyStrategy,并且消息键必须能形成有效删除条件。配置 CDC Handler 时,此设置会被忽略,应按对应 CDC Handler 的删除事件格式检查消息。
配置 errors.tolerance=data 时校验失败
MongoDB Connector 的 ConfigDef 接受data,但 Kafka Connect 3.9.1 的框架配置仅接受 none 或 all。需要只容忍 Connector 内部的数据处理错误时,将 mongo.errors.tolerance=data 配置在 Connector 专用项中,并根据是否需要框架继续处理或 DLQ 决定将 errors.tolerance 保持为 none 或设置为 all。mongo.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。