Skip to main content

概述

Debezium PostgreSQL Source Connector 将 PostgreSQL 表的存量数据和已提交的行级变化发送到 Kafka,适用于增量数据同步、搜索索引更新、缓存刷新和变化事件消费。Connector 在 PostgreSQL 到 Kafka 的链路中负责数据采集,先通过快照建立数据基线,再通过逻辑复制持续捕获插入、更新和删除。 一个 Connector 捕获一个数据库。默认情况下,每张表的事件发送到以主题前缀、Schema 名和表名组成的 Topic,消息键通常来自表的主键。事件包含操作类型、变更前后数据和源端位置等信息;实际序列化格式由 Kafka Connect 的 Converter 决定。快照记录与后续增量共同支持下游重建当前状态,消费者仍需正确处理重复事件、删除和字段不可用的情况。

前置条件

  • PostgreSQL 必须启用逻辑复制,设置 wal_level=logical,并为 max_replication_slotsmax_wal_senders 预留足够容量;服务端认证规则必须允许采集账号建立数据库连接和逻辑复制连接。使用 pgoutput 时依赖 PostgreSQL 内置逻辑解码插件;使用 decoderbufs 时须在服务端安装对应插件。
  • 采集账号需要 LOGINREPLICATION 权限,或托管 PostgreSQL 提供的等效复制角色;执行快照和元数据读取还需要目标数据库的 CONNECT、目标 Schema 的 USAGE 和捕获表的 SELECT 权限。复制权限不自动包含表读取权限。
  • 使用 pgoutput 时,publication 的表集合及发布操作必须覆盖所需增量。快速开始采用 DBA 预建的按表 publication,Connector 不负责创建或修改它;若使用 filtered 自动管理,账号还需相应的数据库 CREATE、publication 所有权及表所有权权限。创建 FOR ALL TABLES publication 需要超级用户权限,不应仅为启动 Connector 而给采集账号授予该权限。
  • 需要同步更新和删除的表必须具备适当的 REPLICA IDENTITY。有主键的表通常可使用 DEFAULT;需要完整旧行时再评估 FULL 及其 WAL 开销。配置消息键不会自动修改数据库复制标识,无可用复制标识时 PostgreSQL 可能拒绝发布表的更新或删除。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、源 PostgreSQL 数据库及需要采集的表,确认网络连通和访问权限,并由 DBA 为这些表准备 publication。Connector 的创建和管理操作参见 管理 Connector
替换数据库地址、账号、密码、数据库名及采集资源占位符。槽名使用小写字母、数字和下划线,不超过 63 个字符,并为每个独立 Connector 选择独立槽名。表筛选表达式匹配完整的 schema.table;在 properties 中匹配字面点可写成 [.],例如 public[.]orders。预建 publication 必须包含匹配的表,否则可能出现有快照但没有增量的情况。密码仅为占位符,实际部署可使用已配置的 ConfigProvider 引用,避免在配置副本中保存明文凭证。 默认快照模式为 initial:没有已保存的 Offset 或首次快照未完成时读取存量,随后持续捕获增量。默认数据库端口为 5432,任务上限为 1,因此示例省略这些默认值。pgoutput 是显式选择,逻辑解码插件的默认值并不是它。未设置 Connector 级 Converter 时沿用 Worker 的序列化配置。

配置

Connector 与序列化

connector.class

指定 Source Connector 实现。
  • 类型STRING
  • 默认值:无
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:使用 io.debezium.connector.postgresql.PostgresConnector

tasks.max

设置任务数上限。
  • 类型INT
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1。此 Connector 只创建一个 Task,提高上限不会并行拆分表或复制流。

key.converter

指定消息键的序列化 Converter。
  • 类型CLASS
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时使用 Worker Converter;指定类须实现 Kafka Connect Converter 接口。

value.converter

指定消息值的序列化 Converter。
  • 类型CLASS
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时使用 Worker Converter;指定类须实现 Kafka Connect Converter 接口。事件并非固定的 JSON 文本。

数据库连接与安全

database.hostname

源数据库的主机名或 IP 地址。
  • 类型STRING
  • 默认值null
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:使用可解析的数据库地址。

database.port

源数据库端口。
  • 类型INT
  • 默认值5432
  • 重要级别:高
  • 有效值 / 注意事项:按实际 PostgreSQL 服务端口设置。

database.user

数据库采集账号。
  • 类型STRING
  • 默认值null
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:同时满足复制、快照读取和所选 publication 管理方式的权限要求。

database.password

数据库账号密码。
  • 类型PASSWORD
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:是否需要密码由数据库认证方式决定;使用 Secret 占位或 ConfigProvider 引用,不在共享配置和日志中保留真实密码。

database.dbname

需要采集的数据库名称。
  • 类型STRING
  • 默认值null
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:一个 Connector 连接一个数据库;表筛选使用 schema.table 而非数据库名前缀。

database.query.timeout.ms

