> ## 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 Source Connector

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

## 概述

MongoDB Source Connector 使用 MongoDB Change Streams 读取数据库、集合或整个部署中的变更，并将事件写入 Kafka Topic，适合构建数据同步、流式处理、搜索索引和审计等数据链路。默认情况下，MongoDB 命名空间 `database.collection` 映射到同名 Kafka Topic，消息值包含完整的变更流事件。

Connector 可以只采集启动后产生的新变更，也可以先复制已有文档，再衔接后续变更。复制的文档会封装为形似插入操作的 Change Stream 事件，并经过与实时事件相同的 Topic、键和值映射流程。

## 前置条件

* MongoDB 部署必须支持 Change Streams；连接账号需要对采集范围具备读取变更流所需权限，使用存量复制时还需要读取目标集合。
* 使用变更前文档、扩展事件或 Stable API 等 MongoDB 能力时，服务端版本和数据库设置必须满足对应功能要求。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和支持 Change Streams 的 MongoDB 部署，并确认网络连通和访问权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSourceConnector
connection.uri=mongodb://<username>:<password>@<mongodb-host>:27017/?replicaSet=<replica-set>
database=<database-name>
collection=<collection-name>
```

替换 MongoDB 连接信息、数据库名和集合名。该配置从当前变更流位置开始采集，默认将事件写入 `<database-name>.<collection-name>` Topic；建议通过 Config Provider 或受保护的配置管理连接凭据。

## 配置

### Connector 与任务

#### `connector.class`

指定 MongoDB Source Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：使用 `com.mongodb.kafka.connect.MongoSourceConnector`。
* **必填**：是

#### `tasks.max`

设置 Connect 允许创建的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须至少为 `1`；此 Connector 始终只创建一个 Source Task，设置为更大的值不会增加并行 Task。

### MongoDB 连接与采集范围

#### `connection.uri`

MongoDB 连接字符串，用于指定主机、认证、Replica Set、TLS 和其他驱动选项。

* **类型**：`password`
* **默认值**：`mongodb://localhost:27017,localhost:27018,localhost:27019`
* **重要级别**：高
* **有效值 / 注意事项**：必须是有效的 MongoDB URI。默认值通常只适合本地开发；不要在明文配置中保存凭据。

#### `database`

限定要监听的数据库。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：留空时监听整个 MongoDB 部署；设置数据库但不设置集合时监听该数据库中的集合。

#### `collection`

限定要监听的单个集合。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：不是正则表达式或列表；设置时必须同时设置 `database`。

#### `offset.partition.name`

覆盖 Connector 用于标识 Source Offset 分区的名称。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：留空时根据连接主机、数据库和集合生成。修改此值会开始使用新的 Offset 身份，原有 Resume Token 不会自动迁移。

### MongoDB Stable API

#### `server.api.version`

启用指定版本的 MongoDB Stable API。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：留空表示不启用；可用值为 `1`，要求 MongoDB 5.0 或更高版本。

#### `server.api.deprecation.errors`

启用 Stable API 弃用错误。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `server.api.version` 非空时生效。

#### `server.api.strict`

启用 Stable API 严格模式。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `server.api.version` 非空时生效。

### TLS 证书

#### `connection.ssl.truststore`

指定 Connect Worker 本地的 Truststore 路径。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空时使用 JVM 默认 KeyStore 类型加载；TLS 本身仍需由连接 URI 或部署设置启用。

#### `connection.ssl.truststorePassword`

指定 Truststore 密码。

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

#### `connection.ssl.keystore`

指定 Connect Worker 本地的 Keystore 路径，用于提供客户端证书和私钥。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空时使用 JVM 默认 KeyStore 类型加载。

#### `connection.ssl.keystorePassword`

指定 Keystore 及其中私钥的密码。

* **类型**：`password`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `connection.ssl.keystore` 非空时使用；同一密码用于加载 KeyStore 和初始化私钥。

### Change Stream

#### `pipeline`

为 Change Stream 设置聚合管道。

* **类型**：`string`
* **默认值**：`[]`
* **重要级别**：中
* **有效值 / 注意事项**：必须是 JSON 文档数组。管道需保留 Topic、Offset、键和值映射所需的 `_id`、`ns`、`documentKey` 或 `fullDocument` 字段；存量复制的高效筛选优先使用 `startup.mode.copy.existing.pipeline`。

#### `batch.size`

设置 MongoDB Change Stream Cursor 的批次大小提示。

* **类型**：`int`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `0`；`0` 表示沿用驱动和服务端行为，与 `poll.max.batch.size` 不同。

