SeaTunnel Greenplum Sink 连接器实战指南:基于 JDBC 的批量写入、Upsert 与 CDC 数据接入
2026/9/19 9:24:56 网站建设 项目流程

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 引擎

  1. 使用org.postgresql.Driver:确保 PostgreSQL JDBC 驱动 已放置到每个执行节点的${SEATUNNEL_HOME}/plugins/Jdbc/lib/
  2. 使用com.pivotal.jdbc.GreenplumDriver:自行下载 Greenplum 原生 JDBC 驱动,同样放置到${SEATUNNEL_HOME}/plugins/Jdbc/lib/

SeaTunnel Zeta 引擎

  1. 使用org.postgresql.Driver:确保 PostgreSQL JDBC 驱动已放置到每个节点的${SEATUNNEL_HOME}/lib/
  2. 使用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 对照)

数据源驱动URLMaven
使用 PostgreSQL 驱动连接 Greenplumorg.postgresql.Driverjdbc:postgresql://localhost:5432/testdb下载
使用 Greenplum 原生驱动连接 Greenplumcom.pivotal.jdbc.GreenplumDriverjdbc:pivotal:greenplum://localhost:5432;DatabaseName=testdb从 Greenplum 官方渠道获取

注意两点细节:

  • urljdbc:pivotal:greenplum:开头时,方言工厂会自动识别为 Greenplum 并复用 PostgreSQL 方言;而jdbc:postgresql://打头的 URL 则直接走 PostgreSQL 方言工厂,二者最终使用的方言实现相同。
  • PostgreSQL 驱动与 Greenplum 服务器版本需要兼容,建议以 Greenplum 官方兼容矩阵为准选择驱动版本。

数据类型映射

Greenplum 沿用 PostgreSQL JDBC 驱动的映射体系,SeaTunnel 侧的类型转换由 PostgresTypeConverter.java 完成(serial/bigserial/money/jsonb/bytea等类型在该转换器中均有显式注册)。常用类型对照如下:

Greenplum 数据类型SeaTunnel 数据类型
BOOLEANBOOLEAN
SMALLINT / INT2SMALLINT
INT / INT4 / SERIALINT
BIGINT / INT8 / BIGSERIALBIGINT
NUMERIC(p, s) / DECIMAL(p, s) / MONEYDECIMAL(p, s)
REAL / FLOAT4FLOAT
DOUBLE PRECISION / FLOAT8DOUBLE
CHAR / VARCHAR / TEXT / JSON / JSONBSTRING
DATEDATE
TIMETIME
TIMESTAMP / TIMESTAMPTZTIMESTAMP
BYTEABYTES

使用建议:

  • 自增主键(SERIAL/BIGSERIAL)映射为INT/BIGINT,在自动建表或 upsert 场景中可直接作为primary_keys使用;
  • NUMERIC(p, s)的精度与标度会随列定义透传,MONEY统一按DECIMAL处理;
  • BYTEA二进制列映射为 SeaTunnelBYTES,写入时以字节数组形式绑定参数。

选项详解

Greenplum 文档中直接列出的是与自身强相关的常用配置;其余 JDBC Sink 参数(batch_sizemax_retriesgenerate_sink_sqldatabasetableprimary_keysconnection_check_timeout_secmax_commit_attempts等)全部继承自 Jdbc Sink。

Greenplum 常用选项

名称类型是否必填默认值描述
urlString-JDBC 连接 URL。PostgreSQL 驱动格式jdbc:postgresql://host:port/database;Greenplum 原生驱动格式jdbc:pivotal:greenplum://host:port;DatabaseName=database
driverString-JDBC 驱动类名,通常为org.postgresql.Drivercom.pivotal.jdbc.GreenplumDriver
usernameString-Greenplum 用户名。
passwordString-Greenplum 密码。
queryString-写入上游数据的参数化 SQL,如insert into sink(age, name) values(?, ?)query优先级高于自动生成的写入 SQL,且不能与generate_sink_sql = true同时使用。
batch_sizeInt1000写入 Greenplum 前最多缓存的记录数;达到该值、checkpoint 准备提交或 writer 关闭时触发 flush。
max_retriesInt0executeBatch失败后的重试次数。
generate_sink_sqlBooleanfalse是否根据databasetable自动生成插入 SQL。
databaseString-generate_sink_sql = true时使用的数据库名。
tableString-generate_sink_sql = true时使用的目标表名;Greenplum 有 schema 概念,推荐写成schema.table(如public.orders),也支持${schema_name}${table_name}占位符。
primary_keysArray-自动生成 SQL 时用于 upsert 语义的主键字段列表。
connection_check_timeout_secInt30验证数据库连接操作的超时时间(秒)。
max_commit_attemptsInt3事务提交失败时的最大重试次数。
transaction_timeout_secInt-1事务超时时间(秒),-1表示无限制。
enable_upsertBooleantrue是否启用基于主键的 upsert 写入;若任务数据无重复键,可设为false提升导入速度。
common-options--Sink 插件通用参数,如plugin_inputparallelism,参考 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 方言支持tablespacefillfactor):自动建表时可追加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(),且不支持MAPARRAYROW类型。

:::tip 两种写入模式,必须二选一

  • 自动生成 SQLgenerate_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的双引号引用、hashModForFieldHASHTEXT分片函数等)。

Upsert(ON CONFLICT)语句生成

当配置了primary_keysenable_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 = trueCOPY 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询