SeaTunnel Greenplum Sink 连接器实战指南:基于 JDBC 的批量写入、Upsert 与 CDC 数据接入
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本指南围绕 SeaTunnel(Zeta/Flink/Spark)中通过 Jdbc 连接器向 Greenplum 写入数据的完整方案展开,涵盖驱动准备、两种写入模式、数据类型映射、常用参数调优、CDC 场景接入,并结合当前仓库源码说明其底层实现(Greenplum 方言如何复用 PostgreSQL 方言、upsert 语句如何生成)。读完本文,你将能够独立完成 Greenplum Sink 任务的配置、驱动部署与排障。
概述:Greenplum Sink 的本质
SeaTunnel 并没有为 Greenplum 单独实现一套 Sink 插件,而是通过 Jdbc 连接器 以 JDBC 协议向 Greenplum 写入数据。由于 Greenplum 完全兼容 PostgreSQL 网络协议,绝大多数场景可以直接使用 PostgreSQL JDBC 驱动(org.postgresql.Driver)连接;若需使用 Greenplum 原生驱动(com.pivotal.jdbc.GreenplumDriver),则需要用户自行提供驱动 jar(受许可证限制,SeaTunnel 不内置该驱动)。
在源码层面,Greenplum 的接入点位于connector-jdbc模块的方言工厂:
- GreenplumDialectFactory.java:通过
url.startsWith("jdbc:pivotal:greenplum:")识别 Greenplum 原生驱动 URL,并直接return new PostgresDialect(),即 Greenplum 完整复用 PostgreSQL 方言实现; - DatabaseIdentifier.java:在方言注册表中登记
GREENPLUM = "Greenplum"。
这意味着 Greenplum Sink 继承的是 JDBC Sink 的全部能力:批写、流写、多表写入、自动生成 SQL / 自定义 SQL、CDC 事件处理以及按 PostgreSQL 方言生成的INSERT ... ON CONFLICT ... DO UPDATEupsert 语句。
支持引擎
| 引擎 | 支持情况 |
|---|---|
| Spark | ✅ 支持 |
| Flink | ✅ 支持 |
| SeaTunnel Zeta | ✅ 支持 |
三者均通过统一的 Jdbc Sink 实现,差异仅体现在驱动 jar 的放置目录(见下文)。
使用依赖:驱动准备是关键一步
不同引擎加载 JDBC 驱动的目录不同,请先按引擎放好驱动再启动任务,否则会抛出“JDBC 驱动未找到”类错误。
Spark / Flink 引擎
- 使用
org.postgresql.Driver:确保 PostgreSQL JDBC 驱动 已放置到每个执行节点的${SEATUNNEL_HOME}/plugins/Jdbc/lib/; - 使用
com.pivotal.jdbc.GreenplumDriver:自行下载 Greenplum 原生 JDBC 驱动,同样放置到${SEATUNNEL_HOME}/plugins/Jdbc/lib/。
SeaTunnel Zeta 引擎
- 使用
org.postgresql.Driver:确保 PostgreSQL JDBC 驱动已放置到每个节点的${SEATUNNEL_HOME}/lib/; - 使用
com.pivotal.jdbc.GreenplumDriver:自行下载驱动并放到${SEATUNNEL_HOME}/lib/,随后重启受影响的 SeaTunnel 进程使驱动进入类路径。
提示:若尚未安装
connector-jdbc插件本身,可在${SEATUNNEL_HOME}下执行sh bin/install-plugin.sh(插件声明见 config/plugin_config 中的connector-jdbc条目)。驱动放置完成后,Zeta 引擎务必重启进程。
主要特性与重要限制
| 特性 | 支持情况 |
|---|---|
| 精确一次(Exactly-Once) | ❌ 不支持 |
| 变更数据捕获(CDC) | ❌ 不支持(指 Sink 端依赖 XA 的精确一次 CDC 语义) |
| 多表写入 | ✅ 支持 |
| 定时刷新 | ❌ 不支持 |
:::tip 为什么 Greenplum Sink 不支持精确一次?
JDBC Sink 的精确一次(is_exactly_once = true)依赖 XA 事务,而Greenplum 不支持 XA 事务,因此无法走 Jdbc.md 中描述的 XA 两阶段提交路径。若业务需要端到端一致性,建议结合外部存储(例如 Kafka、Hudi)使用幂等批写方案兜底。
:::
支持的数据源信息(驱动与 URL 对照)
| 数据源 | 驱动 | URL | Maven |
|---|---|---|---|
| 使用 PostgreSQL 驱动连接 Greenplum | org.postgresql.Driver | jdbc:postgresql://localhost:5432/testdb | 下载 |
| 使用 Greenplum 原生驱动连接 Greenplum | com.pivotal.jdbc.GreenplumDriver | jdbc:pivotal:greenplum://localhost:5432;DatabaseName=testdb | 从 Greenplum 官方渠道获取 |
注意两点细节:
- 当
url以jdbc:pivotal:greenplum:开头时,方言工厂会自动识别为 Greenplum 并复用 PostgreSQL 方言;而jdbc:postgresql://打头的 URL 则直接走 PostgreSQL 方言工厂,二者最终使用的方言实现相同。 - PostgreSQL 驱动与 Greenplum 服务器版本需要兼容,建议以 Greenplum 官方兼容矩阵为准选择驱动版本。
数据类型映射
Greenplum 沿用 PostgreSQL JDBC 驱动的映射体系,SeaTunnel 侧的类型转换由 PostgresTypeConverter.java 完成(serial/bigserial/money/jsonb/bytea等类型在该转换器中均有显式注册)。常用类型对照如下:
| Greenplum 数据类型 | SeaTunnel 数据类型 |
|---|---|
| BOOLEAN | BOOLEAN |
| SMALLINT / INT2 | SMALLINT |
| INT / INT4 / SERIAL | INT |
| BIGINT / INT8 / BIGSERIAL | BIGINT |
| NUMERIC(p, s) / DECIMAL(p, s) / MONEY | DECIMAL(p, s) |
| REAL / FLOAT4 | FLOAT |
| DOUBLE PRECISION / FLOAT8 | DOUBLE |
| CHAR / VARCHAR / TEXT / JSON / JSONB | STRING |
| DATE | DATE |
| TIME | TIME |
| TIMESTAMP / TIMESTAMPTZ | TIMESTAMP |
| BYTEA | BYTES |
使用建议:
- 自增主键(
SERIAL/BIGSERIAL)映射为INT/BIGINT,在自动建表或 upsert 场景中可直接作为primary_keys使用; NUMERIC(p, s)的精度与标度会随列定义透传,MONEY统一按DECIMAL处理;BYTEA二进制列映射为 SeaTunnelBYTES,写入时以字节数组形式绑定参数。
选项详解
Greenplum 文档中直接列出的是与自身强相关的常用配置;其余 JDBC Sink 参数(batch_size、max_retries、generate_sink_sql、database、table、primary_keys、connection_check_timeout_sec、max_commit_attempts等)全部继承自 Jdbc Sink。
Greenplum 常用选项
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL。PostgreSQL 驱动格式jdbc:postgresql://host:port/database;Greenplum 原生驱动格式jdbc:pivotal:greenplum://host:port;DatabaseName=database。 |
| driver | String | 是 | - | JDBC 驱动类名,通常为org.postgresql.Driver或com.pivotal.jdbc.GreenplumDriver。 |
| username | String | 否 | - | Greenplum 用户名。 |
| password | String | 否 | - | Greenplum 密码。 |
| query | String | 否 | - | 写入上游数据的参数化 SQL,如insert into sink(age, name) values(?, ?)。query优先级高于自动生成的写入 SQL,且不能与generate_sink_sql = true同时使用。 |
| batch_size | Int | 否 | 1000 | 写入 Greenplum 前最多缓存的记录数;达到该值、checkpoint 准备提交或 writer 关闭时触发 flush。 |
| max_retries | Int | 否 | 0 | executeBatch失败后的重试次数。 |
| generate_sink_sql | Boolean | 否 | false | 是否根据database和table自动生成插入 SQL。 |
| database | String | 否 | - | generate_sink_sql = true时使用的数据库名。 |
| table | String | 否 | - | generate_sink_sql = true时使用的目标表名;Greenplum 有 schema 概念,推荐写成schema.table(如public.orders),也支持${schema_name}、${table_name}占位符。 |
| primary_keys | Array | 否 | - | 自动生成 SQL 时用于 upsert 语义的主键字段列表。 |
| connection_check_timeout_sec | Int | 否 | 30 | 验证数据库连接操作的超时时间(秒)。 |
| max_commit_attempts | Int | 否 | 3 | 事务提交失败时的最大重试次数。 |
| transaction_timeout_sec | Int | 否 | -1 | 事务超时时间(秒),-1表示无限制。 |
| enable_upsert | Boolean | 否 | true | 是否启用基于主键的 upsert 写入;若任务数据无重复键,可设为false提升导入速度。 |
| common-options | - | 否 | - | Sink 插件通用参数,如plugin_input、parallelism,参考 Sink 通用选项。 |
从 Jdbc Sink 继承的关键参数
schema_save_mode(默认CREATE_SCHEMA_WHEN_NOT_EXIST):表不存在时自动建表、已存在时跳过;RECREATE_SCHEMA会删表重建;ERROR_WHEN_SCHEMA_NOT_EXIST表不存在直接报错;IGNORE跳过建表逻辑;data_save_mode(默认APPEND_DATA):目标表已有数据时的处理策略,DROP_DATA清空数据、APPEND_DATA追加、ERROR_WHEN_DATA_EXISTS有数据即报错;table_options(PostgreSQL/Greenplum 方言支持tablespace与fillfactor):自动建表时可追加TABLESPACE "..."与WITH (fillfactor=<n>),其中fillfactor取值须为[10, 100]的整数,非法取值会在作业提交阶段被 PostgresDialect.validateTableOptions 直接拒绝;batch_interval_ms(默认 0):按时间间隔触发 flush 的写入侧检查,适用于流式低吞吐场景;is_exactly_once/xa_data_source_class_name:XA 精确一次相关,Greenplum 不支持,请勿启用;use_copy_statement:可走 PostgreSQL 的COPY <table> FROM STDIN批量导入路径,但要求驱动支持getCopyAPI(),且不支持MAP、ARRAY、ROW类型。
:::tip 两种写入模式,必须二选一
- 自动生成 SQL:
generate_sink_sql = true+database(+table),SeaTunnel 根据上游 schema 与 RowKind 自动生成 INSERT / UPSERT / UPDATE / DELETE; - 自定义 SQL:不设置
generate_sink_sql(保持默认false),必须提供query,?参数按上游字段顺序绑定。
两者不能同时配置;自定义 SQL 模式下 SaveMode 相关配置不生效。
:::
任务示例
示例一:FakeSource 写入 Greenplum
env { parallelism = 1 job.mode = "BATCH" } source { FakeSource { row.num = 16 schema = { fields { age = "int" name = "string" } } } } sink { Jdbc { driver = "org.postgresql.Driver" url = "jdbc:postgresql://localhost:5432/testdb" username = "tester" password = "pivotal" query = "insert into sink(age, name) values(?, ?)" } }保存为${SEATUNNEL_HOME}/config/greenplum-sink.conf后,通过以下命令提交(以 Zeta local 模式为例):
cd "${SEATUNNEL_HOME}" ./bin/seatunnel.sh --config ./config/greenplum-sink.conf -m local示例二:Greenplum 读写示例(Jdbc Source → Jdbc Sink)
env { parallelism = 1 job.mode = "BATCH" } source { Jdbc { driver = "org.postgresql.Driver" url = "jdbc:postgresql://localhost:5432/testdb" username = "tester" password = "pivotal" query = "select age, name from source" partition_column = "name" split.string_split_mode = charset_based } } sink { Jdbc { driver = "org.postgresql.Driver" url = "jdbc:postgresql://localhost:5432/testdb" username = "tester" password = "pivotal" query = "insert into sink(age, name) values(?, ?)" } }该示例与仓库集成测试中的 jdbc_greenplum_source_and_sink.conf 完全对应,测试容器使用datagrip/greenplum:6.8镜像、tester/pivotal账号,验证了 100 条(age, name)数据从source表读出再写入sink表的完整链路(见 JdbcGreenplumIT.java)。
示例三:自动生成写入 SQL(推荐)
sink { Jdbc { driver = "org.postgresql.Driver" url = "jdbc:postgresql://localhost:5432/testdb" username = "tester" password = "pivotal" generate_sink_sql = true database = "testdb" table = "sink" } }配合schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"可在目标表不存在时自动建表;再配置primary_keys即可生成 PostgreSQL 原生 upsert(见下文源码剖析)。
示例四:CDC 数据流写入 Greenplum
将 MySQL CDC 的变更流实时写入 Greenplum:
env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 10000 } source { MySQL-CDC { username = "cdc_user" password = "cdc_pass" table-list = ["cdc_test.orders"] base-url = "jdbc:mysql://localhost:3306/cdc_test" startup.mode = "initial" } } sink { Jdbc { driver = "org.postgresql.Driver" url = "jdbc:postgresql://localhost:5432/testdb" username = "tester" password = "pivotal" generate_sink_sql = true database = "testdb" table = "orders" primary_keys = ["id"] enable_upsert = true } }要点说明:
enable_upsert = true时,SeaTunnel 针对 CDC 的INSERT/UPDATE_AFTER行生成INSERT ... ON CONFLICT (id) DO UPDATE SET ...,对DELETE行生成DELETE,实现近似“upsert”的同步效果;- 由于 Greenplum 不支持 XA,该链路为**至少一次(at-least-once)**语义,重复的 upsert 由主键幂等收敛;
- 建议为
sink表显式创建主键/唯一约束,避免 ON CONFLICT 因缺少约束而报错。
源码剖析:Greenplum 如何复用 PostgreSQL 方言
方言识别与工厂
GreenplumDialectFactory.java 是理解整个机制的入口:
@Override public String dialectFactoryName() { return DatabaseIdentifier.GREENPLUM; } @Override public boolean acceptsURL(@NonNull String url) { // Support greenplum native driver: com.pivotal.jdbc.GreenplumDriver return url.startsWith("jdbc:pivotal:greenplum:"); } @Override public JdbcDialect create() { return new PostgresDialect(); }可以看到,无论使用哪条 URL,最终拿到的都是PostgresDialect实例。这意味着 Greenplum Sink 在 SQL 生成、类型转换、JDBC 行为上与 PostgreSQL Sink 完全一致(包括quoteIdentifier的双引号引用、hashModForField的HASHTEXT分片函数等)。
Upsert(ON CONFLICT)语句生成
当配置了primary_keys且enable_upsert = true时,PostgresDialect.getUpsertStatement 会生成 PostgreSQL 原生 upsert:
String conflictAction = updateClause.isEmpty() ? "DO NOTHING" : String.format("DO UPDATE SET %s", updateClause); String upsertSQL = String.format( "%s ON CONFLICT (%s) %s", getInsertIntoStatement(database, tableName, fieldNames), uniqueColumns, conflictAction);对应生成的 SQL 形如:
INSERT INTO "public"."orders" ("id", "customer_name", "amount") VALUES (?, ?, ?) ON CONFLICT ("id") DO UPDATE SET "customer_name" = EXCLUDED."customer_name", "amount" = EXCLUDED."amount"细节:若所有字段都是主键(没有可更新列),DO UPDATE SET子句为空,则自动退化为DO NOTHING。未显式配置primary_keys时,SeaTunnel 会尝试从上游 Catalog 元数据继承主键或第一组唯一键,仍无可用键则退化为普通 INSERT。
集成测试佐证
仓库在 connector-jdbc-e2e-part-5 中提供了 Greenplum 的端到端测试:JdbcGreenplumIT通过 Testcontainers 拉起 Greenplum 6.8 容器,创建age INT, name VARCHAR(255)的source/sink表,注入 100 条测试数据后运行jdbc_greenplum_source_and_sink.conf,验证“Jdbc Source 读 Greenplum → Jdbc Sink 写 Greenplum”全链路可用。该测试同时证明了默认使用org.postgresql.Driver驱动、jdbc:postgresql://URL 连接 Greenplum 是被官方验证过的标准路径。
常见问题排查
Q1:任务报“JDBC 驱动未找到”怎么办?
确认驱动 jar 已按引擎放入对应目录:Spark/Flink 放${SEATUNNEL_HOME}/plugins/Jdbc/lib/,Zeta 放${SEATUNNEL_HOME}/lib/并重启进程。PostgreSQL 驱动常见文件名为postgresql-42.x.x.jar。
Q2:upsert 报ON CONFLICT相关错误?
先确认目标表存在与primary_keys匹配的主键或唯一约束;Greenplum 的ON CONFLICT依赖该约束,否则 PostgreSQL 会报 “there is no unique or exclusion constraint matching the ON CONFLICT specification”。
Q3:Greenplum Sink 能做到精确一次吗?
不能。Greenplum 不支持 XA 事务,is_exactly_once在该场景下不可用。需要端到端一致时,请结合外部存储(例如 Kafka、Hudi)使用幂等批写。
Q4:大批量导入如何提速?
可尝试:调大batch_size(如 5000~10000);在无重复键时设置enable_upsert = false;对 PostgreSQL 兼容场景启用use_copy_statement = true走COPY FROM STDIN;迁移大表时关闭create_index以提升写入速度(迁移完成后手动建索引)。
Q5:自动建表时能否指定表空间/填充因子?
可以,Greenplum 沿用 PostgreSQL 方言的table_options:配置"tablespace" = "pg_default"、"fillfactor" = "70",会在CREATE TABLE时追加TABLESPACE "..."与WITH (fillfactor=70);非法取值会在提交阶段被校验拦截。
变更日志
Greenplum Sink 的功能演进随connector-jdbc插件同步发布,详细变更记录请查阅 connector-jdbc 变更日志。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考