数据库查询执行超时,单位毫秒。
  • 类型INT
  • 默认值600000
  • 重要级别:低
  • 有效值 / 注意事项0 表示不限制查询时长;不等同于复制连接的状态更新或 LSN 刷新超时。

database.initial.statements

在建立 JDBC 连接时执行初始化语句。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:多条语句以分号分隔,字面分号用双分号转义。连接可能多次建立,通常仅用于会话设置,不用于业务 DML。

database.sslmode

控制 PostgreSQL TLS 连接与证书校验。
  • 类型STRING
  • 默认值prefer
  • 重要级别:中
  • 有效值 / 注意事项disableallowpreferrequireverify-caverify-fullprefer 允许回退到非加密连接;require 要求加密但不等于主机身份校验,verify-full 还校验证书主机名。

database.sslcert

客户端 TLS 证书文件路径。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:在数据库要求客户端证书时设置,文件须可由执行 Task 的进程读取。

database.sslkey

客户端 TLS 私钥文件路径。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:须符合 PostgreSQL JDBC 私钥格式要求,限制文件访问权限,不把私钥内容写入配置。

database.sslpassword

访问客户端私钥的密码。
  • 类型PASSWORD
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:配合 database.sslkey 使用;只通过安全凭证管理方式提供。

database.sslrootcert

验证服务端证书的根证书文件路径。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:配合证书校验模式使用,执行 Task 的进程须能读取文件。

database.sslfactory

指定创建 SSL Socket 的工厂类。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:使用符合 PostgreSQL JDBC 要求且已安装的工厂类;不应通过不校验证书的工厂规避生产环境的身份校验。

database.tcpKeepAlive

控制 TCP Keepalive 探测。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项truefalse;实际探测间隔受操作系统与连接环境影响。

connection.validation.timeout.ms

连接校验的最长等待时间,单位毫秒。
  • 类型LONG
  • 默认值60000
  • 重要级别:低
  • 有效值 / 注意事项:连接校验通过不代表所有快照、publication 创建或更新操作均已获授权。

driver.<PostgreSQL JDBC property>

透传 PostgreSQL JDBC 连接属性。
  • 类型STRING 属性透传
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:用实际 JDBC 属性名替换名称中的占位部分,driver. 前缀会被移除。类型、默认值和校验规则由 PostgreSQL JDBC 定义;不通过连接属性嵌入明文凭证。

逻辑复制与 publication

plugin.name

选择服务端逻辑解码插件。
  • 类型STRING
  • 默认值decoderbufs
  • 重要级别:中
  • 有效值 / 注意事项decoderbufspgoutput;使用内置 pgoutput 时必须显式设置。

slot.name

指定读取变化的逻辑复制槽。
  • 类型STRING
  • 默认值debezium
  • 重要级别:中
  • 有效值 / 注意事项:匹配 [a-z0-9_]{1,63}。各独立 Connector 使用独立槽,两个活跃复制连接不能共享同一槽。

slot.drop.on.stop

控制正常停止时是否删除复制槽。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:长期采集保留槽有助于恢复;启用后可能失去停机期间所需 WAL。此选项不表示创建临时槽。

slot.failover

控制是否创建 failover slot。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:创建 failover slot 要求连接 PostgreSQL 17 或更新版本的 primary;不满足条件时创建普通槽。不会自动升级已有槽或完成 standby 同步和故障切换配置。

slot.stream.params

向逻辑解码插件提供可选参数。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:参数以分号分隔;仅使用所选插件支持的参数,不把其他解码插件的参数套用于 pgoutput

slot.max.retries

连接复制槽时的最大重试次数。
  • 类型INT
  • 默认值6
  • 重要级别:低
  • 有效值 / 注意事项:与 slot.retry.delay.ms 配合;重试不能修复权限不足、槽占用设计错误或 WAL 丢失。

slot.retry.delay.ms

复制槽连接重试间隔,单位毫秒。
  • 类型LONG
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:按整数毫秒设置;运行时按整型读取,避免使用超出整型范围的值。

publication.name

指定 pgoutput 使用的 publication。
  • 类型STRING
  • 默认值dbz_publication
  • 重要级别:中
  • 有效值 / 注意事项:仅适用于 pgoutput;publication 范围与 Connector 过滤范围应协调一致。

publication.autocreate.mode

控制 publication 的创建与表集合管理。
  • 类型STRING
  • 默认值all_tables
  • 重要级别:中
  • 有效值 / 注意事项disabled 要求预建;all_tables 在不存在时创建全表 publication;filtered 创建或更新按表 publication 的表集合;no_tables 在不存在时创建空 publication。filtered 拒绝已有的 FOR ALL TABLES publication 和空匹配集合,更新可能影响其他使用者,需具备数据库端权限。

replica.identity.autoset.values

按表规则自动修改复制标识。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔的表正则与复制标识映射,格式为 schema.table:replica identity;支持 DEFAULTINDEX index_nameFULLNOTHING。仅适用于 pgoutput,会覆盖数据库现有设置;指定索引需满足 PostgreSQL 复制标识索引要求,并具备修改表的权限。

