> ## 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.

# YugabyteDB Source Connector

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

## 概述

YugabyteDB Source Connector 通过 YSQL 逻辑复制读取 YugabyteDB 表的变更，并将 INSERT、UPDATE、DELETE 和 TRUNCATE 事件写入 Kafka Topic。普通表变更使用由 `topic.prefix`、Schema 和表名组成的 Topic 名称，消息值采用包含 `before`、`after`、`source`、`op` 和时间戳等字段的 Debezium 事件结构。

Connector 可以先为现有数据创建快照，再持续读取增量变更，也可以在已有数据基线时直接进入流式读取。它适合将 YugabyteDB 中的业务数据同步到流处理、搜索、分析或异构数据系统。Kafka 记录键通常来自表主键；没有可用主键时，记录键可能为空。

## 前置条件

* 使用具有 `LOGIN` 和 `REPLICATION` 属性的专用 YugabyteDB YSQL 账号。根据部署方式预先创建或允许 Connector 创建匹配的 replication slot 和 publication；如果允许 Connector 创建 publication，还需授予数据库 `CREATE` 权限以及目标表所需的所有权。确认 `plugin.name` 对应服务端支持的逻辑解码插件。目标 Connector 制品声明适用于 YugabyteDB 2024.1.x。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和 YugabyteDB YSQL 数据库，并确认网络连通和访问权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=io.debezium.connector.postgresql.YugabyteDBConnector
topic.prefix=ybdb
database.hostname=<database-host>
database.port=5433
database.user=<database-user>
database.password=<database-password>
database.dbname=<database-name>
plugin.name=yboutput
slot.name=debezium
publication.name=dbz_publication
snapshot.mode=initial
```

将数据库地址、账号、密码和数据库名占位符替换为实际环境值，并确认端口、逻辑解码插件、slot 和 publication 与目标集群一致。使用敏感配置管理方式提供数据库密码。应用配置后，Connector 会先建立初始快照，再在该快照之后持续发送变更事件。

## 配置

### 运行与序列化

#### `connector.class`

指定 Kafka Connect 加载的 YugabyteDB Source Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。
* **必填**：是

#### `tasks.max`

设置 Connector 可创建的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须为大于或等于 `1` 的整数。

#### `tasks.max.enforce`

控制 Kafka Connect 是否强制执行 `tasks.max` 上限。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：Kafka Connect 3.9.1 已将此项标记为弃用，但仍接受该配置；当前 ConfigDef 未声明替代项。
* **已弃用**：是

#### `key.converter`

覆盖 Worker 级配置，指定记录键的 Converter。

* **类型**：`class`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：必须是具有公共无参构造方法的 `Converter` 具体实现类。

#### `value.converter`

覆盖 Worker 级配置，指定记录值的 Converter。

* **类型**：`class`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：必须是具有公共无参构造方法的 `Converter` 具体实现类。

#### `header.converter`

覆盖 Worker 级配置，指定记录 Header 的 Converter。

* **类型**：`class`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：必须是具有公共无参构造方法的 `HeaderConverter` 具体实现类。

#### `transforms`

指定按顺序应用于记录的单消息转换别名。

* **类型**：`list`
* **默认值**：`[]`
* **重要级别**：低
* **有效值 / 注意事项**：填写互不重复且非空的转换别名列表；每个别名还需配置对应的 `transforms.<alias>.type`。

#### `predicates`

指定供单消息转换引用的 Predicate 别名。

* **类型**：`list`
* **默认值**：`[]`
* **重要级别**：低
* **有效值 / 注意事项**：填写互不重复且非空的 Predicate 别名列表；每个别名还需配置对应的 `predicates.<alias>.type`。

#### `config.action.reload`

控制外部配置提供程序中的值变化时是否重新加载 Connector。

* **类型**：`string`
* **默认值**：`restart`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`none`、`restart`。

### 错误处理与 Topic 创建

#### `errors.retry.timeout`

设置失败操作的总重试时长，单位毫秒。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：`0` 表示不重试，`-1` 表示无限重试，正整数表示总重试时长。

#### `errors.retry.delay.max.ms`

设置连续重试之间的最大等待时间，单位毫秒。

* **类型**：`long`
* **默认值**：`60000`
* **重要级别**：中
* **有效值 / 注意事项**：必须为非负毫秒值；达到上限后会在延迟中加入抖动。

#### `errors.tolerance`

控制遇到记录处理错误时是立即失败还是跳过错误记录。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`none`、`all`。

#### `errors.log.enable`

控制是否把可容忍错误及失败操作详情写入 Connect 日志。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `errors.log.include.messages`

控制错误日志是否包含导致失败的记录内容和元数据。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：启用后可能把记录键、值、Header 和源 Offset 等内容写入日志，应评估敏感数据暴露风险。

#### `topic.creation.groups`

指定 Source Connector 自动创建 Topic 时使用的配置组别名。

* **类型**：`list`
* **默认值**：`[]`
* **重要级别**：低
* **有效值 / 注意事项**：填写互不重复且非空的 Topic 创建组别名。

### Source 事务与 Offset

#### `exactly.once.support`

控制创建 Connector 前是否要求验证 Source Connector 的恰好一次支持。

* **类型**：`string`
* **默认值**：`requested`
* **重要级别**：中
* **有效值 / 注意事项**：可用值（不区分大小写）：`required`、`requested`。

#### `transaction.boundary`

指定 Source Task 提交 Producer 事务时采用的边界。

* **类型**：`string`
* **默认值**：`poll`
* **重要级别**：中
* **有效值 / 注意事项**：可用值（不区分大小写）：`interval`、`poll`、`connector`。

#### `transaction.boundary.interval.ms`

设置按时间间隔提交 Source Producer 事务的周期，单位毫秒。

* **类型**：`long`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：必须大于或等于 `0`；仅在 `transaction.boundary=interval` 时生效，未设置时使用 Worker 的 `offset.flush.interval.ms`。

#### `offsets.storage.topic`

为当前 Connector 指定独立的 Offset 存储 Topic。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：必须为非空 Topic 名称；仅 Distributed 模式生效，未设置时使用 Worker 的全局 Offset Topic。

### Topic 与数据库连接

#### `topic.prefix`

设置数据 Topic 和 Connector 系统 Topic 使用的唯一前缀。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。
* **必填**：是

#### `database.hostname`

设置 YugabyteDB YSQL 服务的主机名或 IP 地址。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。
* **必填**：是

#### `database.port`

设置 YugabyteDB YSQL 服务端口。

* **类型**：`int`
* **默认值**：`5433`
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `database.user`

设置连接 YugabyteDB YSQL 的数据库用户。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。
* **必填**：是

#### `database.password`

设置数据库用户密码。

* **类型**：`password`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：使用敏感配置管理方式提供，不要把真实密码写入普通配置文件或日志。

#### `database.dbname`

设置 Connector 连接并捕获变更的数据库名称。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。
* **必填**：是

#### `database.initial.statements`

指定建立数据库连接后执行的初始化 SQL 语句。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `database.tcpKeepAlive`

控制数据库连接是否启用 TCP keepalive。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

### SSL 与连接路由

#### `database.sslmode`

设置连接 YugabyteDB YSQL 时使用的 SSL 模式。

* **类型**：`string`
* **默认值**：`prefer`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`allow`、`prefer`、`disable`、`verify-ca`、`require`、`verify-full`。

#### `database.sslcert`

设置客户端 SSL 证书文件路径。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `database.sslpassword`

设置访问客户端 SSL 私钥所需的密码。

* **类型**：`password`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：使用敏感配置管理方式提供，不要把真实密码写入普通配置文件或日志。

#### `database.sslrootcert`

设置用于验证服务端证书的根 CA 证书文件路径。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `database.sslkey`

设置客户端 SSL 私钥文件路径。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `database.sslfactory`

设置 PostgreSQL JDBC 连接使用的 SSL Factory 类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `yb.load.balance.connections`

控制 YugabyteDB JDBC 连接的负载均衡和节点偏好。

* **类型**：`string`
* **默认值**：`only-primary`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

### 逻辑复制与 Slot

#### `plugin.name`

设置服务端逻辑解码插件。

* **类型**：`string`
* **默认值**：`yboutput`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`decoderbufs`、`yboutput`、`pgoutput`。目标推荐路径使用 `yboutput`。

#### `slot.name`

设置默认流式读取使用的 replication slot 名称。

* **类型**：`string`
* **默认值**：`debezium`
* **重要级别**：中
* **有效值 / 注意事项**：使用小写字母、数字和下划线，最长 63 个字符。

#### `slot.lsn.type`

设置 replication slot 位置使用的 LSN 表示类型。

* **类型**：`string`
* **默认值**：`SEQUENCE`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`sequence`、`hybrid_time`。

#### `publication.name`

设置默认流式读取使用的 publication 名称。

* **类型**：`string`
* **默认值**：`dbz_publication`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `publication.autocreate.mode`

控制 Connector 是否以及如何自动创建 publication。

* **类型**：`string`
* **默认值**：`all_tables`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`filtered`、`disabled`、`all_tables`；使用 `disabled` 时需预先创建 publication。

#### `replica.identity.autoset.values`

按表名模式自动设置目标表的 replica identity。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `slot.drop.on.stop`

控制 Connector 正常停止时是否删除 replication slot。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `slot.stream.params`

设置启动逻辑复制流时传递给服务端插件的参数。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `slot.max.retries`

设置获取 replication slot 时的最大重试次数。

* **类型**：`int`
* **默认值**：`6`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `slot.retry.delay.ms`

设置获取 replication slot 的重试间隔，单位毫秒。

* **类型**：`long`
* **默认值**：`10000`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `status.update.interval.ms`

设置向数据库发送复制状态更新的时间间隔，单位毫秒。

* **类型**：`int`
* **默认值**：`10000`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `xmin.fetch.interval.ms`

设置查询 replication slot `xmin` 的时间间隔，单位毫秒。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `flush.lsn.source`

控制已处理 LSN 是否回写到源端 replication slot。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `streaming.mode`

选择默认单流式读取或 YugabyteDB 并行流式读取。

* **类型**：`string`
* **默认值**：`DEFAULT`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`default`、`parallel`；`parallel` 需同时配置数量匹配的 `slot.names`、`publication.names` 和 `slot.ranges`。

#### `slot.names`

设置并行流式读取使用的 replication slot 列表。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：仅用于并行流式读取，以逗号分隔；条目数量必须与 `publication.names` 和 `slot.ranges` 一致。

#### `publication.names`

设置并行流式读取使用的 publication 列表。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：仅用于并行流式读取，以逗号分隔；条目数量必须与 `slot.names` 和 `slot.ranges` 一致。

#### `slot.ranges`

设置并行流式读取中各 slot 对应的哈希范围。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：以分号分隔哈希范围；范围数量必须与 slot 和 publication 数量一致，并完整覆盖 `0` 到 `65536`。

#### `ysql.major.upgrade`

控制 Connector 是否按 YSQL 主版本升级场景处理复制状态。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

### 捕获范围与消息选择

#### `skipped.operations`

指定不生成变更事件的数据库操作类型。

* **类型**：`list`
* **默认值**：`["t"]`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `message.prefix.include.list`

设置要包含的逻辑消息前缀正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：与 `message.prefix.exclude.list` 互斥。

#### `message.prefix.exclude.list`

设置要排除的逻辑消息前缀正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：与 `message.prefix.include.list` 互斥。

#### `table.include.list`

设置要捕获的表的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：使用完全限定表名的正则表达式；与 `table.exclude.list` 互斥。

#### `table.exclude.list`

设置不捕获的表的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：使用完全限定表名的正则表达式；与 `table.include.list` 互斥。

#### `table.ignore.builtin`

控制是否忽略 YugabyteDB 或 PostgreSQL 内置表。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `schema.include.list`

设置要捕获的 Schema 的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：与 `schema.exclude.list` 互斥。

#### `schema.exclude.list`

设置不捕获的 Schema 的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：与 `schema.include.list` 互斥。

#### `column.include.list`

设置要包含在变更事件中的列的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：使用完全限定列名的正则表达式；与 `column.exclude.list` 互斥。

#### `column.exclude.list`

设置要从变更事件中排除的列的正则表达式列表。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：使用完全限定列名的正则表达式；与 `column.include.list` 互斥。

#### `message.key.columns`

为指定表设置用于生成 Kafka 消息键的列。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

### 快照与增量快照

#### `snapshot.mode`

控制启动时是否以及如何执行现有数据快照。

* **类型**：`string`
* **默认值**：`initial`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`always`、`never`、`initial_only`、`initial`、`parallel`、`custom`。

#### `snapshot.custom.class`

设置 `snapshot.mode=custom` 时使用的自定义 Snapshotter 类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `snapshot.mode=custom` 时使用。

#### `snapshot.delay.ms`

设置 Connector 启动后开始快照前的延迟，单位毫秒。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `snapshot.include.collection.list`

限制快照阶段要读取的表集合。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `snapshot.fetch.size`

设置快照查询每批从数据库读取的最大行数。

* **类型**：`int`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `snapshot.max.threads`

设置快照阶段可并行处理的最大线程数。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `snapshot.lock.timeout.ms`

设置快照获取表锁时的最大等待时间，单位毫秒。

* **类型**：`long`
* **默认值**：`10000`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `yb.consistent.snapshot`

控制 YugabyteDB 快照是否使用一致性快照行为。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `incremental.snapshot.chunk.size`

设置增量快照每个 Chunk 读取的最大行数。

* **类型**：`int`
* **默认值**：`1024`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `incremental.snapshot.watermarking.strategy`

设置增量快照窗口使用的水位标记策略。

* **类型**：`string`
* **默认值**：`INSERT_INSERT`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`insert_delete`、`insert_insert`。

#### `snapshot.select.statement.overrides`

指定需要覆盖默认快照 SELECT 语句的表。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `snapshot.tables.order.by.row.count`

控制快照多个表时是否按表行数排序。

* **类型**：`string`
* **默认值**：`disabled`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`disabled`、`ascending`、`descending`。

### 队列、轮询与恢复

#### `event.processing.failure.handling.mode`

控制 Connector 处理变更事件失败时的行为。

* **类型**：`string`
* **默认值**：`fail`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`warn`、`fail`、`ignore`、`skip`。

#### `max.batch.size`

设置每次迭代可处理并提交给 Kafka Connect 的最大事件数。

* **类型**：`int`
* **默认值**：`2048`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `max.queue.size`

设置阻塞队列可容纳的最大事件数。

* **类型**：`int`
* **默认值**：`8192`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `poll.interval.ms`

设置队列中没有可用事件时的轮询等待时间，单位毫秒。

* **类型**：`long`
* **默认值**：`500`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `max.queue.size.in.bytes`

按字节限制阻塞队列大小。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `retriable.restart.connector.wait.ms`

设置可重试异常后重启 Connector 前的等待时间，单位毫秒。

* **类型**：`long`
* **默认值**：`10000`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `query.fetch.size`

设置流式查询每批从数据库读取的记录数。

* **类型**：`int`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `errors.max.retries`

设置连接错误后的最大重试次数。

* **类型**：`int`
* **默认值**：`-1`
* **重要级别**：低
* **有效值 / 注意事项**：`-1` 表示不限制重试次数，`0` 表示禁用重试，正整数表示最大重试次数。

### 事件格式与数据类型

#### `provide.transaction.metadata`

控制事件中是否包含数据库事务边界和事务元数据。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `decimal.handling.mode`

设置 DECIMAL 和 NUMERIC 值在变更事件中的表示方式。

* **类型**：`string`
* **默认值**：`precise`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`string`、`double`、`precise`。

#### `time.precision.mode`

设置时间、日期和时间戳值的精度与 Kafka Connect 类型映射。

* **类型**：`string`
* **默认值**：`adaptive`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`adaptive`、`adaptive_time_microseconds`、`connect`。

#### `primary.key.hash.columns`

指定参与 YugabyteDB 主键哈希计算的列。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `hstore.handling.mode`

设置 PostgreSQL `hstore` 值在事件中的表示方式。

* **类型**：`string`
* **默认值**：`json`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`json`、`map`。

#### `binary.handling.mode`

设置二进制列值在事件中的表示方式。

* **类型**：`string`
* **默认值**：`bytes`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`bytes`、`base64`、`hex`、`base64-url-safe`。

#### `schema.name.adjustment.mode`

控制生成的 Kafka Connect Schema 名称如何调整为合法名称。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`none`、`avro_unicode`、`avro`。

#### `interval.handling.mode`

设置 PostgreSQL `INTERVAL` 值在事件中的表示方式。

* **类型**：`string`
* **默认值**：`numeric`
* **重要级别**：低
* **有效值 / 注意事项**：可用值：`string`、`numeric`。

#### `schema.refresh.mode`

控制检测到表结构变化时刷新内存 Schema 的条件。

* **类型**：`string`
* **默认值**：`columns_diff`
* **重要级别**：中
* **有效值 / 注意事项**：可用值：`columns_diff`、`columns_diff_exclude_unchanged_toast`。

#### `unavailable.value.placeholder`

设置源端未提供列值时写入事件的占位值。

* **类型**：`string`
* **默认值**：`__debezium_unavailable_value`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `converters`

指定要应用于事件字段的自定义 Converter 别名。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `post.processors`

指定在事件发送前执行的后处理器别名。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `tombstones.on.delete`

控制删除事件后是否发送相同消息键的 tombstone 记录。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `topic.naming.strategy`

设置生成数据 Topic 名称的 TopicNamingStrategy 实现类。

* **类型**：`class`
* **默认值**：`io.debezium.schema.SchemaTopicNamingStrategy`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `include.schema.comments`

控制生成的 Kafka Connect Schema 是否包含数据库列注释。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `include.unknown.datatypes`

控制是否在事件中包含 Connector 未知的数据类型。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `sourceinfo.struct.maker`

设置构建事件 `source` 结构的 SourceInfoStructMaker 实现类。

* **类型**：`class`
* **默认值**：`io.debezium.connector.postgresql.PostgresSourceInfoStructMaker`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

### 信号、通知与字段处理

#### `heartbeat.interval.ms`

设置生成心跳事件的时间间隔，单位毫秒。

* **类型**：`int`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `heartbeat.topics.prefix`

设置心跳 Topic 名称使用的前缀。

* **类型**：`string`
* **默认值**：`__debezium-heartbeat`
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `heartbeat.action.query`

设置每次生成心跳事件时在源数据库执行的 SQL。

* **类型**：`string`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `signal.data.collection`

设置 Connector 读取运行时信号的数据集合。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `signal.poll.interval.ms`

设置轮询信号数据集合的时间间隔，单位毫秒。

* **类型**：`long`
* **默认值**：`5000`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `signal.enabled.channels`

设置启用的 Debezium 信号通道。

* **类型**：`list`
* **默认值**：`["source"]`
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `notification.enabled.channels`

设置启用的 Debezium 通知通道。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `notification.sink.topic.name`

设置 Sink 通知通道发送通知的 Kafka Topic。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `custom.metric.tags`

设置附加到 Connector 指标的自定义标签。

* **类型**：`list`
* **默认值**：无
* **重要级别**：低
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `column.mask.hash.([^.]+).with.salt.(.+)`

按列规则使用指定哈希算法和盐值对字段值进行哈希遮蔽。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `column.mask.with.(d+).chars`

按列规则用固定长度的星号遮蔽字段值。

* **类型**：`string`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `column.truncate.to.(d+).chars`

按列规则把字段值截断到指定字符数。

* **类型**：`int`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `column.propagate.source.type`

把匹配列的源端类型属性传播到生成的 Schema 参数。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

#### `datatype.propagate.source.type`

把匹配数据类型的源端类型属性传播到生成的 Schema 参数。

* **类型**：`list`
* **默认值**：无
* **重要级别**：中
* **有效值 / 注意事项**：该 ConfigDef 未声明额外的公开取值范围。

## 最佳实践

### 首次接入时建立存量基线并持续同步

**适用业务场景**：首次接入已有业务表，需要先把当前数据写入 Kafka，再无缝接收后续变更。该方式适合下游尚未建立完整基线的初始化阶段。

**配置示例**：

在快速开始配置中保留以下属性：

```properties theme={null}
snapshot.mode=initial
```

**关键说明**：在启动 Connector 前确认 slot、publication、逻辑解码插件和账号权限均已准备完成。快照记录的操作类型为 READ；快照完成后 Connector 转入流式读取。下游应按主键执行幂等处理，并允许在故障恢复或 offset 提交边界出现记录重放。

### 已有数据基线时直接读取后续变更

**适用业务场景**：下游已经通过其他受控流程完成存量初始化，只需要从逻辑复制位置开始接收新的变更，避免 Connector 再次发送全量快照。

**配置示例**：

在快速开始配置中将快照模式改为：

```properties theme={null}
snapshot.mode=never
```

**关键说明**：使用该模式前必须确认下游基线与复制起点一致，并确保 replication slot 仍保留所需历史。该配置不会补发 slot 起点之前的存量数据；slot 被删除或所需复制历史已清理时，不能保证从预期位置恢复。

## 监控

### 监控内容

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

## 限制条件

* 并行流式读取要求只捕获一张表，并为每个 slot、publication 和哈希范围组合创建 Task；三类列表数量必须一致，范围必须完整覆盖 `0` 到 `65536`。
* Connector 不保证跨 Topic、Kafka Partition 或多个并行 Task 的全局顺序。
* Kafka Connect offset 尚未持久化时发生故障可能导致记录重放；不能将普通运行方式视为无条件的端到端恰好一次，下游应提供幂等或去重能力。
* 恢复依赖 replication slot 持续存在并保留所需复制历史。slot 被删除或历史已清理后，不能保证从原位置继续读取。
* `initial_only` 等不进入流式读取的快照模式只提供存量快照，不提供持续 CDC。

## 常见问题

### 为什么 Connector 无法连接 YugabyteDB？

主机名、端口、数据库名、账号权限或逻辑解码插件不匹配都可能导致启动失败。检查 `database.hostname`、`database.port`、`database.dbname` 和 `plugin.name`，确认账号具有登录及逻辑复制所需权限，并确认 Worker 可以访问 YSQL 服务。修正后重启 Connector 并检查 Task 状态。

### 为什么没有生成预期的数据 Topic？

先检查 `topic.prefix` 是否符合命名规则，以及 Connector 是否实际捕获到匹配的数据变更。普通表 Topic 名称还包含 Schema 和表名；确认下游订阅的是完整 Topic 名称，并检查 publication 是否包含目标表、replication slot 是否处于可用状态。

### 为什么重启后出现重复记录？

数据库变更已发送到 Kafka、但对应 offset 尚未持久化时发生故障，恢复后可能重新读取这些变更。检查 Connect 的 offset 提交状态和 replication slot 位置，并让下游按稳定业务键或消息键执行幂等写入或去重；不要仅依赖 Connector 重启来避免重复。

### 为什么初始快照完成后没有持续收到变更？

检查 `snapshot.mode` 是否设置为 `initial_only`，该模式完成快照后不会进入持续流式读取。若需要持续 CDC，使用包含流式阶段的模式，并确认 slot、publication、插件和数据库权限仍然有效。