#### `publish.full.document.only`

仅将变更事件的 `fullDocument` 作为消息值发布。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：高
* **有效值 / 注意事项**：设为 `true` 时强制使用 `updateLookup`；没有 `fullDocument` 的事件会被过滤，除非同时启用删除 Tombstone。

#### `publish.full.document.only.tombstone.on.delete`

在仅发布完整文档模式下，为缺少文档型 `fullDocument` 的事件发送空值 Tombstone；删除事件是最常见的情形。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `publish.full.document.only=true` 时生效。实现按是否存在文档型 `fullDocument` 判断，不只检查事件的 `operationType` 是否为 `delete`。

#### `change.stream.document.key.as.key`

控制非 Schema 输出是否使用 Change Stream 的 `documentKey` 作为 Kafka 消息键。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：设为 `false` 或事件没有 `documentKey` 时，键使用 Resume Token 文档；`output.format.key=schema` 时忽略此项。

#### `change.stream.full.document.before.change`

设置 Change Stream 返回变更前文档的方式。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：可用值为 `default`、`off`、`whenAvailable`、`required` 或空字符串。需要 MongoDB 6.0 或更高版本并为集合启用 Pre-image；不影响存量复制记录。

#### `change.stream.full.document`

设置 Change Stream 返回完整文档的方式。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：可用值为 `default`、`updateLookup`、`whenAvailable`、`required` 或空字符串；`publish.full.document.only=true` 时强制使用 `updateLookup`。

#### `change.stream.show.expanded.events`

控制是否请求扩展 Change Stream 事件。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：扩展 DDL 事件需要 MongoDB 6.0 或更高版本，消除歧义的更新路径需要 MongoDB 6.1 或更高版本。

#### `collation`

为 Change Stream 设置 Collation。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：留空表示不设置；否则必须是有效的 MongoDB Collation JSON 文档。该设置不应用于存量复制聚合。

#### `poll.max.batch.size`

限制单次 `poll` 返回的最大记录数。

* **类型**：`int`
* **默认值**：`1000`
* **重要级别**：低
* **有效值 / 注意事项**：必须至少为 `1`；只限制记录数量，不限制单批总字节数。

#### `poll.await.time.ms`

设置 Change Stream 等待新事件的最长时间。

* **类型**：`long`
* **默认值**：`5000`
* **重要级别**：低
* **有效值 / 注意事项**：必须至少为 `1`，单位为毫秒；该值不是 Kafka Connect 的 Poll Timeout。

### Topic 映射

#### `topic.mapper`

指定将 MongoDB 命名空间映射为 Kafka Topic 的实现类。

* **类型**：`string`
* **默认值**：`com.mongodb.kafka.connect.source.topic.mapping.DefaultTopicMapper`
* **重要级别**：高
* **有效值 / 注意事项**：必须是可加载且实现 `TopicMapper` 的完整 Java 类名。自定义实现可以读取自身定义的附加属性。

#### `topic.separator`

指定默认 Topic 名各部分之间的分隔符。

* **类型**：`string`
* **默认值**：`.`
* **重要级别**：低
* **有效值 / 注意事项**：用于连接前缀、数据库、集合和后缀；Connector 不校验最终 Topic 名是否合法。

#### `topic.prefix`

为默认映射产生的所有 Topic 添加前缀。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：低
* **有效值 / 注意事项**：非空前缀后会追加 `topic.separator`。

#### `topic.suffix`

为默认映射产生的所有 Topic 添加后缀。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：低
* **有效值 / 注意事项**：非空后缀前会添加 `topic.separator`。

#### `topic.namespace.map`

使用 JSON 对象将 MongoDB 命名空间映射到指定 Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：支持完整命名空间、数据库、按声明顺序匹配的 `/正则/` 和 `*` 映射；映射结果仍会应用前缀和后缀。模板可使用 `db`、`sep`、`coll`、`sep_coll`、`coll_sep` 和 `sep_coll_sep` 变量；需自行避免非法名称和 Topic 冲突。

### 输出格式与 Schema

#### `output.format.key`

设置消息键的输出格式。

* **类型**：`string`
* **默认值**：`json`
* **重要级别**：高
* **有效值 / 注意事项**：可用值为 `json`、`bson`、`schema`，不区分大小写；应与 `key.converter` 兼容。

#### `output.format.value`

设置消息值的输出格式。

* **类型**：`string`
* **默认值**：`json`
* **重要级别**：高
* **有效值 / 注意事项**：可用值为 `json`、`bson`、`schema`，不区分大小写；应与 `value.converter` 兼容。

#### `output.schema.infer.value`