publish.via.partition.root

控制是否通过分区根表发布变化。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:适用于 pgoutputtrue 使用分区根表身份,false 直接发布分区变化。需与 publication 的数据库端设置协调。

status.update.interval.ms

向 PostgreSQL 发送复制连接状态的间隔,单位毫秒。
  • 类型INT
  • 默认值10000
  • 重要级别:中
  • 有效值 / 注意事项:正整数;不等同于 Kafka Offset 提交间隔。

flush.lsn.source

控制是否向 PostgreSQL 确认已处理的 LSN。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:设为 false 时须由外部机制管理 LSN 推进及 WAL 回收,不能仅靠 Connector 心跳替代。

lsn.flush.timeout.ms

LSN 刷新操作的最长等待时间,单位毫秒。
  • 类型LONG
  • 默认值30000
  • 重要级别:中
  • 有效值 / 注意事项:正整数;超时后使用 lsn.flush.timeout.action 决定处理方式。

lsn.flush.timeout.action

控制 LSN 刷新超时后的处理方式。
  • 类型STRING
  • 默认值fail
  • 重要级别:中
  • 有效值 / 注意事项fail 停止 Connector,warn 记录警告并继续,ignore 忽略超时并继续;继续运行不表示源端已确认消费位置。

xmin.fetch.interval.ms

控制复制槽 xmin 的查询间隔,单位毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 禁用跟踪;更频繁查询增加数据库开销。

捕获范围与字段处理

schema.include.list

仅采集匹配的 Schema。
  • 类型LIST
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:逗号分隔的合法正则表达式;不能与 schema.exclude.list 同时设置。还需满足表过滤和 publication 范围。

schema.exclude.list

排除匹配的 Schema。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔的合法正则表达式;不能与 schema.include.list 同时设置。

table.include.list

仅采集匹配的表。
  • 类型LIST
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:逗号分隔的正则表达式,匹配完整 schema.table;不能与 table.exclude.list 同时设置。不自动扩展预建 publication。

table.exclude.list

排除匹配的表。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔的正则表达式,匹配完整 schema.table;不能与 table.include.list 同时设置。

table.ignore.builtin

控制是否忽略内置表。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:启用时排除 pg_cataloginformation_schema、以 pg_temp 开头的 Schema 及 spatial_ref_sys 表。

column.include.list

仅输出匹配的列。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:正则匹配 schema.table.column;不能与 column.exclude.list 同时设置。列过滤不代替数据库权限隔离。

column.exclude.list

从事件字段中排除匹配列。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:正则匹配 schema.table.column;不能与 column.include.list 同时设置,不应据此推断数据库不读取这些列。

message.prefix.include.list

仅捕获指定前缀的逻辑解码消息。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:合法正则表达式列表;不能与 message.prefix.exclude.list 同时设置。未过滤时捕获所有前缀。

message.prefix.exclude.list

排除指定前缀的逻辑解码消息。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:合法正则表达式列表;不能与 message.prefix.include.list 同时设置。

skipped.operations

指定流式阶段跳过的操作。
  • 类型LIST
  • 默认值t
  • 重要级别:低
  • 有效值 / 注意事项c 插入、u 更新、d 删除、t 清空表,none 表示不跳过。跳过操作会改变下游能够重建的数据状态。

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

通过哈希和盐对匹配列进行掩码处理。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:名称是参数化模式,实际键用哈希算法和盐替换模式部分;值为匹配完整列名的正则列表。盐按敏感配置管理,不使用真实盐作为共享示例。

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

将匹配列替换为指定数量的星号。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:实际键使用数字字符数,例如 column.mask.with.8.chars;值用于指定列匹配规则。名称中的模式并非可直接提交的属性名。

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

将匹配列截断到指定字符数。
  • 类型INT
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:实际键使用数字字符数,例如 column.truncate.to.16.chars。这是参数化键模式,字符数在实际键名中,列选择由对应匹配规则指定;截断不可逆,须评估业务字段影响。

快照与存量初始化

snapshot.mode

选择启动时的快照与增量衔接策略。
  • 类型STRING
  • 默认值initial
  • 重要级别:中
  • 有效值 / 注意事项initial 首次初始化后持续采集;always 每次启动重新快照后采集;initial_only 只完成初始化;no_data 不扫描存量行;when_needed 允许在所需日志位置不可用时重新快照;configuration_based 使用配套开关;custom 使用已安装扩展。never 仍接受但已弃用,替代项为 no_data。重快照恢复当前状态,不补回已丢失的历史变化。

snapshot.isolation.mode

设置快照事务隔离级别。
  • 类型STRING
  • 默认值serializable
  • 重要级别:低
  • 有效值 / 注意事项serializablerepeatable_readread_committedread_uncommitted。不同隔离级别的一致性不能视为等价,调整前评估快照期间业务写入影响。

