对于 AI 代理:可在 https://www.mongodb.com/zh-cn/docs/llms.txt 获取文档索引—通过在任何 URL 路径后添加 .md 可获取所有页面的 Markdown 版本。
Docs 菜单

所有 Sink 连接器配置属性

在此页面上,您可以查看 MongoDB Kafka connector 的所有可用属性。此页面重复其他接收器连接器配置属性页面的内容。

要查看所有connector配置属性页面的列表,请参阅connector配置属性页面。

使用以下配置设置指定MongoDB Kafka Sink connector如何与MongoDB cluster连接并通信。

要仅查看与配置 MongoDB 连接相关的选项,请参阅MongoDB 连接配置属性页面。

名称
说明

connection.uri

必需

类型: string

说明:
用于连接 MongoDB 实例或集群的 MongoDB 连接 URI string
有关更多信息,请参阅 连接 MongoDB 指南

重要提示:为避免在 connection.uri 设置中暴露身份验证凭证,请使用ConfigProvider并设立适当的配置参数。

默认mongodb://localhost:27017
接受值:MongoDB 连接 URI string

server.api.version

类型:string

说明:
要与 MongoDB 服务器一起使用的 Stable API 版本。有关 Stable API 和支持它的服务器版本的更多信息,请参阅 Stable API MongoDB 服务器手册指南。

默认""
接受的值:空 string 或有效的 Stable API 版本。

server.api.deprecationErrors

类型:布尔值

说明:
设置为 true 时,如果连接器在 MongoDB 实例上调用已声明 Stable API 版本中已弃用的命令,则会引发异常。

您可以使用 server.api.version 配置选项设置 API 版本。有关 Stable API 的更多信息,请参阅 MongoDB 手册中有关 Stable API 的条目。

默认值false
可接受值truefalse

server.api.strict

类型:布尔值

说明:
设置为 true 时,如果连接器针对 MongoDB 实例调用已声明 Stable API 版本中未包含的命令,则会引发异常。

您可以使用 server.api.version 配置选项设置 API 版本。有关 Stable API 的更多信息,请参阅 MongoDB 手册中有关 Stable API 的条目。

默认值false
可接受值truefalse

使用以下配置设置指定 MongoDB Kafka connector 将数据写入哪个 MongoDB database 和 collection。您可以使用默认的DefaultNamespaceMapper或指定自定义类。

要仅查看与指定connector写入数据的位置相关的选项,请参阅MongoDB 命名空间映射配置属性页面。

名称
说明

namespace.mapper

类型:string

说明
指定将数据汇总到其中的数据库或集合的类的完全合格类名。默认 DefaultNamespaceMapper 使用 databasecollection 属性中指定的值。


连接器包含一个用于指定数据库和集合的替代类,名为FieldPathNamespaceMapper
有关详细信息,请参阅 FieldPathNamespaceMapper 设置


默认:

com.mongodb.kafka.connect.sink.namespace.mapping.DefaultNamespaceMapper

接受值:实现NamespaceMapper接口的类的完全限定 Java 类名称。

database

必需

类型:string

说明
汇入连接器写入的 MongoDB 数据库名称。

已接受值:MongoDB 数据库名称

集合

类型:string

说明
汇入连接器写入的 MongoDB 集合的名称。如果汇入连接器跟随多个主题,则这是未另外指定的任何写入的默认集合。

默认:主题名称。

接受值:MongoDB collection名称

如果将接收器连接器配置为使用FieldPathNamespaceMapper ,则可以根据数据的字段值指定用于接收文档的数据库和集合。

要启用此映射行为,请将接收器connectornamespace.mapper配置属性设置为完全限定的类名称,如下所示:

namespace.mapper=com.mongodb.kafka.connect.sink.namespace.mapping.FieldPathNamespaceMapper

FieldPathNamespaceMapper要求您指定以下设置:

  • 将属性映射到数据库和集合的一个或两个

  • 到数据库的keyvalue映射之一

  • 到collection的keyvalue映射之一

您可以使用以下设置自定义FieldPathNamespaceMapper的行为:

名称
说明

namespace.mapper.key.database.field

类型:string

说明
指定要写入的数据库名称的关键文档字段名称。

namespace.mapper.key.collection.field

类型:string

说明
指定要写入的集合名称的关键文档字段名称。

namespace.mapper.value.database.field

类型:string