为每条消息值自动推断 Connect Schema。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `output.format.value=schema` 时生效。文档结构变化可能使同一 Topic 出现多个 Schema。

#### `output.schema.key`

定义 Schema 格式消息键使用的 Avro Schema。

* **类型**：`string`
* **默认值**：`{  "type": "record",  "name": "keySchema",  "fields" : [{"name": "_id", "type": "string"}]}`
* **重要级别**：高
* **有效值 / 注意事项**：必须是有效且与完整 Change Stream 文档转换兼容的 Connector Avro Schema；仅在 `output.format.key=schema` 时读取。

#### `output.schema.value`

定义 Schema 格式消息值使用的 Avro Schema。

* **类型**：`string`
* **默认值**：内置 `ChangeStream` Avro Schema
* **重要级别**：高
* **有效值 / 注意事项**：必须是有效且与消息值兼容的 Connector Avro Schema；`output.schema.infer.value=true` 或输出格式不是 `schema` 时忽略。

#### `output.json.formatter`

指定 BSON 到 JSON 的格式化实现类。

* **类型**：`string`
* **默认值**：`com.mongodb.kafka.connect.source.json.formatter.DefaultJson`
* **重要级别**：中
* **有效值 / 注意事项**：必须是可加载且实现 `JsonWriterSettingsProvider` 的完整 Java 类名；JSON 输出及 Schema 转换会使用，原始 BSON 输出会忽略。

#### `key.converter`

覆盖 Worker 级消息键 Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：`null` 表示继承 Worker 配置；所选 Converter 必须兼容 `output.format.key` 产生的 String、BSON Bytes 或 Schemaful Connect Data。

#### `value.converter`

覆盖 Worker 级消息值 Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：`null` 表示继承 Worker 配置；所选 Converter 必须兼容 `output.format.value` 产生的 String、BSON Bytes 或 Schemaful Connect Data。

### 启动位置与存量复制

#### `startup.mode`

设置没有可用 Source Offset 时的启动方式。

* **类型**：`string`
* **默认值**：空字符串，实际按 `latest` 处理
* **重要级别**：中
* **有效值 / 注意事项**：可用值为 `latest`、`timestamp`、`copy_existing` 或空字符串。已有 Resume Token 优先于该设置；显式非空值优先于弃用的 `copy.existing`。

#### `startup.mode.timestamp.start.at.operation.time`

设置 Timestamp 启动模式的起始操作时间。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `startup.mode=timestamp` 且没有可用 Offset 时应用；支持十进制 Epoch 秒、秒精度 ISO-8601 Instant 或 Canonical Extended JSON BSON Timestamp。

#### `startup.mode.copy.existing.max.threads`

设置存量复制读取线程数。

* **类型**：`int`
* **默认值**：`Runtime.getRuntime().availableProcessors()`（运行环境可用处理器数量）
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `1`，仅在 `startup.mode=copy_existing` 时使用；实际线程数不会超过选中的命名空间数量。显式设置时覆盖弃用项 `copy.existing.max.threads`。

#### `startup.mode.copy.existing.queue.size`

设置存量复制线程与 Source Task 之间的内存队列容量。

* **类型**：`int`
* **默认值**：`16000`
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `1`，仅在 `startup.mode=copy_existing` 时使用；显式设置时覆盖弃用项 `copy.existing.queue.size`。

#### `startup.mode.copy.existing.pipeline`

设置只用于读取存量文档的聚合管道。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：必须为空或 JSON 文档数组，仅在 `startup.mode=copy_existing` 时使用；它在合成变更事件和通用 `pipeline` 之前执行。显式设置时覆盖弃用项 `copy.existing.pipeline`。

#### `startup.mode.copy.existing.namespace.regex`

使用正则表达式筛选启动时发现的存量复制命名空间。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：匹配完整的 `database.collection`，使用 Java 正则的子串匹配语义，仅在 `startup.mode=copy_existing` 时使用。显式设置时覆盖弃用项 `copy.existing.namespace.regex`。

#### `startup.mode.copy.existing.allow.disk.use`

控制存量复制聚合是否允许 MongoDB 使用磁盘。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `startup.mode=copy_existing` 时使用；显式设置时覆盖弃用项 `copy.existing.allow.disk.use`。

#### `copy.existing`

使用旧式开关启用存量复制。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `startup.mode` 为空时生效。
* **已弃用**：是
* **替代项**：`startup.mode=copy_existing`

#### `copy.existing.max.threads`

设置旧式存量复制线程数。