snapshot.locking.mode

控制读取快照元数据时的锁策略。
  • 类型STRING
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项nonesharedcustomnone 要求快照期间避免 Schema 变化;shared 获取 ACCESS SHARE 锁,阻止冲突 DDL 而不阻止普通行写入;custom 使用已安装的锁策略。

snapshot.lock.timeout.ms

快照开始时获取表锁的最长等待时间,单位毫秒。
  • 类型LONG
  • 默认值10000
  • 重要级别:中
  • 有效值 / 注意事项:在所选锁策略需要锁时生效;获取锁超时会终止快照,应先检查冲突事务或 DDL。

snapshot.delay.ms

快照开始前的延迟,单位毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:延迟开始不代替资源容量规划。

streaming.delay.ms

快照完成后开始增量采集前的延迟,单位毫秒。
  • 类型LONG
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:延迟期间变化仍依赖复制槽保留 WAL,需考虑日志保留空间。

snapshot.include.collection.list

限制参与快照的表集合。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:指定表匹配规则,需处于捕获范围内;只影响快照范围,不代替增量的表过滤配置。

snapshot.fetch.size

控制快照读取时的批次取数规模。
  • 类型INT
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:未设置时 PostgreSQL 运行时回退为 10240,该回退不是配置定义的默认值;调整时考虑行大小和 Task 内存。

snapshot.max.threads

设置快照线程数上限。
  • 类型INT
  • 默认值1
  • 重要级别:中
  • 有效值 / 注意事项:影响快照阶段,不增加流式采集 Task 数;提高并行度需评估数据库读取负载。

snapshot.tables.order.by.row.count

控制初始化快照的表处理顺序。
  • 类型STRING
  • 默认值disabled
  • 重要级别:中
  • 有效值 / 注意事项disabled 不按行数排序,ascending 按行数升序,descending 按行数降序。

snapshot.select.statement.overrides

指定需要覆盖快照查询的表。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔 schema.table;每张表需配套同名动态查询配置。只改变存量读取,不改变后续逻辑复制范围。

snapshot.select.statement.overrides.<schema>.<table>

为指定表提供快照 SQL 查询。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:实际表名须列在 snapshot.select.statement.overrides 中,值为 SQL SELECT;缺少值时跳过该覆盖。过滤存量会改变下游初始基线,不能据此推断增量也被按同一条件过滤。

snapshot.query.mode

选择快照查询构建方式。
  • 类型STRING
  • 默认值select_all
  • 重要级别:低
  • 有效值 / 注意事项select_allcustom;自定义模式需提供相应 SPI 实现。

snapshot.query.mode.custom.name

指定自定义快照查询实现名称。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:在 snapshot.query.mode=custom 时指定已安装 io.debezium.snapshot.spi.SnapshotQuery 实现的 name() 值,不是任意类名。

snapshot.locking.mode.custom.name

指定自定义快照锁实现名称。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:在 snapshot.locking.mode=custom 时指定已安装 io.debezium.snapshot.spi.SnapshotLock 实现的 name() 值。

snapshot.mode.custom.name

指定自定义快照策略名称。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:在 snapshot.mode=custom 时指定已安装 Snapshotter 实现的 name() 值。

snapshot.mode.configuration.based.snapshot.data

控制配置驱动策略是否快照数据。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只用于 snapshot.mode=configuration_based

snapshot.mode.configuration.based.snapshot.schema

控制配置驱动策略是否快照 Schema。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只用于 snapshot.mode=configuration_based

snapshot.mode.configuration.based.start.stream

控制配置驱动策略是否在快照后进入增量采集。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只用于 snapshot.mode=configuration_based;该模式下不应假定自动持续采集。

snapshot.mode.configuration.based.snapshot.on.schema.error

控制配置驱动策略遇到 Schema 错误时是否快照。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只用于 snapshot.mode=configuration_based,不代替具体错误的诊断。

snapshot.mode.configuration.based.snapshot.on.data.error

控制配置驱动策略遇到数据位置错误时是否快照。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只用于 snapshot.mode=configuration_based;重新读取当前数据不代表恢复丢失的历史事件。

事件、消息键与数据类型

message.key.columns

按表指定自定义消息键列。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:映射以分号分隔,形式为 schema.table:column1,column2;未指定的表通常使用主键。自定义消息键不提供数据库 WAL 中缺失的旧键值,也不修改复制标识。

tombstones.on.delete

控制删除事件后是否发送同键空值记录。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项true 先发送删除事件,再发送 tombstone;false 只发送删除事件。tombstone 用于 Kafka 日志压缩,不保证下游自动执行数据库删除。

provide.transaction.metadata

控制事务元数据和事件计数输出。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:启用会增加事务上下文和元数据输出;不表示跨 Topic 原子消费或端到端事务保证。

transaction.metadata.factory