说明
指定要写入的数据库名称的值文档字段名称。

namespace.mapper.value.collection.field

类型:string

说明
指定要写入的集合名称的值文档字段名称。

namespace.mapper.error.if.invalid

类型:布尔值

说明
当文档缺少映射字段或其具有无效的 BSON 类型时是否抛出异常。

当设置为 true 时,连接器不会处理缺少映射字段或包含无效 BSON 类型的文档。连接器可能会根据相关错误处理配置设置停止或跳过处理。

当设置为 false 时,如果文档缺少映射字段或其具有无效的 BSON 类型,连接器将默认写入到指定的 databasecollection 设置。

默认false
可接受值truefalse

使用以下配置设置指定 MongoDB Kafka 接收器连接器应监视哪些 Kafka 主题的数据。

要仅查看与指定 Kafka 主题相关的选项,请参阅 Sink Connector 的 Kafka 主题属性 页面。

名称
说明

话题

必需

类型:list

说明:
汇总连接器监视的 Kafka 主题列表。

您可以定义topicstopics.regex设置,但不能同时定义两者。

接受值:以逗号分隔的有效 Kafka 主题列表

topic.regex

必需

类型:string

说明:
匹配汇总连接器监视的 Kafka 主题的正则表达式。

示例,以下正则表达式匹配主题名称,例如“action.landing.clicks” 和“活动.支持.点击”。 它与主题名称“action.landing.views”不匹配 和“Activity.Clicks”。

topics.regex=activity\\.\\w+\\.clicks$

您可以定义topicstopics.regex设置,但不能同时定义两者。

接受值:使用java.util.regex.Pattern的有效正则表达式模式。

使用此页面上的设置可配置 MongoDB Kafka connector 的消息处理行为,包括:

  • 消息批量大小

  • 速率限制

  • 并行任务数

要仅查看与变更数据捕获处理程序相关的选项,请参阅connector消息处理Properties页面。

名称
说明

max.batch.size

类型: int

说明:
批处理的汇总水槽记录的最大数量。

考虑包含以下记录的批处理:

[ 1, 2, 3, 4, 5 ]

当设置为 0 时,连接器会对整个批处理执行单个批量写入。

当设置为 1 时,连接器会对批处理中的每条记录执行一次批量写入,共执行五次批量写入,如以下示例所示:

[1], [2], [3], [4], [5]

默认值0
可接受值:整数

bulk.write.ordered

类型:布尔值

说明:
连接器是将一批记录写入作为有序或无序批量写入操作。当设置为 true(默认值)时,连接器会将一批记录写入作为有序批量写入操作。

要了解有关批量写入操作的更多信息,请参阅批量写入操作

默认true
可接受值truefalse

rate.limiting.every.n

类型: int

说明:
汇集连接器处理的记录批处理数量,以 trigger 限速超时。值为 0 表示没有限速。

默认0
接受值:整数

rate.limiting.timeout

类型:int

说明:
在达到限速阈值后,汇总连接器应在多长时间(以毫秒为单位)后恢复处理。

默认值0
可接受值:整数

tasks.max

类型:int

说明:
为此连接器创建的最大任务数。如果连接器无法处理您指定的并行处理级别,则它可能创建的任务少于指定的最大任务数。

重要:如果您指定的值大于 1,则连接器将启用任务的并行处理。如果您的主题有多个分区日志,这使连接器可以并行从主题中读取,则任务可能会无序处理消息。

默认值1
接受值:整数

使用以下配置设置指定 MongoDB Kafka connector 如何处理错误并配置死信队列(DLQ)。

要仅查看与处理错误相关的选项,请参阅connector错误处理属性页面。

名称
说明

mongo.errors.tolerance

类型: string

说明:
连接器出错时是否继续处理消息。允许连接器覆盖 errors.tolerance Kafka 集群设置。

当设置为 none 时,连接器会报告任何错误,并阻止对其余消息的进一步处理。

当设置为 all 时,连接器会忽略任何有问题的消息。

当设置为 data 时,连接器只能容忍数据错误,并对所有其他错误失败。

要了解有关错误处理策略的更多信息,请参阅 处理错误 页面。

此属性覆盖 errors.tolerance
Connect 框架的属性。

默认:errors.tolerance 设置继承值。
可接受值: "none""all"

mongo.errors.log.enable

类型:布尔值

