Flink CDC Postgres Pipeline Connector全库同步配置指南与源码级原理剖析【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文围绕 Flink CDC 的 Postgres Pipeline Connector 展开系统讲解如何配置 Postgres 数据源实现整库快照 增量同步、全部连接与快照调优参数、启动位点选择、元数据列与类型映射规则并结合flink-cdc-connect下的源码印证参数校验逻辑、类型推断实现与元数据列落地过程。读完后你可以独立编写一份可运行的 Postgres 全库同步 YAML 管道并理解每个配置项在源码中的实际作用。一、Connector 定位与端到端示例Postgres Connector 允许从 PostgreSQL 数据库读取快照数据snapshot与增量数据incremental提供端到端的整库数据同步能力。它属于 Pipeline Connector管道连接器运行于 Flink CDC 3.x 的 Pipeline 架构之上Source 产出统一的 Event 流可经过 Transform/Route 处理后写入各类 Sink。一个典型的“Postgres 读、Fluss 写”的管道配置如下type: postgres由工厂类 PostgresDataSourceFactory 的identifier()注册为postgressource: type: postgres name: Postgres Source hostname: 127.0.0.1 port: 5432 username: admin password: pass # 所有表必须属于同一个数据库 tables: adb.\.*.\.* decoding.plugin.name: pgoutput slot.name: pgtest schema-change.enabled: true sink: type: fluss name: Fluss Sink bootstrap.servers: localhost:9123 # Fluss 客户端安全相关属性 properties.client.security.protocol: sasl properties.client.security.sasl.mechanism: PLAIN properties.client.security.sasl.username: developer properties.client.security.sasl.password: developer-pass pipeline: name: Postgres to Fluss Pipeline parallelism: 4 schema.change.behavior: lenient示例中的要点tables使用正则匹配整库所有表schema-change.enabled: true开启 Schema 变更事件推断见后文详解Sink 侧通过properties.client.*透传客户端属性pipeline段控制并行度与 Schema 变更处理行为lenient该示例依赖 Postgres 端已创建好 replication slotslot.name指定的pgtest以及wal_level等逻辑解码前置条件。二、Connector Options 全量参数说明下表完整收录官方文档列出的全部选项相对路径均可在当前仓库中查证OptionRequiredDefaultTypeDescriptionhostnamerequired(none)StringPostgres 数据库服务器的 IP 地址或主机名。portoptional5432IntegerPostgres 数据库服务器的整数端口号。usernamerequired(none)String连接 Postgres 数据库服务器使用的用户名。passwordrequired(none)String连接 Postgres 数据库服务器使用的密码。tablesrequired(none)String要监控的 Postgres 表名支持正则匹配多张表。所有表必须属于同一个数据库。注意点号.是 database、schema、table 名之间的分隔符若要用.匹配任意字符必须用反斜杠转义。例如bdb.user_schema_[0-9].user_table_[0-9]、bdb.schema_\.*.order_\.*slot.namerequired(none)String为流式读取变更而创建的 PostgreSQL 逻辑解码 slot 名称。服务器使用该 slot 向连接器推送事件。Slot 命名须符合 PostgreSQL replication slot 命名规则仅可包含小写字母、数字和下划线。decoding.plugin.nameoptional(none)String服务器端安装的 Postgres 逻辑解码插件名称。支持decoderbufs和pgoutput。tables.excludeoptional(none)String要排除的表名在tables参数之后起排除作用同样支持正则用法与tables相同。server-time-zoneoptional(none)String数据库服务器的会话时区如Asia/Shanghai。它控制 Postgres 中 TIMESTAMP 值如何转换为字符串未设置时回退为系统默认时区ZoneId.systemDefault()。scan.incremental.snapshot.chunk.sizeoptional8096Integer表快照的 chunk 行数快照读取时表被切分为多个 chunk。scan.snapshot.fetch.sizeoptional1024Integer快照读取阶段每次 poll 的最大 fetch size。scan.startup.modeoptionalinitialString启动模式可选initial、latest-offset、committed-offset、snapshot详见“启动读取位点”一节。scan.incremental.close-idle-reader.enabledoptionalfalseBoolean是否在快照阶段结束后关闭空闲 reader。开启该功能要求 Flink 版本 ≥ 1.14配合execution.checkpointing.checkpoints-after-tasks-finish.enabledtrueFlink ≥ 1.15 时该 checkpoint 配置默认为 true无需再显式设置。scan.lsn-commit.checkpoints-num-delayoptional3Integer开始提交 LSN 位点前的 checkpoint 延迟个数。checkpoint 的 LSN 位点按滚动方式提交最早被延迟的 checkpoint 标识先提交。connect.timeoutoptional30sDuration连接器尝试连接 Postgres 服务器的最大等待时间不能小于 250ms。connect.max-retriesoptional3Integer构建 Postgres 连接的最大重试次数。connection.pool.sizeoptional20Integer连接池大小。jdbc.properties.*optional(none)String透传自定义 JDBC URL 属性例如jdbc.properties.useSSL false。heartbeat.intervaloptional30sDuration发送心跳事件的间隔用于追踪最新的可用 WAL 日志位点。debezium.*optional(none)String透传给 Debezium Embedded Engine 的属性用于捕获 Postgres 变更。例如debezium.snapshot.mode never完整列表参见 Debezium 的 Postgres Connector 属性文档。chunk-meta.group.sizeoptional1000Integerchunk 元数据的分组大小元数据超过该大小时会被拆分为多组。metadata.listoptional(none)String需要从 SourceRecord 读取并传递给下游可在 transform 模块中使用的元数据列以,分隔。可用值op_ts、table_name、database_name、schema_name详见“支持的元数据列”一节。scan.incremental.snapshot.unbounded-chunk-first.enabledoptionalfalseBoolean快照读取阶段是否优先分配无界 chunk。可能降低对最大无界 chunk 做快照时 TaskManager OOM 的风险。实验性选项。table-id.include-databaseoptionalfalseBoolean生成的 Table ID 是否包含 database。为 true 时 Table ID 为 (database, schema, table)为 false 时为 (schema, table)。schema-change.enabledoptionalfalseBoolean是否为 Postgres source 开启 Schema 变更推断。开启后连接器会将 pgoutput 的 Relation 消息与缓存 Schema 比对推断加列、删列、重命名列、改列类型等变更事件。要求decoding.plugin.name设置为pgoutput。源码印证与补充说明必填项校验PostgresDataSourceFactory 的requiredOptions()明确返回HOSTNAME、USERNAME、PASSWORD、TABLES、SLOT_NAME五个必填项与文档一致。decoding.plugin.name默认值差异Pipeline 工厂的选项定义PostgresDataSourceOptions中默认值为pgoutput而较老的 DataStream API SourcePostgresSourceOptions中默认值为decoderbufs。由于schema-change.enabled依赖 pgoutput建议像示例那样显式写出该配置。单库约束的实现文档 Note 强调“所有表必须属于同一个数据库”。从 PostgresDataSourceFactory 的getValidateDatabaseName()可以看到其实现将tables按逗号拆分、再按非转义点拆成三段校验每段为合法的 Postgres 库名不超过 63 字符、以字母/下划线/$开头且所有表名解析出的库名必须一致否则抛出The value of option tables ... not all table names have the same database name异常。表选择逻辑工厂通过PostgresSchemaUtils.listTables()先拉取库内全部表再用tables正则构造Selectors做包含匹配若tables.exclude非空则做同样的匹配并从结果中剔除两者任一导致最终表列表为空都会抛出IllegalArgumentExceptioncreateDataSource。数值类参数存在下限校验scan.incremental.snapshot.chunk.size、chunk-meta.group.size、scan.snapshot.fetch.size、connection.pool.size必须大于 1connect.max-retries必须大于 0工厂中的validateIntegerOption()。server-time-zone未设置时工厂会打印告警日志“server-time-zoneis not set, which might cause data inconsistencies for time-related fields”并回退到ZoneId.systemDefault()getServerTimeZone。Schema 变更推断schema-change.enabled对应工厂中的SCHEMA_CHANGE_ENABLED由 PostgresSchemaDataTypeInference 等类配合 PostgresSchemaChangeEventHandler 处理把 DDL 变更转换为 SchemaChangeEvent 下发。参数默认值与文档一致的佐证文件PostgresDataSourceOptions.javachunk size 8096、fetch size 1024、connect timeout 30s、pool size 20、heartbeat 30s、LSN 提交延迟 3 等。三、启动读取位点scan.startup.modescan.startup.mode指定 PostgreSQL CDC 消费端的启动模式可选枚举值initial默认首次启动时对监控的表执行一次初始快照随后继续读取 replication slot 中的变更日志。latest-offset首次启动时不做快照直接从复制流末尾开始读即只消费连接器启动之后发生的变更。committed-offset跳过快照阶段从 replication slot 的confirmed_flush_lsn位点开始读取事件。snapshot仅执行快照阶段快照读取完成后作业退出。源码印证PostgresDataSourceFactory.getStartupOptions() 将字符串解析为StartupOptions.initial() / snapshot() / latest() / committed()传入非法值时抛出ValidationException并列出全部合法取值。LSN 位点相关的偏移管理位于 PostgresOffset 与 OffsetCommitEventEnumerator 通过滚动延迟多个 checkpoint 后提交 LSN以支持流阶段持续回收 WAL 文件对应scan.lsn-commit.checkpoints-num-delay的行为说明。四、Source 可用 MetricsMetrics 有助于理解快照/增量进度Postgres Source 支持以下 Flink 指标Group 为namespace.schema.table即真实的 database.schema.table 名GroupNameTypeDescriptionnamespace.schema.tableisSnapshottingGauge该表是否正在做快照namespace.schema.tableisStreamReadingGauge该表是否正在做流式读取namespace.schema.tablenumTablesSnapshottedGauge已完成快照的表数量namespace.schema.tablenumTablesRemainingGauge尚未完成快照的表数量namespace.schema.tablenumSnapshotSplitsProcessedGauge正在/已处理中的 split 数量namespace.schema.tablenumSnapshotSplitsRemainingGauge尚未处理的 split 数量namespace.schema.tablenumSnapshotSplitsFinishedGauge已处理完成的 split 数量namespace.schema.tablesnapshotStartTimeGauge快照开始时间namespace.schema.tablesnapshotEndTimeGauge快照结束时间这些指标通常用于监控整库快照进度例如结合numTablesRemaining与numSnapshotSplitsRemaining判断全库初始化何时完成。五、支持的元数据列metadata.listPostgreSQL CDC 连接器支持从 source record 中读取元数据列这些列可以在 transform 操作中使用或直接透传给下游 sink。与 Transform 内置表达式的区别op_ts只能通过metadata.list获取提供数据库中变更事件发生时的真实操作时间戳table_name、database_name、schema_name既可走metadata.list也可通过 Transform 表达式__table_name__、__namespace_name__、__schema_name__获取。使用metadata.list的优势是不用写 transform 规则即可把这些值直接传递给下游 sink基础场景下更简洁。配置方式逗号分隔source: type: postgres # ... 其他配置 metadata.list: op_ts,table_name,database_name,schema_name支持的元数据列Metadata ColumnData TypeDescriptionop_tsBIGINT NOT NULL变更事件在数据库中发生的时间戳毫秒级 epoch。快照记录该值为 0。table_nameSTRING NOT NULL变更行所在表名。替代方案Transform 表达式中使用__table_name__。database_nameSTRING NOT NULL变更行所在数据库名。替代方案__namespace_name__。schema_nameSTRING NOT NULL变更行所在 schema 名PostgreSQL 特有。替代方案__schema_name__。使用示例将元数据列加入输出投影source: type: postgres hostname: localhost port: 5432 username: postgres password: postgres tables: mydb.public.orders slot.name: flink_slot metadata.list: op_ts,table_name,schema_name transform: - source-table: mydb.public.orders projection: order_id, customer_id, op_ts, table_name, schema_name description: Include metadata columns in output源码印证工厂的 listReadableMetadata() 将metadata.list按逗号拆分、去空格后与PostgreSQLReadableMetadata枚举逐一匹配出现无法识别的列名会抛出cannot be found in postgresSQL metadata异常每个元数据列由独立类实现OpTsMetadataColumnBIGINT NOT NULL从 record metadata 中解析op_ts为 Long、TableNameMetadataColumn、DatabaseNameMetadataColumn、SchemaNameMetadataColumn并统一通过 PostgresMetadataAccessor 访问底层 record 元数据相关的序列化行为可通过测试 PostgresPipelineRecordEmitterTest 与 PostgresEventDeserializerTest 查看验证。六、数据类型映射Data Type Mapping以下为 PostgreSQL 类型到 Flink CDC 类型的完整映射PostgreSQL 类型CDC 类型BOOLEAN、BIT(1)BOOLEANBIT( 1)BYTESSMALLINT、INT2、SMALLSERIAL、SERIAL2SMALLINTINTEGER、SERIALINTBIGINT、BIGSERIAL、OIDBIGINTREAL、FLOAT4FLOATNUMERICDECIMAL(38, 0)DOUBLE PRECISION、FLOAT8DOUBLECHAR[(M)]、VARCHAR[(M)]、CHARACTER[(M)]、BPCHAR[(M)]、CHARACTER VARYING[(M)]STRINGTIMESTAMPTZ、TIMESTAMP WITH TIME ZONEZonedTimestampTypeINTERVAL [P]BIGINTdebezium.interval.handling.mode为 string 时为 STRINGBYTEABYTES或 STRING当debezium.binary.handling.mode为 base64 / base64-url-safe / hex 时JSON、JSONB、XML、UUID、POINT、LTREE、CITEXT、INET、INT4RANGE、INT8RANGE、NUMRANGE、TSRANGE、DATERANGE、ENUMSTRING注意由于所使用的 Debezium 版本不支持多维数组目前仅支持 PostgreSQL 的一维数组如ARRAY[1,2,3]、int[]。源码印证上述映射由 PostgresTypeUtils.convertFromColumn() 依据 Debezium 的PgOid实现——例如BOOL→BOOLEAN、BIT/VARBIT按 precision1 判定 BOOLEAN 否则BINARY(precision)、TIMESTAMPTZ→TIMESTAMP_LTZ(scale)、POINT/UUID/JSON/JSONB/XML/INET/CIDR/MACADDR(MACADDR8)/各 RANGE 类型→STRING未登记的 OID 会继续经TypeRegistry解析命中 ltree、geometry、geography、citext、hstore、enum 等扩展类型仍不支持的类型会抛出Doesnt support Postgres type ... yet。所有数组*_ARRAY均映射为ARRAY(...)一维数组类型。测试 PostgresTypeUtilsTest 对这些转换提供了逐类型断言。时间类型映射Temporal Types除自带时区信息的 TIMESTAMPTZ 外其余时间类型的映射取决于debezium.time.precision.mode的取值adaptive默认、adaptive_time_microseconds、connect。注意受当前 CDC 限制TIME 类型的精度固定为 3——无论debezium.time.precision.mode设为 adaptive、adaptive_time_microseconds 还是 connectTIME 类型都会转换为 TIME(3)。debezium.time.precision.modeadaptive_time_microseconds默认TIME 精度为 3TIMESTAMP 精度为 6。PostgreSQL 类型CDC 类型DATEDATETIME([P])TIME(3)TIMESTAMP([P])TIMESTAMP([P])debezium.time.precision.modeadaptiveTIME 精度 3TIMESTAMP 精度 6映射表同上。debezium.time.precision.modeconnectTIME 与 TIMESTAMP 精度均为 3TIMESTAMP([P]) → TIMESTAMP(3)DATE → DATE。源码印证PostgresTypeUtils.handleTimeWithTemporalMode() / handleTimestampWithTemporalMode() 对三种模式都返回以列 scale 为精度的TIME(scale)/TIMESTAMP(scale)实际精度由 Debezium 侧在adaptive与connect模式下的 value converter 行为决定与文档描述对应。十进制类型映射Decimal Typesdebezium.decimal.handling.mode决定 DECIMAL/NUMERIC/MONEY 的映射方式precise默认所有 DECIMAL、NUMERIC、MONEY 使用精确的 Decimal 逻辑类型。PostgreSQL 类型CDC 类型NUMERIC[(M[,D])]DECIMAL[(M[,D])]NUMERICDECIMAL(38,0)DECIMAL[(M[,D])]DECIMAL[(M[,D])]DECIMALDECIMAL(38,0)MONEY[(M[,D])]DECIMAL(38,digits)scale 由money.fraction.digits连接属性决定表示小数点移动的位数doubleNUMERIC[(M[,D])]、DECIMAL[(M[,D])]、MONEY[(M[,D])] 全部映射为 DOUBLE。stringNUMERIC[(M[,D])]、DECIMAL[(M[,D])]、MONEY[(M[,D])] 全部映射为 STRING。此外当debezium.decimal.handling.mode为 string 或 double 时PostgreSQL 支持将 NaN非数作为 DECIMAL/NUMERIC 的特殊值存储连接器将其编码为Double.NaN或字符串常量NAN。源码印证handleNumericWithDecimalMode() 中 PRECISE 模式下 precision 超过默认值且不超过最大精度时生成DECIMAL(precision, scale)否则回退为DECIMAL(38, 0)handleMoneyWithDecimalMode()用money.fraction.digits作为 scale 生成DECIMAL(38, digits)与文档一致。HSTORE 类型映射debezium.hstore.handling.mode决定 HSTORE 值的映射设为json默认时以 JSON 值的字符串表示编码设为map时使用 MAP schema。PostgreSQL 类型CDC 类型HSTORESTRINGdebezium.hstore.handling.modestring/json 语义HSTOREMAPdebezium.hstore.handling.modemap源码印证handleHstoreWithHstoreMode() 中 JSON 模式返回STRINGMAP 模式返回MAPSTRING, STRING。网络地址类型映射PostgreSQL 提供存储 IPv4、IPv6 与 MAC 地址的专用类型相比普通文本类型可提供输入校验与专用操作符/函数推荐用于存储网络地址PostgreSQL 类型CDC 类型INETSTRINGCIDRSTRINGMACADDRSTRINGMACADDR8STRING源码印证PgOid.INET_OID / CIDR_OID / MACADDR_OID / MACADDR8_OID在 PostgresTypeUtils 中统一映射为DataTypes.STRING()。PostGIS 空间类型映射PostgreSQL 通过 PostGIS 扩展支持空间数据类型GEOMETRY(POINT, xx)表示笛卡尔坐标系中的点EPSG:xx 定义坐标系统适用于局部平面计算。 GEOGRAPHY(MULTILINESTRING)以经纬度存储多线段基于球面模型适用于全球范围的空间分析。前者用于小范围平面数据后者用于需要计及地球曲率的大范围数据。Postgres 空间数据Flink 中转换后的 Json StringGEOMETRY(POINT, xx){coordinates:[[174.9479, -36.7208]],type:Point,srid:3187}GEOGRAPHY(MULTILINESTRING){coordinates:[[169.1321, -44.7032],[167.8974, -44.6414]],type:MultiLineString,srid:4326}源码印证PostgresSchemaDataTypeInference 明确把 Debezium 的io.debezium.data.geometry.Point / Geography / Geometry逻辑结构推断为STRING即上述 JSON 字符串形态PostgresTypeUtils亦将geometryOid、geographyOid映射为STRING。七、关键源码与测试索引围绕本 Connector 的核心实现与验证文件便于进一步深入参数定义PostgresDataSourceOptions.java工厂与校验逻辑PostgresDataSourceFactory.java数据源主体PostgresDataSource.java事件反序列化与 Schema 变更PostgresEventDeserializer.java、PostgresSchemaChangeEventHandler.java底层 Source分片、枚举、位点提交PostgresChunkSplitter.java、PostgresSourceEnumerator.java集成测试全库同步全类型PostgresFullTypesITCase.java、PostgresPipelineITCase.java八、小结Postgres Pipeline Connector 通过hostname/port/username/password/tables/slot.name五项必填配置即可跑通“快照 增量”的整库同步scan.startup.mode控制启动位点语义scan.incremental.snapshot.*系列控制快照并行切分行为debezium.*与jdbc.properties.*提供透传扩展schema-change.enabled需 pgoutput开启 DDL 变更推断metadata.list让op_ts等元数据直接进入下游。所有参数默认值、校验规则与类型映射均能在flink-cdc-connect模块的源码中逐一对应查证配置时建议结合本文的源码链接对照确认适用前提如单库约束、pgoutput 依赖、Flink 版本要求等。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考