指定事务上下文和结构工厂。
  • 类型CLASS
  • 默认值io.debezium.pipeline.txmetadata.DefaultTransactionMetadataFactory
  • 重要级别:低
  • 有效值 / 注意事项:自定义实现须可由插件类加载器访问且符合事务元数据接口要求。

extended.headers.enabled

控制是否添加 Debezium 上下文消息头。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:消息头提供事件追踪与源端识别信息,不代替消息键和事件位置。

decimal.handling.mode

选择 DECIMALNUMERIC 的表示方式。
  • 类型STRING
  • 默认值precise
  • 重要级别:中
  • 有效值 / 注意事项precise 使用 Connect Decimal 与精确数值,string 使用字符串,double 使用浮点数且可能损失精度;需与消费者类型模型协调。

time.precision.mode

选择日期、时间和时间戳的精度及表示方式。
  • 类型STRING
  • 默认值adaptive
  • 重要级别:中
  • 有效值 / 注意事项adaptiveadaptive_time_microsecondsisostringmicrosecondsnanosecondsconnectadaptive 按源列精度转换,connect 使用 Connect 毫秒精度表示,可能不能保留更高精度。

hstore.handling.mode

选择 HSTORE 的表示方式。
  • 类型STRING
  • 默认值json
  • 重要级别:低
  • 有效值 / 注意事项json 使用 JSON 字符串,map 使用键值映射。

binary.handling.mode

选择二进制列的表示方式。
  • 类型STRING
  • 默认值bytes
  • 重要级别:低
  • 有效值 / 注意事项bytesbase64hexbase64-url-safe;最终字节序列化还受 Converter 影响。

interval.handling.mode

选择 INTERVAL 的表示方式。
  • 类型STRING
  • 默认值numeric
  • 重要级别:低
  • 有效值 / 注意事项numeric 近似换算为微秒,string 使用精确的 ISO 格式字符串;包含月等单位时不能假定数值换算保留原始区间语义。

include.unknown.datatypes

控制是否输出无法识别的数据类型字段。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项false 省略不支持字段;true 使用实现相关的二进制表示,不保证所有自定义类型都有对应业务逻辑类型。

unavailable.value.placeholder

表示数据库未提供的 TOAST 字段值。
  • 类型STRING
  • 默认值__debezium_unavailable_value
  • 重要级别:中
  • 有效值 / 注意事项:不是 SQL NULL 或真实业务值;以 hex: 开头时剩余内容按十六进制字节解析。消费者应识别不可用值,避免用占位符覆盖已有的有效字段值。

schema.refresh.mode

控制内存表 Schema 的刷新条件。
  • 类型STRING
  • 默认值columns_diff
  • 重要级别:中
  • 有效值 / 注意事项columns_diff 按列差异刷新;columns_diff_exclude_unchanged_toast 忽略由未变化 TOAST 造成的差异,可能掩盖 TOAST 列删除并使缓存过时。

schema.name.adjustment.mode

调整输出 Schema 名称以适配序列化格式。
  • 类型STRING
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项noneavroavro_unicode;转换前评估名称碰撞及下游兼容性,不改变数据库实际名称。

include.schema.comments

控制是否将表和列注释加入元数据。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:启用会增加内存使用,应评估注释体积及是否包含不应传播的信息。

column.propagate.source.type

为匹配列的字段 Schema 附加源类型和长度信息。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔的完整列名正则表达式,用于类型元数据传播,不改变字段值本身。

datatype.propagate.source.type

按数据库类型名称传播源类型和长度信息。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:逗号分隔的数据库类型名正则表达式,不等同于列名匹配。

sourceinfo.struct.maker

指定构建事件 source Schema 与结构的实现。
  • 类型CLASS
  • 默认值io.debezium.connector.postgresql.PostgresSourceInfoStructMaker
  • 重要级别:低
  • 有效值 / 注意事项:自定义实现须符合 SourceInfoStructMaker 接口;更改结构会影响依赖源端元数据的消费者。

converters

注册 Debezium 自定义字段转换扩展。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:扩展别名配合各自的 别名.type别名.选项 使用,所需类必须安装;此项不等同于 Kafka Connect 的消息序列化 Converter。

post.processors

注册事件后处理扩展。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:处理器别名配合各自的 别名.type 和配置项使用,须安装兼容的实现并评估处理开销。

Topic 命名

topic.prefix

为当前采集源设置 Topic 命名空间。
  • 类型STRING
  • 默认值null
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:只允许字母、数字、点、下划线和短横线;不同采集源使用独立前缀。维护时保留原前缀,避免意外改变输出身份。

topic.naming.strategy

指定 Topic 命名策略。
  • 类型CLASS
  • 默认值io.debezium.schema.SchemaTopicNamingStrategy
  • 重要级别:中
  • 有效值 / 注意事项:默认表 Topic 为主题前缀、Schema、表名以分隔符连接;以下四项用于默认策略,自定义策略可能采用不同配置。

topic.delimiter