说明:
连接器是否应将错误详情(包括失败的操作)写入日志文件。连接器使用 errors.tolerancemongo.errors.tolerance 设置将错误分类为“可容忍”或“不可容忍”。

当设置为 true 时,连接器会记录“可容忍”和“不可容忍”错误。
当设置为 false 时,连接器会记录“不可容忍”错误。

此属性会覆盖 Connect Framework 的 errors.log.enable
Connect 框架的属性。

默认值: false
可接受值truefalse

errors.log.include.messages

类型:布尔值

说明:
连接器是否应在记录错误时包含无效消息。无效消息包含记录密钥、值和头部等数据。

默认: false
接受值truefalse

errors.deadletterqueue.topic.name

类型: string

说明:
用作死信队列(DLQ)的主题名称。如果为空,连接器将不会将任何无效消息发送到死信队列(DLQ)。

要学习有关死信队列(DLQ)的更多信息,请参阅死信队列(DLQ)配置示例

默认值: ""
接受的值:有效的 Kafka 主题名称

errors.deadletterqueue.context.headers.enable

类型:布尔值

说明:
连接器在将消息写入死信队列(DLQ)时是否包含上下文标头。

要了解有关死信队列(DLQ)的更多信息,请参阅死信队列(DLQ)配置示例

要了解连接器通过上下文标头定义和报告的异常,请参阅批量写入异常

默认: false
可接受值truefalse

errors.deadletterqueue.topic.replication.factor

类型:整数

说明:
复制死信队列(DLQ)主题的节点数。如果您运行单节点 Kafka 集群,则必须将其设置为 1

要了解有关死信队列(DLQ)的更多信息,请参阅 死信队列(DLQ)配置示例

默认值: 3
可接受值:有效的节点数

使用以下配置设置指定 MongoDB Kafka connector 在将 Kafka 数据插入 MongoDB 之前应如何转换该数据。

要仅查看与帖子处理器相关的选项,请参阅Sink connector 帖子处理器属性页面。

名称
说明

post.processor.chain

类型:list

说明:
连接器应在将数据保存到 MongoDB 之前应用的后处理器类列表。

要进一步了解后处理器并查看其用法示例,请参阅

Sink Connector 后处理器

默认:

com.mongodb.kafka.connect.sink.processor.DocumentIdAdder

接受的值:以逗号分隔的完全限定 Java 类名称列表

field.renamer.mapping

类型:string

说明:
键和值字段的字段名称映射列表。以以下格式在内联 JSON 数组中定义映射:

[ { "oldName":"key.fieldA", "newName":"field1" }, { "oldName":"value.xyz", "newName":"abc" } ]

默认[]
接受值:有效的 JSON 数组

field.renamer.regexp

类型:string

说明:
使用正则表达式对关键字段和值字段进行字段名称映射的列表。以以下格式在内联 JSON 数组中定义映射:

[ {"regexp":"^key\\\\..*my.*$", "pattern":"my", "replace":""}, {"regexp":"^value\\\\..*$", "pattern":"\\\\.", "replace":"_"} ]

默认[]
接受值:有效的 JSON 数组

key.projection.list

类型:string

说明:
连接器应包含在关键投影中的字段名称列表。

默认""
可接受值:字段名称的逗号分隔列表

key.projection.type

类型:string

说明:
连接器应使用的密钥投影类型。

默认 : 可接受值none
noneBlockListAllowList (已弃用:黑名单、白名单)

value.projection.list

类型:string

说明:
连接器应包含在值投影中的字段名称列表。

默认""
可接受值:字段名称的逗号分隔列表

value.projection.type

类型:string

说明:
连接器应使用的值投影类型。

默认 : 可接受值none
noneBlockListAllowList (已弃用:黑名单、白名单)

writemodel.strategy

类型: string

说明:
指定连接器应用于 批量写入WriteModelStrategy的类。

要了解如何创建自己的策略,请参阅
自定义写入模型策略

默认:

com.mongodb.kafka.connect.sink.writemodel.strategy.DefaultWriteModelStrategy

接受值:完全限定的 Java 类名称

使用以下配置设置,指定 MongoDB Kafka Sink 连接器应如何确定写入 MongoDB 的每个文档的 _id 值。

要仅查看与确定文档_id字段相关的选项,请参阅“connector ID 策略属性”页面。

名称
说明

document.id.strategy

类型: string

