概述
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。配置
运行与序列化
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,再无缝接收后续变更。该方式适合下游尚未建立完整基线的初始化阶段。 配置示例: 在快速开始配置中保留以下属性:已有数据基线时直接读取后续变更
适用业务场景:下游已经通过其他受控流程完成存量初始化,只需要从逻辑复制位置开始接收新的变更,避免 Connector 再次发送全量快照。 配置示例: 在快速开始配置中将快照模式改为:监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号。仅在启用了相应错误处理时关注 DLQ 活动。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 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、插件和数据库权限仍然有效。