指定默认命名策略的分隔符。
  • 类型STRING
  • 默认值.
  • 重要级别:低
  • 有效值 / 注意事项:须满足 Topic 命名字符要求;更改后会改变 Topic 名称。

topic.cache.size

设置 Topic 名称缓存容量。
  • 类型INT
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:用于数据集合到 Topic 名称的缓存,不是 Kafka Topic 数量上限。

topic.heartbeat.prefix

设置默认命名策略的心跳 Topic 前缀。
  • 类型STRING
  • 默认值__debezium-heartbeat
  • 重要级别:低
  • 有效值 / 注意事项:心跳 Topic 默认由该前缀、分隔符及 topic.prefix 组成;启用心跳后才会发送对应记录。

topic.transaction

设置默认命名策略的事务 Topic 后缀。
  • 类型STRING
  • 默认值transaction
  • 重要级别:低
  • 有效值 / 注意事项:事务 Topic 由 topic.prefix、分隔符和该后缀组成;事务元数据需另行启用。

心跳、信号与增量快照

heartbeat.interval.ms

设置发送心跳记录的间隔,单位毫秒。
  • 类型INT
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:非负整数,0 禁用;有复制进度但表事件被过滤时可帮助推进 Offset。Kafka 心跳本身不保证数据库产生新的 WAL。

heartbeat.topics.prefix

设置通用心跳前缀属性。
  • 类型STRING
  • 默认值__debezium-heartbeat
  • 重要级别:低
  • 有效值 / 注意事项:该属性仍可配置,但默认 SchemaTopicNamingStrategy 采用独立的 topic.heartbeat.prefix,不能仅凭本项推断默认策略的 Topic 名。

heartbeat.action.query

指定数据库心跳执行的查询。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:需要非零 heartbeat.interval.ms。若通过写入推进低流量数据库的 WAL,账号需要相应写权限,相关表须进入 pgoutput publication;查询对象位于当前连接数据库。

signal.data.collection

指定通过源端接收信号的表。
  • 类型STRING
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:PostgreSQL 使用 schema.table;使用源端信号时需准备相应信号表、捕获范围和权限。

signal.enabled.channels

选择启用的信号通道。
  • 类型LIST
  • 默认值source
  • 重要级别:中
  • 有效值 / 注意事项:通道名称列表;所选通道需有对应实现和配置,不能只设置名称而忽略通道资源。

signal.poll.interval.ms

检查已注册信号通道的间隔,单位毫秒。
  • 类型LONG
  • 默认值5000
  • 重要级别:中
  • 有效值 / 注意事项:缩短间隔增加轮询开销;不是业务变更的采集间隔。

incremental.snapshot.chunk.size

设置增量快照每个数据块的行数上限。
  • 类型INT
  • 默认值1024
  • 重要级别:中
  • 有效值 / 注意事项:增量快照用于流式运行期间补采数据;块大小影响内存和读取负载,不增加 streaming Task 数。

incremental.snapshot.watermarking.strategy

选择增量快照水位信号的管理策略。
  • 类型STRING
  • 默认值INSERT_INSERT
  • 重要级别:低
  • 有效值 / 注意事项:配置值接受 insert_insertinsert_delete;前者写入开闭两个水位信号,后者写入开启信号并通过删除关闭。涉及信号表写入,需相应权限。

notification.enabled.channels

选择通知输出通道。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:通道名称列表,所选通道须有可用实现;包含 sink 时须指定通知 Topic。此处是 Source 的通知输出,不是另一类 Connector 的方向设置。

notification.sink.topic.name

指定 Topic 通知通道的输出位置。
  • 类型STRING
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项notification.enabled.channels 包含 sink 时必需;不改变业务表事件 Topic。

队列、吞吐与错误恢复

max.batch.size

设置一次向 Worker 交付的记录数上限。
  • 类型INT
  • 默认值2048
  • 重要级别:中
  • 有效值 / 注意事项:批次上限应小于 max.queue.size;增大批次需评估内存和延迟,不保证提升源端提交能力。

max.queue.size

设置尚未交付的变更事件队列条数上限。
  • 类型INT
  • 默认值8192
  • 重要级别:中
  • 有效值 / 注意事项:必须大于 0 且严格大于 max.batch.size;达到上限产生背压而非静默丢弃,积压仍可能增加源端 WAL 保留压力。

max.queue.size.in.bytes

设置变更事件队列的字节容量限制。
  • 类型LONG
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 禁用字节限制;启用后与条数限制共同约束队列,不能当作 Worker JVM 总内存上限。

poll.interval.ms

队列中没有新事件时的等待时间,单位毫秒。
  • 类型LONG
  • 默认值500
  • 重要级别:中
  • 有效值 / 注意事项:不是按该间隔全表查询数据库;减小等待时间会改变空闲轮询开销。

query.fetch.size

控制流式阶段 JDBC 查询的取数规模。
  • 类型INT
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项0 使用 JDBC 默认取数规模;不等同于快照批次或复制消息批次。