* **类型**：`int`
* **默认值**：`Runtime.getRuntime().availableProcessors()`（运行环境可用处理器数量）
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `1`；只有未显式设置替代项时才作为回退值。
* **已弃用**：是
* **替代项**：`startup.mode.copy.existing.max.threads`

#### `copy.existing.queue.size`

设置旧式存量复制内存队列容量。

* **类型**：`int`
* **默认值**：`16000`
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `1`；只有未显式设置替代项时才作为回退值。
* **已弃用**：是
* **替代项**：`startup.mode.copy.existing.queue.size`

#### `copy.existing.pipeline`

设置旧式存量复制聚合管道。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：必须为空或 JSON 文档数组；只有未显式设置替代项时才作为回退值。
* **已弃用**：是
* **替代项**：`startup.mode.copy.existing.pipeline`

#### `copy.existing.namespace.regex`

设置旧式存量复制命名空间正则表达式。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：必须是有效 Java 正则；只有未显式设置替代项时才作为回退值。
* **已弃用**：是
* **替代项**：`startup.mode.copy.existing.namespace.regex`

#### `copy.existing.allow.disk.use`

设置旧式存量复制聚合是否允许使用磁盘。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：只有未显式设置替代项时才作为回退值。
* **已弃用**：是
* **替代项**：`startup.mode.copy.existing.allow.disk.use`

### 错误处理与心跳

#### `errors.tolerance`

设置 Kafka Connect 框架及 MongoDB Connector 的错误容忍模式。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：可用值为 `none`、`all`，不区分大小写。`all` 可跳过转换失败并允许部分 Change Stream 恢复分支无 Resume Token 重建；`mongo.errors.tolerance` 仅覆盖 Connector 自身对该值的读取。

#### `mongo.errors.tolerance`

仅覆盖 MongoDB Connector 内部的错误容忍模式。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：可用值为 `none`、`all`；只要显式设置，就优先于 Connector 对 `errors.tolerance` 的读取，但不会替代 Kafka Connect 框架的同名设置。

#### `errors.log.enable`

控制 Kafka Connect 框架及 MongoDB Connector 是否记录被容忍的错误。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：容忍模式为 `none` 时 Connector 仍会记录转换错误；`mongo.errors.log.enable` 可仅覆盖 Connector 自身行为。

#### `mongo.errors.log.enable`

仅覆盖 MongoDB Connector 是否记录被容忍的转换错误。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：显式设置时优先于 Connector 对 `errors.log.enable` 的读取，不改变 Kafka Connect 框架日志设置。

#### `errors.deadletterqueue.topic.name`

设置 MongoDB 转换失败记录写入的 DLQ Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 Connector 的有效容忍模式为 `all` 时使用；Connector 不创建或校验 Topic。`mongo.errors.deadletterqueue.topic.name` 可覆盖 Connector 对该值的读取。

#### `mongo.errors.deadletterqueue.topic.name`

仅覆盖 MongoDB Connector 转换失败记录使用的 DLQ Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：要求 Connector 的有效错误容忍模式为 `all`；仅影响 Connector 转换错误处理。

#### `heartbeat.interval.ms`

设置无业务记录时生成心跳记录的最小间隔。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：必须至少为 `0`，单位为毫秒；`0` 表示禁用。只有出现新的 Post-batch Resume Token 且当前 Poll 没有业务记录时才会生成心跳。

#### `heartbeat.topic.name`

指定心跳记录写入的 Kafka Topic。

* **类型**：`string`
* **默认值**：`__mongodb_heartbeats`
* **重要级别**：中
* **有效值 / 注意事项**：必须为非空字符串，仅在 `heartbeat.interval.ms` 大于 `0` 且生成心跳时使用；Topic 需预先满足命名、权限和可用性要求。

### 自定义凭据

#### `mongo.custom.auth.mechanism.enable`

启用自定义 MongoDB 凭据提供器。

* **类型**：`string`
* **默认值**：无
* **重要级别**：未声明
* **有效值 / 注意事项**：仅字符串值 `true`（不区分大小写）会启用；该属性未注册到 ConfigDef，因此不会出现在 ConfigDef 元数据中，也不会接受 ConfigDef 类型校验。

#### `mongo.custom.auth.mechanism.providerClass`

指定自定义 MongoDB 凭据提供器实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：未声明
* **有效值 / 注意事项**：启用自定义认证时必填。类必须可加载、实现 `CustomCredentialProvider`、具有可访问的无参构造方法，并在所有可运行 Task 的 Worker 上提供实现及依赖。

## 最佳实践

### 首次接入已有集合并持续采集增量

