ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

SeaTunnel Vertica Sink Connector 实战指南:JDBC 写入、类型映射与 Exactly-Once 配置

SeaTunnel Vertica Sink Connector 实战指南:JDBC 写入、类型映射与 Exactly-Once 配置 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文以仓库中 docs/en/connector-v2/sink/Vertica.md 为核心骨架结合 connector-jdbc 模块下 Vertica 方言源码与 JdbcVerticaIT.java 端到端测试系统讲解 Vertica Sink 的部署依赖、数据源信息、类型映射、全部 Sink 选项与三种任务配置示例并深入解读底层方言实现帮助你直接上手把数据写入 Vertica 数据仓库。一、Vertica Sink Connector 概述Vertica Sink 是 SeaTunnel 基于 JDBC 协议实现的 Vertica 写入连接器文档定位为JDBC Vertica Sink Connector它通过标准 JDBC 接口将上游数据写入 Vertica 数据库。从 docs/en/connector-v2/sink/Vertica.md 可以确认该连接器具备以下核心能力多引擎支持可运行在 Spark、Flink、SeaTunnel Zeta 三种引擎之上Batch / Streaming 双模式既支持有界批式任务也支持无界流式任务并发写入支持多任务并行写入Exactly-once 语义通过 XA 分布式事务保证精确一次写入要求数据库支持 XA 事务。在连接器能力矩阵上connector-v2-features.md 中定义的 Sink 特性里Vertica Sink 已支持exactly-once打勾而cdc基于主键写入 INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE 行类型尚未支持。这一点在实现层面也有呼应Vertica 的方言类 VerticaDialect.java 提供了getUpsertStatement的 MERGE 语句实现但连接器整体仍不标记 CDC 能力。二、环境依赖与 Driver 准备2.1 依赖放置位置按引擎区分使用 Vertica Sink 前必须先准备 Vertica JDBC 驱动 jar可从 Vertica 官方客户端驱动下载页获取放置目录取决于运行时引擎引擎Driver 放置目录Spark / Flink${SEATUNNEL_HOME}/plugins/SeaTunnel Zeta${SEATUNNEL_HOME}/lib/2.2 数据库驱动依赖文档明确说明请将 Maven 坐标对应的驱动 jar 拷贝到$SEATNUNNEL_HOME/plugins/jdbc/lib/工作目录Vertica 示例如下cp vertica-jdbc-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/从仓库 Maven 配置可确认项目实际使用的驱动版本connector-jdbc/pom.xml 中声明了vertica.version12.0.3-0/vertica.version并以com.vertica.jdbc:vertica-jdbc作为测试/依赖坐标JdbcVerticaIT.java 的driverUrl()方法同样指向vertica-jdbc-12.0.3-0.jar。也就是说vertica-jdbc12.0.3-0 是当前仓库测试链路中实际验证过的驱动版本。三、支持的 DataSource 信息Datasource支持的版本DriverURLMavenVertica不同驱动版本对应不同 Driver 类com.vertica.jdbc.Driverjdbc:vertica://localhost:5433/vertica参考官方下载页其中jdbc:vertica://localhost:5433/vertica中的5433是 Vertica 默认 JDBC 端口vertica为默认数据库名。在 JdbcVerticaIT.java 的端到端测试里使用的连接串模板为jdbc:vertica://%s:%s/%s数据库名为VMart、schema 为public、用户为DBADMIN与文档示例相互印证。3.1 方言识别机制在实现层面Vertica 连接器是 JDBC 连接器家族的方言之一。其方言工厂 VerticaDialectFactory.java 通过AutoService(JdbcDialectFactory.class)注册并在acceptsURL方法中判断 URL 前缀Override public boolean acceptsURL(String url) { return url.startsWith(jdbc:vertica:); }也就是说只要 JDBC URL 以jdbc:vertica:开头SeaTunnel 就会自动匹配到 Vertica 方言无需额外指定方言名称。方言名常量定义于 DatabaseIdentifier.javaVERTICA Vertica。四、数据类型映射Vertica Sink 的字段类型映射由 VerticaTypeMapper.java 实现。该映射器读取ResultSetMetaData的列类型名、精度precision与小数位scale按 Vertica 官方数据类型文档源码注释引用 Vertica 12.0.x SQL 数据类型参考进行转换。映射规则汇总如下Vertica 数据类型SeaTunnel 数据类型BIT(1)、INT UNSIGNEDBOOLEANTINYINT、TINYINT UNSIGNED、SMALLINT、SMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINTINT UNSIGNED、INTEGER UNSIGNED、BIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20,0)DECIMAL(x,y)精度 38DECIMAL(x,y)DECIMAL(x,y)精度 38DECIMAL(38,18)DECIMAL UNSIGNEDDECIMAL(精度1, 小数位)FLOAT、FLOAT UNSIGNEDFLOATDOUBLE、DOUBLE UNSIGNEDDOUBLECHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、JSONSTRINGDATEDATETIMETIMEDATETIME、TIMESTAMPTIMESTAMPTINYBLOB、MEDIUMBLOB、BLOB、LONGBLOB、BINARY、VARBINARY、BIT(n)BYTESGEOMETRY、UNKNOWN暂不支持源码中几个值得注意的实现细节BIT 映射为布尔VERTICA_BIT直接返回BasicType.BOOLEAN_TYPE大精度 DECIMAL 防溢出当 DECIMAL 精度大于 38 时源码会记录will probably cause value overflow告警并降级为DECIMAL(38,18)避免超出 SeaTunnel 类型系统的精度上限LONGTEXT 精度裁剪Vertica 的 LONGTEXT 最大精度为 536870911受限于 SeaTunnel 类型系统会裁剪为 2147483647源码同样有对应告警UNSIGNED 类型告警FLOAT UNSIGNED、DOUBLE UNSIGNED在转换时会打印可能溢出的告警不支持类型直接报错GEOMETRY与UNKNOWN及其他未识别类型会抛出CommonError.convertToSeaTunnelTypeError提示无法转换为 SeaTunnel 类型。行数据转换则由 VerticaJdbcRowConverter.java 提供它继承通用的AbstractJdbcRowConverter仅需返回方言名Vertica具体的 Java 类型 ↔ JDBC 类型互转逻辑复用 JDBC 连接器公共实现。五、Sink 选项详解以下为 Vertica Sink 支持的全部配置项默认值与说明均以 Vertica.md 为准参数名类型是否必填默认值说明urlString是-JDBC 连接 URL例如jdbc:vertica://localhost:5433/verticadriverString是-连接驱动类名Vertica 为com.vertica.jdbc.DriveruserString否-连接用户名passwordString否-连接密码queryString否-自定义写入 SQL如INSERT ...优先级最高databaseString否-配合table自动生成写入 SQL与query互斥且优先级更高tableString否-配合database自动生成写入 SQL与query互斥且优先级更高primary_keysArray否-自动生成 SQL 时支持 insert/delete/update 等操作的主键support_upsert_by_query_primary_key_existBoolean否false数据库不支持 upsert 语法时按主键是否存在选择 INSERT 或 UPDATE SQL 处理更新事件。注意该方法性能较低connection_check_timeout_secInt否30连接校验操作等待超时时间秒max_retriesInt否0提交失败executeBatch的重试次数batch_sizeInt否1000批量写入时缓冲记录数达到batch_size或时间到达checkpoint.interval即刷入数据库is_exactly_onceBoolean否false是否开启 exactly-once 语义使用 XA 事务开启后需配置xa_data_source_class_namegenerate_sink_sqlBoolean否false根据目标表自动生成 SQL 语句xa_data_source_class_nameString否-数据库驱动的 XA 数据源类名max_commit_attemptsInt否3事务提交失败重试次数transaction_timeout_secInt否-1事务开启后的超时时间默认 -1 表示永不超时注意设置超时可能影响 exactly-once 语义auto_commitBoolean否true默认开启自动事务提交propertiesMap否-附加连接参数与 URL 中同名参数冲突时优先级由驱动具体实现决定common-options-否-Sink 插件通用参数详见 Sink Common Optionsenable_upsertBoolean否true根据主键是否存在启用 upsert若任务无主键重复数据设为false可加速导入5.1 关键选项的底层原理Exactly-onceXA 事务开启is_exactly_oncetrue后写入流程切换为两阶段提交Two-phase Commit通过 XA 事务保证数据精确写入一次。这与 connector-v2-features.md 中目标支持 XA 事务时可用两阶段提交保证 exactly-once的说明一致。前提是目标数据库支持 XA 事务Vertica 驱动需提供对应的 XA DataSource 类。Upsert / MERGEVertica 方言在 VerticaDialect.java 中实现了getUpsertStatement自动生成基于MERGE INTO的 UPSERT SQL以USING (SELECT ... FROM DUAL) SOURCE构造数据源按唯一键uniqueKeyFields匹配ON (TARGET.colSOURCE.col)命中则UPDATE SET非键列未命中则INSERT。同时quoteIdentifier使用双引号包裹标识符符合 Vertica 的 SQL 语法习惯。因此配置primary_keysenable_upserttrue即可实现存在则更新、不存在则插入的写入语义。批量写入与重试batch_size默认 1000与checkpoint.interval共同决定刷盘时机max_retries控制executeBatch失败后的重试次数。5.2 Tips并发与分区文档特别提示若未设置分区列任务将以单并发运行若设置了分区列则按任务并发度并行执行。实际配置时请根据目标表结构合理规划并行度与分区策略。六、任务配置示例6.1 简单写入示例以下任务通过 FakeSource 自动生成 16 行数据每行含name字符串与age整型两个字段经 Vertica Sink 写入test_table。运行前需在 Vertica 中预先创建test数据库与test_table表若尚未部署 SeaTunnel请先参照 Install SeaTunnel 安装部署再按 Quick Start With SeaTunnel Engine 运行任务# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { # This is a example source plugin **only for test and demonstrate the feature source plugin** FakeSource { parallelism 1 result_table_name fake row.num 16 schema { fields { name string age int } } } } transform { # 如需了解 transform 插件配置请参阅 transform-v2 相关文档 } sink { jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver user root password 123456 query insert into test_table(name,age) values(?,?) } }上述任务运行后test_table中同样会有 16 行数据。这种写法在端到端测试中也有对应体现jdbc_vertica_source_and_sink.conf 使用INSERT INTO e2e_table_sink (id, name, age) VALUES (?, ?, ?);将 100 行测试数据写入 Vertica。6.2 自动生成 Sink SQL无需手写复杂 SQL只需配置数据库名与表名即可自动生成写入语句sink { jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver user root password 123456 # Automatically generate sql statements based on database table names generate_sink_sql true database test table test_table } }当开启generate_sink_sql后连接器会根据database与table推断目标表结构并自动构造 INSERT SQL配合primary_keys时则生成 upsert 语义 SQL。注意database/table与query互斥且自动生成方式优先级更高。6.3 Exactly-once 写入示例对于要求精确一次写入的场景开启is_exactly_once并配置 XA 数据源类名sink { jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.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 com.vertical.cj.jdbc.VerticalXADataSource } }需要说明的是文档中给出的 XA 数据源类名com.vertical.cj.jdbc.VerticalXADataSource以及描述文字中的com.vertical.cj.jdbc.VerticalXADataSource疑似存在笔误Vertica 的 XA 类名应以官方驱动文档为准同时在 Jdbc.md 附录的参考表中Vertica 行的xa_data_source_class_name标记为/即未提供参考值因此实际部署前请务必核对所用 vertica-jdbc 版本驱动包内真实的 XA DataSource 类名。此外开启 exactly-once 时建议将max_retries设为 0如示例所示避免与事务提交重试机制max_commit_attempts相互干扰。开启 XA 事务的通用注意事项详见 Jdbc.md 的 Tips 小节部分数据库需要额外配置才能支持 XA例如 PostgreSQL 需设置max_prepared_transactions 1MySQL 需 8.0.29 且授权XA_RECOVER_ADMINVertica 场景请以官方 XA 支持文档为准。七、验证与排错7.1 端到端测试验证仓库在 JdbcVerticaIT.java 中提供了完整的 Vertica 读写验证用例可作为配置与行为参考使用vertica/vertica-ce:latestDocker 镜像启动容器端口映射5433数据库VMart、schemapublic、用户DBADMIN建表 SQL 为create table if not exists ... (id int, name varchar, age int)并插入 100 行(id, name, age)测试数据通过 jdbc_vertica_source_and_sink.conf 执行Jdbc 源 → Jdbc Sink的整链路任务Sink 侧配置了connection_check_timeout_sec 1000。这意味着你可以用相同思路在本地搭建 Vertica 容器快速验证自己的 Sink 配置。7.2 常见问题定位建议驱动找不到检查驱动 jar 是否已放入对应引擎要求的${SEATUNNEL_HOME}/plugins/Spark/Flink或${SEATUNNEL_HOME}/lib/Zeta以及$SEATNUNNEL_HOME/plugins/jdbc/lib/URL 匹配失败确认 URL 以jdbc:vertica:开头否则方言工厂无法识别类型转换报错目标表包含GEOMETRY、UNKNOWN等不支持类型时会直接抛错可在上游先做类型改写Exactly-once 不生效确认is_exactly_oncetrue且xa_data_source_class_name填写了驱动包内真实存在的 XA 类名同时数据库支持 XA 事务。八、参考资料本文主文档docs/en/connector-v2/sink/Vertica.md通用 JDBC Sink 文档选项附录、XA 注意事项docs/en/connector-v2/sink/Jdbc.md连接器能力定义docs/en/concept/connector-v2-features.md方言工厂VerticaDialectFactory.java方言实现MERGE UPSERT、标识符引用VerticaDialect.java类型映射VerticaTypeMapper.java行转换器VerticaJdbcRowConverter.java端到端测试JdbcVerticaIT.java 与其任务配置 jdbc_vertica_source_and_sink.conf赞分享数据工程大数据批处理流处理【免费下载链接】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 Vertica JDBC Sink Connector 实战指南从依赖配置到 MERGE Upsert 写入SeaTunnel Vertica JDBC Sink Connector 实战指南从依赖配置到 MERGE Upsert 写入 Vertica 作为一款面向数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入 本文围绕 Kingbase Sink 连接器文档 https数据集成ETL大数据批处理流处理变更数据捕获上一篇Spack Build Caches 完全指南创建、分发与使用预编译二进制包加速安装下一篇MPAndroidChart版本历史重要更新与特性演进创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表