ARTICLE DETAIL

资讯详情

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

Apache SeaTunnel Vertica Source 实战指南:JDBC 连接器中的参数、并行切分与方言实现

Apache SeaTunnel Vertica Source 实战指南:JDBC 连接器中的参数、并行切分与方言实现 Apache SeaTunnel Vertica Source 实战指南JDBC 连接器中的参数、并行切分与方言实现【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 仓库中的 Vertica Source 文档及其在connector-jdbc模块中的源码实现系统讲解如何通过通用Jdbc插件读取 Vertica 数据库包括驱动部署方式、完整参数与默认值、Vertica 到 SeaTunnel 的数据类型映射、按partition_column的并行切分读取以及三个可直接复用的作业配置示例。读完后你可以独立完成一个 Vertica 批量同步任务的配置并理解 SeaTunnel 在底层通过VerticaDialect、VerticaTypeMapper等类支撑该数据源的具体方式。1. Vertica Source 定位通用 JDBC 连接器的一个方言Vertica Source 并不是一个独立的连接器模块而是 SeaTunnel 通用 JDBC 连接器connector-jdbc下针对 Vertica 数据库的方言实现。在作业配置中它使用统一的插件名JdbcSeaTunnel 根据 JDBC URL 自动选择 Vertica 方言。从源码结构看方言选择由 VerticaDialectFactory 完成其acceptsURL方法判断 URL 是否以jdbc:vertica:开头命中后创建 VerticaDialect。该工厂通过AutoService(JdbcDialectFactory.class)注册到 SPI与 MySQL、PostgreSQL 等方言并列因此你只需把url写成 Vertica 的 JDBC 格式无需额外指定dialect。Vertica 方言在连接器中的三个组成部分组件职责VerticaDialect提供标识符引用规则双引号包裹、UpsertMERGE INTOSQL 生成、Collation SQL 等VerticaTypeMapper将ResultSetMetaData中的 Vertica 列类型映射为 SeaTunnel 类型VerticaJdbcRowConverter基于通用AbstractJdbcRowConverter完成行数据与 SeaTunnel 类型系统间的转换Vertica 方言自 2.3.2 版本引入对应 changelog 中的Add vertica connector (#4303)后续在 2.3.12 修复了 Vertica 无法执行 upsert 写入的问题Fixed Vertica data source cannot upsert data (#9607)相关记录见 connector-jdbc 变更日志。功能支持矩阵该数据源支持的引擎与能力如下支持引擎Spark、Flink、SeaTunnel Zetabatch批处理支持stream流式不支持exactly-once支持column projection列投影支持通过查询 SQL 实现投影parallelism并行度支持support user-defined split自定义切分支持驱动部署方式Vertica JDBC 驱动未随发行版内置需要手动放置驱动 jarSpark / Flink 引擎将 Vertica 官方 JDBC 驱动 jar 放到${SEATUNNEL_HOME}/plugins/目录SeaTunnel Zeta 引擎将驱动 jar 放到${SEATUNNEL_HOME}/lib/目录。驱动需从 Vertica 官方客户端驱动下载页获取官网地址请查阅 Vertica 官方文档。不同驱动版本对应不同的 driver classSeaTunnel 官方 e2e 测试统一使用com.vertica.jdbc.Driver。2. 数据源信息与数据类型映射数据源信息DatasourceSupported versionsDriverUrlMavenVertica不同依赖版本对应不同驱动类com.vertica.jdbc.Driverjdbc:vertica://localhost:5433/vertica从 Vertica 官方下载数据类型映射官方文档给出的 Vertica 到 SeaTunnel 类型映射如下Vertica Data TypeSeaTunnel Data TypeBITBOOLEANTINYINT、TINYINT UNSIGNED、SMALLINT、SMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINTINT UNSIGNED、INTEGER UNSIGNED、BIGINTLONGBIGINT 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暂不支持上述映射与 VerticaTypeMapper 源码逐项对应其中有几处源码级细节值得注意DECIMAL 精度保护当列精度precision 38时映射器会记录 warn 日志“will probably cause value overflow”并降级为DECIMAL(38,18)因为 SeaTunnel 的DecimalType精度上限为 38DECIMAL UNSIGNED映射为DECIMAL(precision 1, scale)通过提升一位精度容纳无符号数的取值上限LONGTEXTVertica 中该类型最大精度为 536870911源码会打印 warn 说明受 SeaTunnel 类型系统限制精度按 2147483647 处理最终映射为STRINGGEOMETRY / UNKNOWN映射到这些类型会直接抛出convertToSeaTunnelTypeError转换异常而不是静默降级。因此建表或查询时如果包含空间类型字段需要提前在query中剔除。3. Source Options 完整参数说明以下是文档给出的全部 Source 参数并补充了从 JdbcCommonOptions 与 JdbcSourceOptions 源码确认的默认值与别名信息NameTypeRequiredDefaultDescriptionurlStringYes-JDBC 连接 URL例如jdbc:vertica://localhost:5433/vertica。源码中该选项存在 fallback 键base-url两者等价driverStringYes-连接远端数据源使用的 JDBC 驱动类名Vertica 场景下为com.vertica.jdbc.DriverusernameStringNo-连接实例用户名源码中亦接受别名userpasswordStringNo-连接实例密码queryStringYes-查询语句。通过查询 SQL 即可实现列投影只输出需要的字段connection_check_timeout_secIntNo30校验数据库连接的操作超时时间秒partition_columnStringNo-并行切分列仅支持数值类型只允许配置一个列partition_lower_boundBigDecimalNo-切分列扫描下界不设置时 SeaTunnel 会查询数据库获取 min 值partition_upper_boundBigDecimalNo-切分列扫描上界不设置时 SeaTunnel 会查询数据库获取 max 值partition_numIntNo作业并行度切分片数量仅支持正整数默认取作业并行度fetch_sizeIntNo0大结果集查询时的行抓取大小减少数据库往返次数以提升性能0 表示使用 JDBC 默认值propertiesMapNo-附加连接配置参数。当 properties 与 URL 中同名参数冲突时优先级由具体驱动实现决定如 MySQL 中 properties 优先common-options-No-Source 插件通用参数详见 Source 通用参数提交作业时的参数校验由 JdbcSourceFactory 的optionRule()执行url与driver为必填项query、partition_column、partition_num、fetch_size、properties等均为可选项。此外该工厂还实现了 dry-run 校验SupportSourceDryRunValidation可以通过真实连接读取表元数据推断 Schema便于在正式运行前验证凭证与表是否存在。Tips如果未设置partition_column作业将以单并发执行设置后将按任务并发度并行执行切分读取。4. 并行切分读取Partitioned Read对大表做批量同步时单连接顺序扫描往往成为瓶颈。SeaTunnel JDBC Source 的并行读取机制是以partition_column数值列为切分键在[partition_lower_bound, partition_upper_bound]区间内均匀切出partition_num个分片每个分片由一个子任务并行执行query 分片条件。参数配合要点只设partition_column下上界未给定时SeaTunnel 会额外发起 min/max 查询探测边界省去人工预估显式给出边界当切分列是主键或自增列且范围已知时直接指定partition_lower_bound/partition_upper_bound可跳过边界探测查询减少一次全表 min/max 扫描控制分片数partition_num不设置时默认等于作业并行度显式设置时要保证为正整数分片过多会加剧数据库压力分片过少则并行收益有限。仓库中的 e2e 测试 JdbcVerticaIT 基于 Testcontainers 启动vertica/vertica-ce容器端口 5433、账号DBADMIN通过 jdbc_vertica_source_and_sink.conf 验证了 Vertica 端到端读写链路Source 侧执行select id, name, age from e2e_table_sourceSink 侧以INSERT INTO e2e_table_sink (id, name, age) VALUES (?, ?, ?)落库并逐行比对 100 条测试数据。可以推断Vertica 场景下最小可运行的 Source 配置就是url driver username password query五项。5. 作业配置示例以下三个示例完整继承自官方文档可直接复制到.conf作业文件中使用。5.1 简单查询单并发从 Vertica 的type_bin表中读取 16 行数据到 Console也可以自行修改query指定要查询的字段实现投影输出。# Defining the runtime environment env { parallelism 2 job.mode BATCH } source{ Jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 query select * from type_bin limit 16 } } transform { # 如需了解 transform 插件配置方式与完整插件列表 # 可参考仓库文档目录 docs/en/transforms 下的说明 } sink { Console {} }5.2 并行读取按配置的切分字段并行读取整张表适合全表同步场景source { Jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 # Define query logic as required query select * from type_bin # Parallel sharding reads fields partition_column id # Number of fragments partition_num 10 } }5.3 并行读取 指定边界当你已知切分列的取值范围时显式指定上下界比让 SeaTunnel 探测 min/max 更高效source { Jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 # Define query logic as required query select * from type_bin partition_column id # Read start boundary partition_lower_bound 1 # Read end boundary partition_upper_bound 500 partition_num 10 } }三个示例的共同点connection_check_timeout_sec 100放宽了连接校验超时适用于首次连接较慢如网络抖动、数据库冷启动的场景默认值 30 秒通常已够用。6. 深入源码VerticaDialect 的 Upsert 实现虽然本文聚焦 Source但理解 Sink 侧实现有助于完整把握 Vertica 方言。从 VerticaDialect 源码看标识符引用采用双引号quoteIdentifier(name)生成namegetUpsertStatement基于主键字段名列表的重载返回Optional.empty()而getUpsertStatementByTableSchema则生成完整的MERGE INTO语句。生成的 SQL 形如MERGE INTO database.table TARGET USING (CAST(:col1 AS sourceType) AS col1, ...) SOURCE ON (TARGET.pkSOURCE.pk) WHEN MATCHED THEN UPDATE SET col2SOURCE.col2, ... WHEN NOT MATCHED THEN INSERT (col1, col2, ...) VALUES (SOURCE.col1, SOURCE.col2, ...)值得注意的是USING子句中对每个绑定参数显式CAST(:field AS sourceType)——源码注释说明这是 Vertica JDBC 驱动的要求必须显式指定数据类型。这段实现正是 changelog 中 2.3.12 版本Fixed Vertica data source cannot upsert data (#9607)修复的核心内容因此如果你的作业涉及 Vertica 写入尤其是带主键更新场景建议升级到 2.3.12 及以上版本。7. 版本与适用性说明引擎支持Spark、Flink、SeaTunnel Zeta 三引擎均可使用驱动 jar 放置位置按引擎区分plugins/或lib/模式限制仅支持 BATCH 模式job.mode BATCH不支持流式读取切分列限制partition_column仅支持数值类型且只能配置一列如需按字符串列切分属于 JDBC Source 的其他方言能力Vertica 文档未承诺该支持类型限制GEOMETRY / UNKNOWN 类型暂不支持遇到时会在类型映射阶段抛出转换异常而非静默处理。参考路径官方文档Vertica SourceSource 通用参数source-common-options.md连接器变更日志connector-jdbc.md方言实现VerticaDialect.java、VerticaTypeMapper.java、VerticaDialectFactory.java参数定义JdbcCommonOptions.java、JdbcSourceOptions.javae2e 验证JdbcVerticaIT.java、jdbc_vertica_source_and_sink.conf【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表