Skip to main content

概述

YugabyteDB Source Connector 通过 YSQL 逻辑复制读取 YugabyteDB 表的变更,并将 INSERT、UPDATE、DELETE 和 TRUNCATE 事件写入 Kafka Topic。普通表变更使用由 topic.prefix、Schema 和表名组成的 Topic 名称,消息值采用包含 beforeaftersourceop 和时间戳等字段的 Debezium 事件结构。 Connector 可以先为现有数据创建快照,再持续读取增量变更,也可以在已有数据基线时直接进入流式读取。它适合将 YugabyteDB 中的业务数据同步到流处理、搜索、分析或异构数据系统。Kafka 记录键通常来自表主键;没有可用主键时,记录键可能为空。

前置条件

  • 使用具有 LOGINREPLICATION 属性的专用 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
将数据库地址、账号、密码和数据库名占位符替换为实际环境值,并确认端口、逻辑解码插件、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
  • 重要级别:低
  • 有效值 / 注意事项:可用值:nonerestart

错误处理与 Topic 创建

errors.retry.timeout

设置失败操作的总重试时长,单位毫秒。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 表示不重试,-1 表示无限重试,正整数表示总重试时长。

errors.retry.delay.max.ms

设置连续重试之间的最大等待时间,单位毫秒。
  • 类型long
  • 默认值60000
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负毫秒值;达到上限后会在延迟中加入抖动。

errors.tolerance

控制遇到记录处理错误时是立即失败还是跳过错误记录。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可用值:noneall

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
  • 重要级别:中
  • 有效值 / 注意事项:可用值(不区分大小写):requiredrequested

transaction.boundary

指定 Source Task 提交 Producer 事务时采用的边界。
  • 类型string
  • 默认值poll
  • 重要级别:中
  • 有效值 / 注意事项:可用值(不区分大小写):intervalpollconnector

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
  • 重要级别:中
  • 有效值 / 注意事项:可用值:allowpreferdisableverify-carequireverify-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
  • 重要级别:中
  • 有效值 / 注意事项:可用值:decoderbufsyboutputpgoutput。目标推荐路径使用 yboutput

slot.name

设置默认流式读取使用的 replication slot 名称。
  • 类型string
  • 默认值debezium
  • 重要级别:中
  • 有效值 / 注意事项:使用小写字母、数字和下划线,最长 63 个字符。

slot.lsn.type

设置 replication slot 位置使用的 LSN 表示类型。
  • 类型string
  • 默认值SEQUENCE
  • 重要级别:中
  • 有效值 / 注意事项:可用值:sequencehybrid_time

publication.name

设置默认流式读取使用的 publication 名称。
  • 类型string
  • 默认值dbz_publication
  • 重要级别:中
  • 有效值 / 注意事项:该 ConfigDef 未声明额外的公开取值范围。

publication.autocreate.mode

控制 Connector 是否以及如何自动创建 publication。
  • 类型string
  • 默认值all_tables
  • 重要级别:中
  • 有效值 / 注意事项:可用值:filtereddisabledall_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
  • 重要级别:低
  • 有效值 / 注意事项:可用值:defaultparallelparallel 需同时配置数量匹配的 slot.namespublication.namesslot.ranges

slot.names

设置并行流式读取使用的 replication slot 列表。
  • 类型string
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:仅用于并行流式读取,以逗号分隔;条目数量必须与 publication.namesslot.ranges 一致。

publication.names

设置并行流式读取使用的 publication 列表。
  • 类型string
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:仅用于并行流式读取,以逗号分隔;条目数量必须与 slot.namesslot.ranges 一致。

slot.ranges

设置并行流式读取中各 slot 对应的哈希范围。
  • 类型string
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:以分号分隔哈希范围;范围数量必须与 slot 和 publication 数量一致,并完整覆盖 065536

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
  • 重要级别:中
  • 有效值 / 注意事项:可用值:alwaysneverinitial_onlyinitialparallelcustom

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_deleteinsert_insert

snapshot.select.statement.overrides

指定需要覆盖默认快照 SELECT 语句的表。
  • 类型string
  • 默认值:无
  • 重要级别:中
  • 有效值 / 注意事项:该 ConfigDef 未声明额外的公开取值范围。

snapshot.tables.order.by.row.count

控制快照多个表时是否按表行数排序。
  • 类型string
  • 默认值disabled
  • 重要级别:中
  • 有效值 / 注意事项:可用值:disabledascendingdescending

队列、轮询与恢复

event.processing.failure.handling.mode

控制 Connector 处理变更事件失败时的行为。
  • 类型string
  • 默认值fail
  • 重要级别:中
  • 有效值 / 注意事项:可用值:warnfailignoreskip

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
  • 重要级别:中
  • 有效值 / 注意事项:可用值:stringdoubleprecise

time.precision.mode

设置时间、日期和时间戳值的精度与 Kafka Connect 类型映射。
  • 类型string
  • 默认值adaptive
  • 重要级别:中
  • 有效值 / 注意事项:可用值:adaptiveadaptive_time_microsecondsconnect

primary.key.hash.columns

指定参与 YugabyteDB 主键哈希计算的列。
  • 类型string
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:该 ConfigDef 未声明额外的公开取值范围。

hstore.handling.mode

设置 PostgreSQL hstore 值在事件中的表示方式。
  • 类型string
  • 默认值json
  • 重要级别:低
  • 有效值 / 注意事项:可用值:jsonmap

binary.handling.mode

设置二进制列值在事件中的表示方式。
  • 类型string
  • 默认值bytes
  • 重要级别:低
  • 有效值 / 注意事项:可用值:bytesbase64hexbase64-url-safe

schema.name.adjustment.mode

控制生成的 Kafka Connect Schema 名称如何调整为合法名称。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项:可用值:noneavro_unicodeavro

interval.handling.mode

设置 PostgreSQL INTERVAL 值在事件中的表示方式。
  • 类型string
  • 默认值numeric
  • 重要级别:低
  • 有效值 / 注意事项:可用值:stringnumeric

schema.refresh.mode

控制检测到表结构变化时刷新内存 Schema 的条件。
  • 类型string
  • 默认值columns_diff
  • 重要级别:中
  • 有效值 / 注意事项:可用值:columns_diffcolumns_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 前确认 slot、publication、逻辑解码插件和账号权限均已准备完成。快照记录的操作类型为 READ;快照完成后 Connector 转入流式读取。下游应按主键执行幂等处理,并允许在故障恢复或 offset 提交边界出现记录重放。

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

适用业务场景:下游已经通过其他受控流程完成存量初始化,只需要从逻辑复制位置开始接收新的变更,避免 Connector 再次发送全量快照。 配置示例 在快速开始配置中将快照模式改为:
关键说明:使用该模式前必须确认下游基线与复制起点一致,并确保 replication slot 仍保留所需历史。该配置不会补发 slot 起点之前的存量数据;slot 被删除或所需复制历史已清理时,不能保证从预期位置恢复。

监控

监控内容

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

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

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

常见问题

为什么 Connector 无法连接 YugabyteDB?

主机名、端口、数据库名、账号权限或逻辑解码插件不匹配都可能导致启动失败。检查 database.hostnamedatabase.portdatabase.dbnameplugin.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、插件和数据库权限仍然有效。