event.processing.failure.handling.mode

控制事件处理失败后的行为。
  • 类型STRING
  • 默认值fail
  • 重要级别:中
  • 有效值 / 注意事项failwarnignoreskipfail 停止处理;跳过类策略会遗漏问题事件,告警或继续运行不表示数据完整。

errors.max.retries

控制连接类错误的重试次数上限。
  • 类型INT
  • 默认值-1
  • 重要级别:低
  • 有效值 / 注意事项-1 不限制,0 禁用,正数为重试次数;只影响可重试错误,不能恢复已清理的 WAL。

retriable.restart.connector.wait.ms

可重试异常后重新启动采集逻辑前的等待时间,单位毫秒。
  • 类型LONG
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:与错误重试策略配合;永久权限错误需要修正授权而非持续缩短等待。

executor.shutdown.timeout.ms

Task 执行器关闭的最长等待时间,单位毫秒。
  • 类型LONG
  • 默认值4000
  • 重要级别:中
  • 有效值 / 注意事项:停止过程需要给任务清理留出时间;不代替源端槽和 Offset 的保留规划。

扩展元数据与数据血缘

custom.metric.tags

为 MBean 对象名添加自定义标签。
  • 类型LIST
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:逗号分隔的键值对,例如 environment=production;控制标签数量,不包含凭证或敏感业务内容。

openlineage.integration.enabled

控制是否通过 OpenLineage 输出数据血缘元数据。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:启用前准备相应的 OpenLineage 配置及可用集成资源。

openlineage.integration.config.file.path

指定 OpenLineage 配置文件路径。
  • 类型STRING
  • 默认值./openlineage.yml
  • 重要级别:低
  • 有效值 / 注意事项:文件须可由运行进程读取;配置文件中的敏感信息按凭证管理要求保护。

openlineage.integration.job.namespace

设置血缘任务的命名空间。
  • 类型STRING
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:用于血缘任务身份,不改变采集数据库或 Topic。

openlineage.integration.job.description

设置血缘任务描述。
  • 类型STRING
  • 默认值Debezium change data capture job
  • 重要级别:低
  • 有效值 / 注意事项:仅用于元数据描述,不放入敏感请求或凭证。

openlineage.integration.job.tags

设置血缘任务标签。
  • 类型LIST
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:逗号分隔的键值对;避免过量标签及敏感信息。

openlineage.integration.job.owners

设置血缘任务所有者元数据。
  • 类型LIST
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:逗号分隔的键值对,按组织元数据策略填写。

最佳实践

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

适用业务场景:首次连接已有业务数据的 PostgreSQL,为搜索、缓存或数据同步建立完整初始状态,并继续接收业务变化,而不是只接收启动后的新增数据。 配置示例:沿用快速开始的连接、独立复制槽、预建 publication 和表过滤配置;首次创建 Connector 前确定完整采集范围,无需添加其他配置。
关键说明:启动前先核对表筛选与 publication 集合,再安排快照读取窗口,快照期间避免 DDL。默认 initial 在首次无 Offset 时输出 r 记录,完成后继续输出增量;支持衔接快照期间的变化,但初始化读取会增加数据库负载。不要把“扩大表筛选范围”视为新表一定会自动重新初始化:已有 Offset 时 initial 不会仅因加表重新扫描存量。未完成快照重启可能再次发送已读行,下游需按稳定消息键处理重复。日常停启保留 Connector 身份、主题前缀、Offset 和复制槽,并确保所需 WAL 仍可用,不要为例行维护删除槽。

维护捕获范围时统一管理 publication 表集合

适用业务场景:采集需求从初始表集合扩展到其他业务表,希望表筛选与 PostgreSQL publication 一致,避免只改 Connector 配置却收不到新表增量。适用于由采集团队获授权管理独占按表 publication 的部署。 配置示例:在快速开始上覆盖以下两项,表正则表示维护后的完整表集合,publication 名称仍指向该 Connector 独占的按表 publication;若原先设置了 table.exclude.list,先移除它。
关键说明:先确认账号拥有管理 publication 及加入各表所需的权限,再更新配置并重新启动采集;初始化时 filtered 会按筛选结果创建或替换 publication 的表集合,不是向原集合单纯追加。不能用它改写 FOR ALL TABLES publication,也不要用于被其他消费者共享的 publication。新增表若有存量,需另外规划基线补采;更新 publication 只解决复制范围,不会触发已完成的 initial 快照重做。若采用 DBA 管理方式,则保留 disabled 并由 DBA 协调表集合变更,不需要为自动管理授予额外权限。

监控

监控内容

关注 Kafka Connect 集群健康状态、Connector 与 Task 状态、吞吐和处理延迟、Offset 提交进度、错误与重试,以及 Worker JVM 的堆内存、GC 和线程信号;Task 显示运行中时仍需确认记录和 Offset 持续推进,发现长期积压应结合源端运维信息检查复制槽的 WAL 保留压力。只有部署启用了相应错误处理时才关注 DLQ 活动,不应假定所有采集异常都会写入 DLQ。

