SeaTunnel PostgreSQL JDBC Sink Connector 完全指南:批量写入、Exactly-Once 与 CDC 数据同步实践

发布时间:2026/9/28 3:27:38
SeaTunnel PostgreSQL JDBC Sink Connector 完全指南:批量写入、Exactly-Once 与 CDC 数据同步实践 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载JDBC PostgreSql Sink Connector通过 JDBC 将上游数据写入 PostgreSQL支持批式与流式模式、并发写入并可借助 XA 事务实现精确一次exactly-once语义。本篇技术指南围绕 SeaTunnel 官方文档 docs/en/connector-v2/sink/PostgreSql.md 展开完整覆盖 PostgreSQL Sink 的引擎支持范围、驱动依赖准备、数据类型映射、全部 Sink 参数含义以及手工 SQL、自动生成 SQL、Exactly-Once、CDC 事件、Save Mode 五类实战任务配置同时结合seatunnel-connectors-v2/connector-jdbc模块源码深入讲解 PostgreSQL 方言Dialect、类型转换器TypeConverter与批写入执行器的底层实现帮助读者在掌握配置的同时理解其原理可直接复制示例用于生产环境的数据集成任务。一、支持哪些引擎PostgreSQL JDBC Sink Connector 支持以下三种运行引擎SeaTunnel ZetaSeaTunnel 自研引擎推荐SparkFlink无论使用哪种引擎Sink 插件本身的配置方式是统一的区别主要在于 JDBC 驱动包放置的目录见下一节。二、核心能力概述Connector 通过 JDBC 将上游数据写入 PostgreSQL具备以下能力支持Batch 批式模式与Streaming 流式模式支持并发写入配合并发度与分区列使用支持exactly-once精确一次语义通过XA 事务保证支持CDCChange Data Capture变更数据写入INSERT / UPDATE / DELETE 事件。注意exactly-once 依赖数据库对XA 事务的支持因此只有支持 XA 事务的数据库才能启用精确一次语义。可通过设置is_exactly_oncetrue开启。SeaTunnel 官方特性列表可参考 connector-v2-features.md其中 [x] exactly-once 与 [x] cdc 两项均在此 Connector 中标记为已支持。三、依赖准备驱动 JAR 放置位置3.1 Spark / Flink 引擎需要将 PostgreSQL JDBC 驱动 JARpostgresql-xxx.jar放置到目录${SEATUNNEL_HOME}/plugins/3.2 SeaTunnel Zeta 引擎需要将 PostgreSQL JDBC 驱动 JAR 放置到目录${SEATUNNEL_HOME}/lib/3.3 数据库依赖通用按上述「Maven」列对应的驱动版本下载 JAR并复制到 JDBC Sink 的工作目录$SEATUNNEL_HOME/plugins/jdbc/lib/例如 PostgreSQL 数据源cp postgresql-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/如果需要操作 PostgreSQL 的GEOMETRY 空间数据类型需要同时加入两个 JARpostgresql-xxx.jar postgis-jdbc-xxx.jar并一并放入$SEATUNNEL_HOME/plugins/jdbc/lib/。3.4 支持的驱动与连接信息数据源支持的版本驱动类URLMavenPostgreSQL不同驱动版本对应不同驱动类org.postgresql.Driverjdbc:postgresql://localhost:5432/testDownloadPostgreSQL操作 GEOMETRY 类型需要 PostGIS JDBCorg.postgresql.Driverjdbc:postgresql://localhost:5432/testDownload四、数据类型映射PostgreSQL 数据类型到 SeaTunnel 数据类型的映射关系如下表中ARRAYT表示数组类型PostgreSQL 数据类型SeaTunnel 数据类型BOOLBOOLEAN_BOOLARRAYBOOLEANBYTEABYTES_BYTEAARRAYTINYINTINT2、SMALLSERIAL、INT4、SERIALINT_INT2、_INT4ARRAYINTINT8、BIGSERIALBIGINT_INT8ARRAYBIGINTFLOAT4FLOAT_FLOAT4ARRAYFLOATFLOAT8DOUBLE_FLOAT8ARRAYDOUBLENUMERIC列大小 0DECIMAL(列大小, 小数位数)NUMERIC列大小 0DECIMAL(38, 18)BPCHAR、CHARACTER、VARCHAR、TEXT、GEOMETRY、GEOGRAPHY、JSON、JSONB、UUIDSTRING_BPCHAR、_CHARACTER、_VARCHAR、_TEXTARRAYSTRINGTIMESTAMPTIMESTAMPTIMETIMEDATEDATE其他数据类型暂不支持NOT SUPPORTED YET4.1 源码视角类型映射的实际实现类型映射的核心实现位于 PostgresTypeConverter.java其中convert(BasicTypeDefine)负责 PostgreSQL 类型 → SeaTunnel 类型reconvert(Column)负责反向转换用于建表 DDL 生成。从源码可以看到几个文档之外的细节int2 的实际映射文档表格将INT2/SMALLSERIAL归类为INT而在 PostgresTypeConverter.java 中int2/smallserial实际映射为SMALLINTBasicType.SHORT_TYPEint4/serial映射为INTint8/bigserial映射为BIGINT数组类型_int2/_int4/_int8也按元素类型对应为数组类型建议以实际运行时的类型为准money 类型源码将其映射为DECIMAL(30, 2)PG_MONEY分支并非文档中所说的「其他类型不支持」精度与标度约束NUMERIC最大精度为 1000MAX_PRECISION默认精度 38、默认标度 18DEFAULT_PRECISION/DEFAULT_SCALE超出MAX_PRECISION时会在日志中给出警告并按最大精度收敛时间类型标度截断TIME与TIMESTAMP的标度超过 6 时会被截断为 6并打印The scale of time type is larger than 6的警告日志UUID 特殊处理在 PostgresDialect.java 中查询/分片语句遇到uuid类型字段时会拼接::text转换为文本以便与 SeaTunnel 的STRING类型对齐GEOMETRY / GEOGRAPHY在 PostgresJdbcRowConverter.java 中读取GEOMETRY/GEOGRAPHY列时直接取rs.getObject().toString()转为字符串。类型映射结果由 PostgresTypeMapper.java 在读取ResultSetMetaData列名、原生类型名、可空性、精度、标度后构造BasicTypeDefine并交给PostgresTypeConverter完成转换。五、Sink 参数详解参数名类型是否必填默认值说明urlString是-JDBC 连接 URL示例jdbc:postgresql://localhost:5432/test。若需要插入json或jsonb类型请在 JDBC URL 中追加stringtypeunspecified选项driverString是-连接远程数据源使用的 JDBC 驱动类名PostgreSQL 取值为org.postgresql.DriveruserString否-连接实例的用户名passwordString否-连接实例的密码queryString否-使用该 SQL 将上游数据写入数据库例如INSERT ...query具有更高优先级databaseString否-使用database与table-name自动生成 SQL 并将上游数据写入数据库。该选项与query互斥且优先级更高tableString否-与database配合自动生成 SQL 并写入数据。该选项与query互斥且优先级更高表名支持变量${table_name}、${schema_name}替换规则为${schema_name}替换为目标端传入的 SCHEMA 名${table_name}替换为目标端传入的表名primary_keysArray否-用于支持自动生成 SQL 时的insert、delete、update操作support_upsert_by_query_primary_key_existBoolean否false当数据库不支持 upsert 语法时基于查询主键是否存在来选择使用 INSERT SQL 还是 UPDATE SQL 处理变更事件INSERT、UPDATE_AFTER。注意该方式性能较低connection_check_timeout_secInt否30用于校验连接的数据库操作等待完成的时间秒max_retriesInt否0提交失败executeBatch的重试次数batch_sizeInt否1000批量写入时当缓冲记录数达到batch_size或时间达到checkpoint.interval时数据将被刷入数据库is_exactly_onceBoolean否false是否启用精确一次语义使用 XA 事务。开启后需要设置xa_data_source_class_namegenerate_sink_sqlBoolean否false基于要写入的数据库表自动生成 SQL 语句xa_data_source_class_nameString否-数据库驱动的 XA 数据源类名例如 PostgreSQL 为org.postgresql.xa.PGXADataSource其他数据源请参考附录max_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务开启后的超时时间默认-1永不超时。注意设置超时可能影响 exactly-once 语义auto_commitBoolean否true默认开启自动事务提交field_ideString否-标识从源同步到目标端时字段是否需要进行转换ORIGINAL表示不转换UPPERCASE表示转换为大写LOWERCASE表示转换为小写propertiesMap否-额外的连接配置参数。当properties与 URL 存在相同参数时优先级由驱动的具体实现决定例如在 MySQL 中properties优先于 URLcommon-options-否-Sink 插件通用参数详见 Sink Common Optionsschema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST同步任务开启前对目标端已存在的表结构采取的不同处理方案data_save_modeEnum否APPEND_DATA同步任务开启前对目标端已存在数据采取的不同处理方案custom_sqlString否-当data_save_mode选择CUSTOM_PROCESSING时需要填写custom_sql参数。该参数通常填写一条可执行的 SQL会在同步任务开始前执行enable_upsertBoolean否true基于primary_keys是否已存在来启用 upsert。如果任务没有主键重复数据将该参数设置为false可以加速数据导入以上参数的定义与默认值均可在 JdbcOptions.java 中找到对应实现例如IS_EXACTLY_ONCE默认false、BATCH_SIZE默认1000、MAX_COMMIT_ATTEMPTS默认3、TRANSACTION_TIMEOUT_SEC默认-1并通过 JdbcSinkConfig.java 中的of(ReadonlyConfig)统一解析成 Sink 配置对象。5.1 源码中额外提供的参数除文档表格外JdbcOptions.java 中还定义了几个对 PostgreSQL 写入有实际意义、文档未展开的参数use_copy_statementBoolean默认false开启 PostgreSQL 的COPY批量导入语句底层使用 JDBC 驱动的CopyManager可显著提升大批量数据写入吞吐实现见 CopyManagerBatchStatementExecutor.java将行数据转换为 PostgreSQL CSV 格式后执行COPYBYTES类型会做 Base64 编码is_primary_key_updatedBoolean默认true执行 UPDATE 操作时主键是否被更新support_upsert_by_insert_onlyBoolean默认false是否仅通过 INSERT 实现 upsertfetch_sizeInt默认0查询返回大量对象时设置的批量读取行数0表示使用 JDBC 默认值。5.2 table 参数说明table参数与database配合使用用于自动生成写入 SQL。示例${schema_name}.${table_name}_test—— 动态拼接 schema 与表名dbo.tt_${table_name}_sink—— 固定 schema、动态表名public.sink_table—— 完全静态的表名。5.3 schema_save_mode枚举同步任务开启前对目标端已存在的表结构采取不同处理方案RECREATE_SCHEMA表不存在时创建表已存在时删除并重建CREATE_SCHEMA_WHEN_NOT_EXIST表不存在时创建表已存在时跳过ERROR_WHEN_SCHEMA_NOT_EXIST表不存在时直接报错。5.4 data_save_mode枚举同步任务开启前对目标端已存在数据采取不同处理方案DROP_DATA保留表结构删除已有数据APPEND_DATA保留表结构保留已有数据追加写入CUSTOM_PROCESSING用户自定义处理需配合custom_sqlERROR_WHEN_DATA_EXISTS目标端存在数据时报错。5.5 custom_sql字符串当data_save_mode选择CUSTOM_PROCESSING时通过custom_sql填写一条在同步任务开始前执行的 SQL例如清理历史数据、初始化临时表等。5.6 并发与 Tips若未设置partition_column任务将以单并发运行若设置了partition_column则会根据任务并发度并行执行。六、源码视角PostgreSQL 方言实现原理PostgreSQL 专属的方言与类型处理集中在psql包下理解这些实现有助于排查写入与并发问题。6.1 PostgresDialectSQL 生成与分片PostgresDialect.java 实现了JdbcDialect接口核心能力包括Upsert 语句生成getUpsertStatement使用 PostgreSQL 专有的INSERT ... ON CONFLICT (主键) DO UPDATE SET ...语法将冲突时的更新子句构造成colEXCLUDED.col形式PostgresDialect.java标识符引用quoteIdentifier使用双引号包裹字段名并支持按field_ideORIGINAL/UPPERCASE/LOWERCASE对字段名做大小写转换含点号的名称会逐段引用如schema.table游标式读取creatPreparedStatement将autoCommit关闭并使用TYPE_FORWARD_ONLYCONCUR_READ_ONLY的游标模式默认 fetch size 为 128DEFAULT_POSTGRES_FETCH_SIZE用于避免大结果集一次性载入内存PostgresDialect.java行数估算approximateRowCntStatement优先查询pg_class.reltuples统计信息中的近似行数来估算表行数仅当查询带 WHERE 子句或未配置表路径时才回退到COUNT(*)PostgresDialect.java分片哈希hashModForField使用 PostgreSQL 的HASHTEXT()函数实现按哈希取模的分片并行分片读取/写入。6.2 批量写入执行器批量写入侧支持多种执行器其中与 PostgreSQL 强相关的是CopyManagerBatchStatementExecutor它把SeaTunnelRow逐行转换为 PostgreSQL 的 CSV 格式再通过 JDBCCopyManager的COPY ... FROM STDIN一次性导入非常适合大批量初始化导入场景对应use_copy_statement true。6.3 建表与元数据catalog/psql包下的 PostgresCatalog.java 与PostgresCreateTableSqlBuilder负责目标表的元数据读取与自动建表 SQL 生成是schema_save_mode/generate_sink_sql在目标端落地的底层支撑。七、任务示例7.1 简单示例手工指定 SQL该示例通过 FakeSource 自动生成 16 行数据row.num16每行包含namestring与ageint两个字段最终写入 PostgreSQL 的test_table表。运行前需先在 PostgreSQL 中创建数据库test与表test_table若尚未安装部署 SeaTunnel请先参考 Install SeaTunnel 安装部署再参考 Quick Start With SeaTunnel Engine 运行任务。# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 result_table_name fake row.num 16 schema { fields { name string age int } } } } transform { } sink { jdbc { # if you would use json or jsonb type insert please add jdbc url stringtypeunspecified option url jdbc:postgresql://localhost:5432/test driver org.postgresql.Driver user root password 123456 query insert into test_table(name,age) values(?,?) } }提示query中的?占位符会按上游字段顺序绑定需确保字段顺序与表列顺序一致。7.2 自动生成 Sink SQL无需手写 SQL不需要编写复杂 SQL 语句配置数据库名与表名即可自动生成写入语句sink { Jdbc { # if you would use json or jsonb type insert please add jdbc url stringtypeunspecified option url jdbc:postgresql://localhost:5432/test driver org.postgresql.Driver user root password 123456 generate_sink_sql true database test table public.test_table } }7.3 Exactly-Once精确一次写入对于要求精确写入的场景通过 XA 事务保证精确一次sink { jdbc { # if you would use json or jsonb type insert please add jdbc url stringtypeunspecified option url jdbc:postgresql://localhost:5432/test driver org.postgresql.Driver max_retries 0 user root password 123456 query insert into test_table(name,age) values(?,?) is_exactly_once true xa_data_source_class_name org.postgresql.xa.PGXADataSource } }关键点开启is_exactly_oncetrue后必须指定 PostgreSQL 的 XA 数据源类org.postgresql.xa.PGXADataSource此时建议将max_retries设为 0避免与 XA 提交机制产生重复提交的歧义。XA 相关的两阶段提交、Xid 生成与分组提交逻辑可在internal/xa包XaFacadeImplAutoLoad、XaGroupOpsImpl、SemanticXidGenerator等中进一步查看。7.4 CDCChange Data Capture变更事件CDC 变更数据INSERT / UPDATE / DELETE同样受支持此时需要配置database、table与primary_keys并开启generate_sink_sqlsink { jdbc { # if you would use json or jsonb type insert please add jdbc url stringtypeunspecified option url jdbc:postgresql://localhost:5432/test driver org.postgresql.Driver user root password 123456 generate_sink_sql true # You need to configure both database and table database test table sink_table primary_keys [id,name] field_ide UPPERCASE } }说明database与table必须同时配置primary_keys用于定位 UPDATE/DELETE 事件对应的行field_ide UPPERCASE表示目标端字段统一转为大写若源端字段为小写而目标表为大写可避免标识符大小写不匹配问题底层 upsert 语句由PostgresDialect.getUpsertStatement基于ON CONFLICT语法生成具体见 PostgresDialect.java。7.5 Save Mode保存模式功能在任务启动前按需处理目标端的表结构与存量数据sink { Jdbc { # if you would use json or jsonb type insert please add jdbc url stringtypeunspecified option url jdbc:postgresql://localhost:5432/test driver org.postgresql.Driver user root password 123456 generate_sink_sql true database test table public.test_table schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_modeAPPEND_DATA } }该组合表示目标表不存在时自动建表存在则跳过数据以追加方式写入是日常数据同步最常见的策略。若希望每次任务重建表可将schema_save_mode改为RECREATE_SCHEMA若希望在任务前清空存量数据可将data_save_mode改为DROP_DATA。八、总结与建议依赖先行Spark/Flink 引擎将postgresql-xxx.jar放入${SEATUNNEL_HOME}/plugins/Zeta 引擎放入${SEATUNNEL_HOME}/lib/操作空间类型GEOMETRY需额外加入postgis-jdbc驱动并在 JDBC URL 上按需追加stringtypeunspecifiedjson/jsonb 场景。写 SQL 的方式二选一手写query或配置generate_sink_sql database table自动生成两者互斥query优先级更高。精确一次对强一致性场景开启is_exactly_oncetrue并配置xa_data_source_class_nameorg.postgresql.xa.PGXADataSource。CDC 场景必须配置database、table、primary_keys可按需设置field_ide控制字段大小写。性能优化大批量导入可尝试use_copy_statementtrueCOPY 模式无主键重复数据时建议enable_upsertfalse以加快写入batch_size控制单批刷入行数。如需继续深入可阅读 JdbcOptions.java 中全部参数定义以及 PostgresDialect.java、PostgresTypeConverter.java、CopyManagerBatchStatementExecutor.java 等核心实现Sink 通用参数请参考 Sink Common Options。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel MySQL Sink Connector 实战指南JDBC 批量写入、Exactly-Once 与多表同步SeaTunnel MySQL Sink Connector 实战指南JDBC 批量写入、Exactly Once 与多表同步 SeaTunnel 通过 JD数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel PostgreSQL JDBC Sink 连接器完全指南批量/流式写入、Exactly-Once 与 COPY 高速导入SeaTunnel PostgreSQL JDBC Sink 连接器完全指南批量/流式写入、Exactly Once 与 COPY 高速导入 导读 本文以 A数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel SQL Server JDBC Sink 连接器实战指南批量写入、CDC 事件与 Exactly-Once 语义SeaTunnel SQL Server JDBC Sink 连接器实战指南批量写入、CDC 事件与 Exactly Once 语义 本文以 SqlServe数据集成ETL大数据批处理流处理变更数据捕获上一篇猫抓Cat-Catch如何让浏览器自动帮你发现和下载网页视频资源下一篇3个颠覆性技巧让你零基础掌握浏览器资源嗅探的秘密武器创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考