> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# MongoDB Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 MongoDB Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
topics=mongodb-sink-input
connection.uri=mongodb://<mongodb-host>:27017
database=inventory
collection=events
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
```

将 `<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 配置：

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
topics=mongodb-sink-input
connection.uri=mongodb://<mongodb-host>:27017
database=inventory
collection=events
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
document.id.strategy=com.mongodb.kafka.connect.sink.processor.id.strategy.KafkaMetaDataStrategy
```

`KafkaMetaDataStrategy` 根据 Topic、Partition 和 Offset 生成稳定 `_id`，配合默认替换并 Upsert 的写入模型，可使同一记录的重放定位到同一文档。它不提供端到端 Exactly-Once：MongoDB 写入与 Kafka Offset 提交之间没有原子事务，触发器、校验和其他外部副作用仍可能再次执行。

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

适用业务场景：基础链路已运行，需要利用多个 Kafka 分区提高并行度，并限制单次 MongoDB Bulk Write 的记录数，以平衡吞吐、请求大小、延迟和失败影响范围。

在快速开始配置基础上使用以下完整 Connector 配置，并按实际 Topic 分区数和 MongoDB 容量调整示例值：

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
topics=mongodb-sink-input
connection.uri=mongodb://<mongodb-host>:27017
database=inventory
collection=events
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
tasks.max=4
max.batch.size=1000
```

`tasks.max` 只是允许的最大 Task 数，实际并行度受 Topic 分区分配限制；多个 Task 可能并发写同一集合。`max.batch.size` 按记录数拆分批次而不是按 BSON 字节数；增大该值通常能减少请求次数，但也可能增加单批延迟和错误影响范围。不同 Task、批次和命名空间之间不提供全局顺序或原子性。

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

适用业务场景：链路进入稳定运维阶段，个别格式错误或写入失败的记录不应立即停止整个 Task，同时运维人员需要保留错误记录和上下文以便排查、修复和重放。

在快速开始配置基础上使用以下完整 Connector 配置，并提前创建独立的 DLQ Topic：

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
topics=mongodb-sink-input
connection.uri=mongodb://<mongodb-host>:27017
database=inventory
collection=events
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
errors.tolerance=all
mongo.errors.tolerance=all
errors.log.enable=true
errors.deadletterqueue.topic.name=mongodb-sink-dlq
errors.deadletterqueue.context.headers.enable=true
```

DLQ Topic 不能同时作为输入 Topic。被容忍的记录会从 MongoDB 写入中排除，并允许当前处理继续；这不是自动修复或重试机制。持续监控 DLQ，确认其中不含敏感上下文，并建立修复后重放流程；Write Concern 错误的写入结果可能不确定，不能仅凭 DLQ 记录推断目标端一定未写入。

## 监控

### 监控内容

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

### 导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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。
