Postgres CDC可用于依次读取PostgreSQL数据库全量快照数据和变更数据,保证不多读一条也不少读一条数据。即使发生故障,也能采用Exactly Once方式处理。本文为您介绍如何使用Postgres CDC连接器。
背景信息
Postgres CDC连接器支持的信息如下。
|
类别 |
详情 |
|
支持类型 |
源表 说明
您可以使用JDBC作为结果表和维表连接器。 |
|
运行模式 |
仅支持流模式 |
|
数据格式 |
暂不适用 |
|
特有监控指标 |
|
|
API种类 |
SQL和数据摄入YAML |
|
是否支持更新或删除结果表数据 |
不涉及 |
特色功能
Postgres CDC连接器接入CDC增量快照框架(实时计算引擎VVR 8.0.6及以上版本)。Postgres CDC读取历史全量数据后,自动切换到WAL变更日志读取,保证不多读也不少读数据。即使发生故障,也能保证Exactly Once语义处理数据。Postgres CDC源表提供了并发读取全量数据,无锁读取和断点续传的能力。
作为源表,功能与优势详情如下:
-
流批一体,支持读取全量和增量数据,无需维护两套流程。
-
支持并发读取全量数据,性能水平扩展。
-
全量读取无缝切换增量读取,自动缩容,节省计算资源。
-
全量阶段读取支持断点续传,更稳定。
-
无锁读取全量数据,不影响线上业务。
前提条件
Postgres CDC连接器通过PostgreSQL数据库的逻辑复制读取CDC变更流数据,支持阿里云RDS PostgreSQL、Amazon RDS PostgreSQL和自建PostgreSQL。
阿里云RDS PostgreSQL、Amazon RDS PostgreSQL或者自建PostgreSQL上相应的配置可能有差异,请您在使用之前详细阅读配置Postgres文档进行相关配置。
完成配置后确保有下列的条件:
-
wal_level参数的值需设置为logical,即在预写式日志WAL(Write-ahead logging)中增加支持逻辑编码所需的信息。
-
订阅表的REPLICA IDENTITY为FULL(发出的插入和更新操作事件包含表中所有列的旧值),以保障该表数据同步的一致性。
说明REPLICA IDENTITY是PostgreSQL特有的表级设置,它决定了逻辑解码插件在发生(INSERT)和更新(UPDATE)事件时,是否包含涉及的表列的旧值。REPLICA IDENTITY取值含义详情请参见REPLICA IDENTITY。
-
需要确保max_wal_senders和max_replication_slots的参数值均大于当前数据库复制槽已使用数与Flink作业所需要的slot数量。
-
确保账户系统权限为SUPERUSER或者同时拥有LOGIN和REPLICATION权限,并且具有订阅表的SELECT权限用于全量数据查询。
-
如果您的 Postgres 表中包含计算生成列(Generated Columns),则需要在创建 slot 时将 publish_generated_columns 参数设置为 stored,以免出现全增量阶段读取到的表结构不一致的问题。
注意事项
仅实时计算引擎8.0.6及以上版本支持Postgres CDC增量快照功能。
Replication Slot说明
Flink PostgreSQL CDC 作业依赖 Replication Slot 来确保 WAL(Write-Ahead Log)不被过早清理,从而保障数据一致性。但若管理不当,可能引发磁盘空间浪费或数据读取延迟等问题。请遵循以下建议:
-
请及时清理不再使用的 Slot
-
Flink 不会自动删除 Replication Slot,即使作业已停止(尤其无状态重启场景),以防止因 WAL 被清除而导致数据丢失。
-
若确认某作业不再启动,请手动删除其关联的 Replication Slot,释放磁盘空间。
重要生命周期管理:将 Replication Slot 视为作业资源的一部分,随作业启停同步管理。
-
-
避免复用旧 Slot
-
新作业应使用新的 Slot Name,而非复用旧 Slot。复用可能导致作业启动后需回溯大量历史 WAL,延迟读取最新数据。
-
PostgreSQL的逻辑复制要求一个 Slot 仅能被一个连接使用,不同作业必须使用不同的 Slot 名称。
重要命名规范:自定义slot.name时,避免使用带数字后缀的名称(如 my_slot_1),以防与临时 Slot 冲突。
-
-
启用增量快照下的Slot行为
-
前提条件:必须启用checkpoint,且Source 表必须声明主键。
-
Slot创建规则:
-
未开启增量快照:仅支持单并发,使用 1 个全局 Slot。
-
开启增量快照:
-
全量阶段:每个 Source 并发子任务会创建一个临时 Slot,命名格式为
${slot.name}_${task_id}。 -
增量阶段:自动回收所有临时 Slot,仅保留 1 个全局 Slot。
-
-
-
最大Slot数量:Source 并发数 + 1(全量阶段)
-
-
资源与性能
-
若 PostgreSQL 的 Slot 数量或磁盘空间受限,应适当降低全量阶段的并发度(减少临时 Slot 数量),但会牺牲全量读取速度。
-
若下游支持幂等写入,可设置:
scan.incremental.snapshot.backfill.skip = true,跳过全量阶段的 Binlog 回溯,加快启动速度。此配置仅提供 At-Least-Once 语义。不适用于含聚合、维表 Join 等状态计算的作业(可能丢失中间状态所需的历史变更)。
-
-
不开启增量快照时,不支持在全表扫描阶段执行Checkpoint。
复用postgres订阅
Postgres CDC 使用 pgoutput 时,会通过 Publication 决定哪些表的变更可以被读取。
如果多个作业共用同一个 Publication,请确认该 Publication 已包含所有作业需要读取的表,否则部分表的增量变更可能无法被同步。
注意事项
debezium.publication.autocreate.mode 默认值为 all_tables。在该模式下:
-
如果 Publication 不存在,连接器会自动创建
FOR ALL TABLES的 Publication。 -
如果 Publication 已存在,连接器只会直接复用,不会检查或修改其中的表范围。
如果配置为 filtered,连接器会按当前作业的表过滤规则创建或更新 Publication。多个作业共用同一个 Publication 时,后启动的作业可能会覆盖表范围,影响其他作业读取。
推荐做法
-
生产环境建议手动创建和维护 Publication,并关闭自动创建/更新
-- 创建一个名为 my_flink_pub 的发布,包含所有表(或指定表,每个作业创一个Publication) CREATE PUBLICATION my_flink_pub FOR TABLE table_a, table_b; -- 或者更简单,包含库里所有表 CREATE PUBLICATION my_flink_pub FOR ALL TABLES;说明不建议订阅全库全表,如果数据库非常大,表非常多,这会造成网络带宽浪费和Flink 端 CPU 消耗。
-
添加Flink配置
-
debezium.publication.name = 'my_flink_pub'(指定Publication名称) -
debezium.publication.autocreate.mode = 'disabled'(禁止Flink启动时尝试创建或修改 Publication)
-
这样可以避免多个作业互相覆盖 Publication 配置。
SQL
语法结构
CREATE TABLE postgrescdc_source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>',
'decoding.plugin.name'= 'pgoutput',
'scan.incremental.snapshot.enabled' = 'true',
-- skip backfill 可以加速读取和减少资源,但是会有数据重复。如果下游幂等,建议开启
'scan.incremental.snapshot.backfill.skip' = 'false',
-- 生产上建议设置为filter或者disabled, 手动管理publication而非flink
'debezium-publication.autocreate.mode' = 'disabled'
-- 多Source时,为不同的Source配置不同的publication。
--'debezium.publication.name' = 'my_flink_pub'
);
WITH参数
|
参数 |
说明 |
数据类型 |
是否必填 |
默认值 |
备注 |
|
connector |
connector类型。 |
STRING |
是 |
无 |
固定值为 |
|
hostname |
Postgres数据库的IP地址或者Hostname。 |
STRING |
是 |
无 |
无。 |
|
username |
Postgres数据库服务的用户名。 |
STRING |
是 |
无 |
无。 |
|
password |
Postgres数据库服务的密码。 |
STRING |
是 |
无 |
无。 |
|
database-name |
数据库名称。 |
STRING |
是 |
无 |
数据库名称。 |
|
schema-name |
Postgres Schema名称。 |
STRING |
是 |
无 |
Schema名称支持正则表达式以读取多个Schema的数据。 |
|
table-name |
Postgres表名。 |
STRING |
是 |
无 |
表名支持正则表达式以读取多个表的数据。 |
|
port |
Postgres数据库服务的端口号。 |
INTEGER |
否 |
5432 |
无。 |
|
decoding.plugin.name |
Postgres Logical Decoding插件名称。 |
STRING |
否 |
decoderbufs |
根据Postgres服务上安装的插件确定。支持的插件列表如下:
|
|
slot.name |
逻辑解码槽的名字。 |
STRING |
8.0.1版本之前为非必填,从8.0.1版本开始为必填 |
8.0.1版本之前默认值为flink,从8.0.1版本开始无默认值 |
建议每个表都设置 |
|
debezium.* |
Debezium属性参数。 |
STRING |
否 |
无 |
更细粒度控制Debezium客户端的行为。例如 |
|
scan.incremental.snapshot.enabled |
是否开启增量快照。 |
BOOLEAN |
否 |
false |
参数取值如下:
|
|
scan.startup.mode |
消费数据时的启动模式。 |
STRING |
否 |
initial |
参数取值如下:
|
|
changelog-mode |
用于编码流更改的变更日志(Changelog)模式。 |
String |
否 |
all |
支持的Changelog模式包括:
|
|
heartbeat.interval.ms |
发送心跳包的时间间隔。 |
Duration |
否 |
30s |
单位为毫秒。 Postgres CDC连接器主动向数据库发送心跳包来保证推进Slot的偏移量。当表变更不频繁时,设置该值可以及时回收WAL日志。 |
|
scan.incremental.snapshot.chunk.key-column |
指定某一列作为快照阶段切分分片的切分列。 |
STRING |
否 |
无 |
默认从主键中选择第一列。 |
|
scan.incremental.close-idle-reader.enabled |
是否在快照结束后关闭空闲的Reader。 |
Boolean |
否 |
false |
该配置生效需要设置 |
|
scan.incremental.snapshot.backfill.skip |
是否跳过全量阶段的日志读取。 |
Boolean |
否 |
false |
参数取值如下:
backfill仅在单个分片(chunk)快照查询期间生效,不覆盖整个全量读取过程。跳过backfill后,分片快照SQL执行时读到该时刻表的最新数据;分片已读完之后该分片上发生的更新,不再在全量阶段合并,会在进入增量阶段后从WAL中读取。例如,chunk5快照期间发生的更新会直接体现在chunk5的最新数据中;若已读到chunk80时chunk5才发生更新,该更新会在增量阶段通过WAL补回。 重要
开启后,分片扫描期间及之后的变更在增量阶段仍会通过WAL下发,可能与快照数据重复,仅提供at-least-once语义。请确认下游支持按主键幂等写入后再开启。 |
类型映射
PostgreSQL和Flink字段类型对应关系如下。
|
PostgreSQL字段类型 |
Flink字段类型 |
|
SMALLINT |
SMALLINT |
|
INT2 |
|
|
SMALLSERIAL |
|
|
SERIAL2 |
|
|
INTEGER |
INT |
|
SERIAL |
|
|
BIGINT |
BIGINT |
|
BIGSERIAL |
|
|
REAL |
FLOAT |
|
FLOAT4 |
|
|
FLOAT8 |
DOUBLE |
|
DOUBLE PRECISION |
|
|
NUMERIC(p, s) |
DECIMAL(p, s) |
|
DECIMAL(p, s) |
|
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
|
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
|
CHAR(n) |
STRING |
|
CHARACTER(n) |
|
|
VARCHAR(n) |
|
|
CHARACTER VARYING(n) |
|
|
TEXT |
|
|
BYTEA |
BYTES |
使用示例
CREATE TABLE source (
id INT NOT NULL,
name STRING,
description STRING,
weight DECIMAL(10,3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = '<host name>',
'port' = '<port>',
'username' = '<user name>',
'password' = '<password>',
'database-name' = '<database name>',
'schema-name' = '<schema name>',
'table-name' = '<table name>'
);
SELECT * FROM source;
数据摄入
自实时计算引擎11.4版本起,PostgreSQL连接器作为数据源可以在数据摄入YAML作业中使用。
语法结构
source:
type: postgres
name: PostgreSQL Source
hostname: localhost
port: 5432
username: pg_username
password: pg_password
tables: db.scm.tbl
slot.name: test_slot
scan.startup.mode: initial
server-time-zone: UTC
connect.timeout: 120s
decoding.plugin.name: decoderbufs
sink:
type: ...
配置项
|
参数 |
说明 |
是否必填 |
数据类型 |
默认值 |
备注 |
|
type |
数据源类型。 |
是 |
STRING |
无 |
固定值为postgres。 |
|
name |
数据源名称。 |
否 |
STRING |
无 |
无。 |
|
hostname |
Postgres数据库服务器域名或IP地址。 |
是 |
STRING |
(none) |
无。 |
|
port |
Postgres数据库服务器暴露的端口。 |
否 |
INTEGER |
5432 |
无。 |
|
username |
Postgres用户名。 |
是 |
STRING |
(none) |
无。 |
|
password |
Postgres密码。 |
是 |
STRING |
(none) |
无。 |
|
tables |
需要捕获的Postgres数据库表名。 支持正则表达式,可以监控多个满足该正则表达式的表。 |
是 |
STRING |
(none) |
重要
目前仅支持捕获同一数据库下的表。 点号 (.) 被视为database、schema和table名的分隔符。如果需要在正则表达式中使用点号 (.) 来匹配任何字符,则必须使用反斜杠转义点号。例如: |
|
slot.name |
PostgreSQL复制槽名称。 |
是 |
STRING |
(none) |
名称必须符合 PostgreSQL 复制槽命名规则,可以包含小写字母、数字和下划线字符。 |
|
decoding.plugin.name |
服务器上安装的Postgres逻辑解码插件的名称。 |
否 |
STRING |
|
可选值包括 |
|
tables.exclude |
要排除的 Postgres 数据库表名,此参数将在 tables 参数之后生效。 |
否 |
STRING |
(none) |
表名也支持正则表达式,可以排除多个满足该正则表达式的表。用法与 tables 参数相同。 |
|
server-time-zone |
数据库服务器的会话时区,如“Asia/Shanghai”。 |
否 |
STRING |
(none) |
如果未设置,则将使用系统默认时区 ( |
|
scan.incremental.snapshot.chunk.size |
增量快照框架中每个chunk的大小(包含的行数)。 |
否 |
INTEGER |
8096 |
当开启增量快照读取时,表会被切分成多个chunk读取。在读完chunk的数据之前,chunk的数据会先缓存在内存中。 每个chunk包含的行数越少,则表中的chunk的总数量越大,尽管这会降低故障恢复的粒度,但可能导致内存OOM和整体的吞吐量降低。因此,您需要进行权衡,并设置合理的chunk大小。 |
|
scan.snapshot.fetch.size |
当读取表的全量数据时,每次最多拉取的记录数。 |
否 |
INTEGER |
1024 |
无。 |
|
scan.startup.mode |
消费数据时的启动模式。 |
否 |
STRING |
initial |
参数取值如下:
|
|
scan.incremental.close-idle-reader.enabled |
是否在快照结束后关闭空闲的 Reader。 |
否 |
BOOLEAN |
false |
该配置生效需要设置execution.checkpointing.checkpoints-after-tasks-finish.enabled为true。 |
|
scan.lsn-commit.checkpoints-num-delay |
在开始提交 LSN 偏移量之前,延迟多少个检查点。 |
否 |
INTEGER |
3 |
检查点 LSN 偏移量将滚动提交,以避免无法从状态恢复。 |
|
connect.timeout |
连接器尝试连接到 Postgres 数据库服务器后,超时前应等待的最长时间。 |
否 |
DURATION |
30s |
此值不能小于 250 毫秒。 |
|
connect.max-retries |
连接器尝试建立 Postgres 数据库服务器连接的最大重试次数。 |
否 |
INTEGER |
3 |
无。 |
|
connection.pool.size |
连接池大小。 |
否 |
INTEGER |
20 |
无。 |
|
jdbc.properties.* |
允许用户传递自定义 JDBC URL 属性。 |
否 |
STRING |
20 |
用户可以传递自定义属性,例如 |
|
heartbeat.interval |
用于追踪最新可用 WAL 日志偏移量的心跳事件发送间隔。 |
否 |
DURATION |
30s |
无。 |
|
debezium.* |
将 Debezium 的属性传递给 Debezium Embedded Engine,后者用于捕获来自 PostgreSQL 服务器的数据更改。 |
否 |
STRING |
(none) |
有关 Debezium PostgreSQL 连接器属性的更多信息,请参阅相关文档。 |
|
chunk-meta.group.size |
chunk元信息的大小。 |
否 |
STRING |
1000 |
如果元信息大于该值,元信息会分为多份传递。 |
|
metadata.list |
传递到下游的可读元数据列表,可在transform模块中使用。 |
否 |
STRING |
false |
使用逗号 (,) 分隔。目前可用的元数据有: |
|
scan.incremental.snapshot.unbounded-chunk-first.enabled |
快照读取阶段是否先分发无界的分片。 |
否 |
STRING |
false |
参数取值如下:
重要
实验性功能。开启后能够降低TaskManager在快照阶段同步最后一个分片时遇到内存溢出 (OOM) 的风险,建议在作业第一次启动前添加。 |
相关文档
-
实时计算Flink版支持的连接器列表,请参见支持的连接器。
-
将数据写入PolarDB PostgreSQL版(Oracle语法兼容1.0)结果表,请参见PolarDB PostgreSQL版(Oracle语法兼容1.0)(退役中)。
-
如果您需要读写RDS MySQL、PolarDB for MySQL或者自建MySQL数据库,请使用MySQL连接器。