适用业务场景：目标集合已经包含存量文档，需要先建立 Kafka 数据基线，再持续接收复制期间及之后产生的新变更。该模式适合首次接入，不适合作为每次重启都重复执行的初始化流程。

在快速开始配置基础上使用以下完整 Connector 配置：

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSourceConnector
connection.uri=mongodb://<username>:<password>@<mongodb-host>:27017/?replicaSet=<replica-set>
database=<database-name>
collection=<collection-name>
startup.mode=copy_existing
```

关键说明：Connector 启动时复制所选范围内已存在的文档，并衔接 Change Stream。复制中断后重启会重新执行整轮复制，消费者应使用稳定业务键或其他方式处理重复；复制不是跨多个集合的一致性快照。

### 为业务命名空间设置稳定的 Topic 名称

适用业务场景：MongoDB 的数据库或集合名称不适合作为下游 Topic 契约，或者迁移过程中需要保持现有 Topic 名称不变。显式映射可将源端命名与 Kafka 消费接口解耦。

使用以下 Connector 配置将一个命名空间路由到固定 Topic：

```properties theme={null}
connector.class=com.mongodb.kafka.connect.MongoSourceConnector
connection.uri=mongodb://<username>:<password>@<mongodb-host>:27017/?replicaSet=<replica-set>
database=inventory
collection=orders
topic.namespace.map={"inventory.orders":"commerce.orders"}
```

关键说明：映射在前缀和后缀处理之前完成。调整映射不会改变 Source Offset 身份，但必须提前确认目标 Topic 的命名、权限、分区策略以及是否会与其他映射结果冲突。

## 监控

### 监控内容

关注 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 并选择对应数据源。

## 限制条件

* 每个 Connector 实例只创建一个 Source Task；增大 `tasks.max` 不能水平扩展 Change Stream 消费。
* `copy_existing` 中途重启会从头重新复制，已写入 Kafka 的存量记录可能重复；该过程也不是跨集合的一致性快照。
* 多命名空间复制可能交错，Topic 映射可以拆分数据流，且 Connector 不指定 Kafka Partition，因此不保证跨命名空间、Topic 或 Partition 的全局顺序。
* `publish.full.document.only=true` 会过滤没有 `fullDocument` 的事件，除非同时启用 Tombstone 选项生成空值记录。
* 错误容忍设为 `all` 且未配置 DLQ 时，转换失败记录会被丢弃，后续 Offset 可能永久越过该事件。
* Resume Token 无效或历史已丢失时，容忍分支可能无 Token 重建 Change Stream，无法补回已经离开 Oplog 的数据区间。
* Connector 不提供 MongoDB 到 Kafka 的端到端 Exactly-once、无重复或无丢失保证；Offset 提交窗口、存量复制和错误容忍均可能产生重复或数据缺口。

## 常见问题

### Connector 运行后没有收到新消息

检查 MongoDB 部署是否支持 Change Streams、连接账号是否具有所选范围的权限，以及 `database` 和 `collection` 是否指向实际写入的位置。确认默认 Topic `<database>.<collection>` 或 `topic.namespace.map` 的目标 Topic 正确；若使用 `pipeline`，还要确认它没有过滤全部事件或移除 `_id`、`ns` 等必需字段。

### 更新事件中没有最新完整文档

默认 Change Stream 事件不一定包含更新后的完整文档。需要完整文档时设置 `change.stream.full.document=updateLookup`；如果只希望发布完整文档，可设置 `publish.full.document.only=true`，但需同时评估没有 `fullDocument` 的事件会被过滤，以及删除事件是否需要 Tombstone。

### 重启后从意外位置开始采集

Connector 优先使用已提交的 Resume Token，只有没有可用 Offset 时才应用 `startup.mode`。检查是否修改了连接主机、`database`、`collection` 或 `offset.partition.name`，这些变化可能产生新的 Offset 身份；恢复旧配置后确认原 Offset 仍然存在，再重启 Connector。

### 存量复制后出现重复记录

存量复制期间的文档会作为插入事件发布，复制中断重启会重新执行整轮复制，同时发生的更新还可能作为后续 Change Stream 事件出现。消费者应根据稳定的文档键实现幂等处理，并在重启前判断是否需要继续使用原 Offset 或重新建立数据基线。

### 转换错误没有进入 DLQ

确认有效的 Connector 错误容忍模式为 `all`，DLQ Topic 已存在且 Worker 有写权限，并检查是否设置了 `mongo.errors.*` 覆盖项。Topic 映射为空、仅发布完整文档时缺少 `fullDocument` 等过滤行为不属于转换失败，不会写入 DLQ。