导入 Grafana 大盘

下载 AutoMQ Connect 集群 Grafana 大盘,准备与大盘匹配的 Prometheus 兼容数据源和集群、Worker 等标签,随后在 Grafana 中导入 JSON 并选择相应数据源与集群。

限制条件

  • 一个 Connector 只捕获一个数据库并创建一个流式 Task,提高 tasks.max 不会拆分复制流。
  • 同一层级的 include 与 exclude 配置互斥;Connector 表过滤不能补足 publication 中缺失的表。
  • PostgreSQL 逻辑解码不提供可回放的 SQL DDL 事件;事件 Schema 更新不等于输出完整 DDL 历史。
  • before 的完整程度取决于复制标识和数据库提供的旧元组,不能假定每次更新或删除都有完整旧行。
  • 增量快照执行期间不支持 Schema 变化,应协调 DDL 窗口。
  • 默认跳过 TRUNCATE;启用后输出无键的 t 事件,而非逐行删除或逐行 tombstone。
  • 跨 Kafka Topic 或分区不保证全局消费顺序和跨表原子可见性,尤其不能依赖无键 TRUNCATE 与行事件的全局顺序。
  • 删除复制槽或所需 WAL 被清理后,Kafka Offset 不能重新生成源端历史日志;重新快照只能重建当前状态,不能补回已消失的中间变化。
  • 正常故障恢复和未完成快照重启可能重发记录;不能将默认配置解释为端到端恰好一次同步保证。

常见问题

快照已有数据,后续更新却没有进入 Topic

先检查实际使用的逻辑解码插件,再核对 pgoutput publication 是否包含目标表、是否发布所需操作,以及 schematableskipped.operations 过滤是否排除了事件。快照通过 JDBC 读取,能读到存量不证明 publication 范围正确。DBA 管理方式下由 DBA 补齐 publication;使用 filtered 时确认独占按表 publication 和管理权限,并在更新配置后重新启动。若数据库端更新或删除报错,还需检查表的复制标识。

启动时报 publication 或表权限不足

检查账号是否同时具有复制权限、数据库连接权限、Schema 使用权限和表读取权限。自动创建 publication 还涉及数据库 CREATE,修改 publication 和加入表涉及对象所有权,默认 all_tables 可能尝试创建需要超级用户权限的全表 publication。可由 DBA 预建所需按表 publication 并设置 publication.autocreate.mode=disabled,而不是直接升级采集账号为超级用户。若采用 filtered,确认已有 publication 不是 FOR ALL TABLES 且表筛选非空。

删除后收到两条记录,第二条值为空

这是默认删除语义:先发送 op=d 的删除事件,其 after 为空,再发送同键空值 tombstone。消费者应区分删除 envelope 与 tombstone,不能把空值记录当作解析损坏。需要重建状态的应用应处理删除;日志压缩使用 tombstone 清理同键历史,不会替下游数据库自动删除行。

更新事件中旧值不完整,或大字段出现不可用占位符

检查表的 REPLICA IDENTITYDEFAULT 通常仅提供旧键列,未变化 TOAST 字段也可能不携带实际值。占位符表示不可用,不是 SQL NULL;消费者可在不可用时保留已有字段值。确实需要完整旧行时由 DBA 评估 FULL 及其额外 WAL 和处理成本,不要仅为消除占位符对所有表统一开启。自定义消息键不会补回数据库未发送的旧列值。

增加表筛选后没有收到新表存量

table.include.list 控制捕获范围,不代表请求重新执行初始化。已有 Offset 且首次快照已完成时,默认 initial 不会因加表自动重做存量扫描。先核对 publication 和新表权限,再根据下游基线要求规划受支持的增量快照,或建立独立的初始化采集流程;不要直接清空 Offset、删除槽来尝试修复,以免影响原有恢复边界。

停机后 WAL 占用持续增长

复制槽会保留采集所需 WAL。检查 Connector 与 Task 状态、积压、Offset 提交和 LSN 确认进度,确认没有关闭 flush.lsn.source 却缺少外部回收机制。恢复消费后观察进度,并按允许停机时长规划磁盘与服务端保留上限。过滤大量事件时可评估心跳;低流量数据库可能还需能产生 WAL 的数据库心跳查询,且相关表须在 publication 内。不要仅为释放磁盘删除仍需恢复的槽。

重启后出现重复记录或无法读取原日志位置

重复可能来自尚未持久化的 Offset 或未完成快照重扫,下游应使用稳定键和源端位置等信息设计幂等处理。无法读取原位置时检查是否更改了采集身份、复制槽是否保留,以及所需 WAL 是否已清理;缺失日志无法靠增加重试次数恢复。确需重建时协调新的快照基线、下游已有状态和历史事件缺口,即使选择允许重快照的策略也不能恢复已经丢失的中间变化。