说明:
连接器应用于生成唯一 _id 字段的类。

默认:

com.mongodb.kafka.connect.sink.processor.id.strategy.BsonOidStrategy

接受的值:空字符串或完全限定的 Java 类名

document.id.strategy.overwrite.existing

类型:布尔值

说明:
connector 在应用 document.id.strategy 属性定义的策略时是否应覆盖 _id 字段中的现有值。

默认false
可接受值truefalse

document.id.strategy.uuid.format

类型:string

说明:
连接器应在 _id 字段中以 string 格式还是 BsonBinary 格式输出 UUID。

默认值string
接受值stringbinary

delete.on.null.values

类型:布尔值

说明:
当键值与 MongoDB 中的文档匹配且值字段为 null 时,连接器是否应删除文档。

当您指定对关键文档(如 FullKeyStrategyPartialKeyStrategyProvidedInKeyStrategy)进行操作的 ID 生成策略时,会应用此设置。

默认false
可接受值truefalse

您可以设置配置属性来指定MongoDB Kafka Sink connector如何将数据写入MongoDB 。 以下部分介绍了可以设置的配置属性,可自定义此行为。

writemodel.strategy 配置属性设置为 ,以指定接收器connector在接收接收器记录时如何写入数据。

您可以将writemodel.strategy的值设置为本页“策略”部分中描述的写入模型策略的任何完全限定类名称。 您可以通过设置以下配置来指定策略:

writemodel.strategy=<a write model strategy>

delete.writemodel.strategy 配置属性设置为 ,以指定 Sink connector在收到 tombstone 事件时如何写入数据。 逻辑删除事件是指包含键但不包含值的记录,这表示记录已删除。

您可以将delete.writemodel.strategy的值设置为本页“策略”部分中描述的写入模型策略的任何完全限定类名称。 您可以通过设置以下配置来指定策略:

delete.writemodel.strategy=<a write model strategy>

要仅查看与写入模型策略相关的选项,请参阅接收器连接器写入模型策略页面。

名称
说明

DefaultWriteModelStrategy

描述:
此策略默认使用 ReplaceOneDefaultStrategy,如果设置 timeseries.timefield 选项,则使用 InsertOneDefaultStrategy

这是 writemodel.strategy 配置属性的默认值。

InsertOneDefaultStrategy

描述:
将每个 sink 记录作为文档插入 MongoDB。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.InsertOneDefaultStrategy

ReplaceOneDefaultStrategy

描述:
根据 _id 字段,替换 MongoDB 中至多一个匹配接收记录的文档。如果没有任何文档匹配,连接器会将接收记录作为新文档插入。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.ReplaceOneDefaultStrategy

替换一个业务键策略

描述:
按指定的业务密钥替换与汇总记录匹配的最多一个文档。如果没有文档匹配,连接器会将汇总记录作为新文档插入。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.ReplaceOneBusinessKeyStrategy

要查看如何使用此策略的示例,请参阅我们的 写模型策略 指南

DeleteOneDefaultStrategy

描述:
最多删除一个与您的汇总连接器的关键结构匹配的文档,按 _id 字段,仅当文档包含空值结构时。

这是配置属性 delete.writemodel.strategy 的默认值。

当您设置 delete.on.null.values=true 时,此策略将设置为 writemodel.strategy 属性的默认值。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.DeleteOneDefaultStrategy

删除一个业务键策略

描述:
最多删除一个与业务密钥匹配的汇总记录的 MongoDB 文档。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.DeleteOneBusinessKeyStrategy

要查看如何使用此策略的示例,请参阅我们的 写模型策略 指南

UpdateOneDefaultStrategy

描述:
最多更新 MongoDB 中一个与汇总记录匹配的文档,按 _id 字段。如果没有文档匹配,连接器会将汇总记录作为新文档插入。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneDefaultStrategy

更新一个时间戳策略

描述:
在文档中添加 _insertedTS (插入时间戳) 和 _modifiedTS (修改时间戳) 字段。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneTimestampsStrategy

要查看如何使用此策略的示例,请参阅我们的 写模型策略 指南

UpdateOneBusinessKeyTimestampStrategy

描述:
在符合业务键的文档中添加 _insertedTS(插入时间戳)和 _modifiedTS(修改时间戳)字段。
要指定此策略,请将配置属性设置为以下类名:

com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneBusinessKeyTimestampStrategy

使用以下 MongoDB Kafka 接收器连接器配置设置来覆盖特定主题的全局或默认属性设置。

要仅查看与覆盖主题设置相关的选项,请参阅主题覆盖属性页面。

名称
说明

topic.override.<topicName>.<propertyName>

类型:string

说明:
指定主题和属性名称以覆盖对应的全局或默认属性设置。

例如,topic.override.foo.collection=bar
设置指示汇集连接器将foo 主题中的数据存储在bar
集合中。

您可以在
<propertyName> 段中按主题指定任何有效的配置设置,但
connection.uritopics 除外。

默认""
可接受值:覆盖属性特定的可接受值

使用以下配置设置指定MongoDB Kafka接收器connector用于处理变更数据捕获 (CDC) 事件的类。

有关使用内置ChangeStreamHandler和Debezium及Qlik Replicate事件生成器的处理程序的示例,请参阅变更数据捕获处理程序指南

要仅查看与变更数据捕获处理程序相关的选项,请参阅变更数据捕获属性页面。

名称
说明

change.data.capture.handler

类型:string

说明:
用于将更改转换为事件流的 CDC 处理程序的类名称。有关 CDC 处理程序列表,请参阅可用 CDC 处理程序

默认""
接受的值:空 string 或完全限定的 Java 类名

使用以下配置设置指定 MongoDB Kafka 接收器连接器应如何将数据接收到 MongoDB 时间序列集合。

要仅查看与time-series collection相关的选项,请参阅Kafka time-series 属性页面。

名称
说明

timeseries.timefield

类型:string

说明:
源数据中的顶级字段的名称,其中包含要与时间序列集合中的新文档关联的时间信息。

默认""
可接受值:空字符串或包含 BSON DateTime 值的字段的名称

timeseries.timefield.auto.convert.date.format

类型:string

说明:
连接器应用于转换 timeseries.timefield 设置指定的字段中包含的源数据的日期格式模式。

连接器将日期格式模式传递给 Java DateTimeFormatter.ofPattern(pattern, locale) 方法,以对时间字段进行日期和时间转换。

如果源数据中的日期值只包含日期信息,连接器会将时间信息设置为指定日的开始。如果日期值不包含时区偏移,连接器会将偏移设置为 UTC。

默认

yyyy-MM-dd[['T'][ ]][HH:mm:ss[[.][SSSSSS][SSS]][ ]VV[ ]'['VV']'][HH:mm:ss[[.][SSSSSS][SSS]][ ]X][HH:mm:ss[[.][SSSSSS][SSS]]]

接受的值:有效的DateTimeFormatter格式

timeseries.timefield.auto.convert

类型:布尔值

说明:
是否将字段中的数据转换为 BSON Date 格式。

当设置为 true 时,连接器使用 Unix 纪元后的毫秒,并在值为数字时忽略小数部分。如果值为 string,连接器使用以下配置中的设置来解析日期:

timeseries.timefield.auto.convert.date.format

如果连接器无法转换该值,它会将原始值发送到时间序列集合。

默认false
可接受值truefalse

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

类型:string

说明:
要与日期格式模式一起使用的 DateTimeFormatter 区域设置语言标签(例如:"en-US")。

要了解有关区域设置的更多信息,请参阅 Java SE 文档中的 Locale

默认值ROOT
可接受值:有效的 Locale 语言标签格式

timeseries.metafield

类型: string

说明:
要从源数据中读取哪个顶级字段来描述一组相关的时间序列文档。

重要:此字段不得为 _id 字段,也不得为您在 timeseries.timefield 设置中指定的字段。

默认值""
可接受值:空 string 或包含除 BsonArray 以外的任何 BSON 类型的字段名称。

timeseries.expire.after.seconds

类型:int

说明:
MongoDB 在自动删除时间序列集合数据之前应等待的秒数。当设置值小于 1 时,连接器会禁用定时到期。

要了解更多信息,请参阅 MongoDB 手册中的设置时间序列集合的自动删除

默认0
接受值:整数

timeseries.granularity

类型:string

说明:
源数据后续测量之间的预期间隔。

要了解更多信息,请参阅 MongoDB 手册中的设置时间序列数据的粒度

可选
默认""
可接受值"""seconds""minutes""hours"

有关如何将现有集合转换为时间序列集合的示例,请参阅如何将现有集合迁移到时间序列集合的教程。