SeaTunnel JDBC Vertica Source Connector 完全指南:并行分区读取、类型映射与配置实战
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本篇指南聚焦 Apache SeaTunnel 的 JDBC Vertica Source ConnectorJdbc连接器 Vertica 方言讲解如何通过 JDBC 从 Vertica 分析型数据库读取数据、如何利用partition_column实现并行分片扫描、以及 Vertica 数据类型到 SeaTunnel 类型系统的映射规则。读完本文你将能够独立完成 Vertica 数据源驱动的安装、Source 配置编写单并发/并行/上下界并行三种模式并理解分区并行的底层实现原理。一、连接器概述Vertica Source Connector 是 SeaTunnel 基于通用 JDBC 连接器实现的数据库方言之一通过 JDBC 驱动读取外部数据源数据。它以Jdbc插件名注册在运行时由 VerticaDialectFactory 根据 URL 前缀jdbc:vertica:自动识别并创建 VerticaDialect 方言实例无需额外注册。支持的引擎Spark Flink SeaTunnel Zeta三种引擎均可运行该 Source但驱动加载路径不同见下文。连接器特性Key Features特性支持情况批处理batch✅ 支持流处理stream❌ 不支持精确一次exactly-once✅ 支持列投影column projection✅ 支持并行度parallelism✅ 支持用户自定义分片support user-defined split✅ 支持支持查询 SQL 并可实现投影效果你可以在query中只 select 需要的字段配合列投影特性实现字段裁剪。关于 batch / stream / exactly-once 等特性的详细定义可参考 连接器 v2 特性说明。二、依赖准备JDBC 驱动安装Vertica 官方驱动类为com.vertica.jdbc.Driver需要手动获取驱动 jar 包并放置到对应目录Spark / Flink 引擎将 Vertica JDBC 驱动 jar 包放入${SEATUNNEL_HOME}/plugins/目录SeaTunnel Zeta 引擎将 Vertica JDBC 驱动 jar 包放入${SEATUNNEL_HOME}/lib/目录。驱动可从 Vertica 官方客户端驱动下载页面获取。从仓库的端到端测试可以看到SeaTunnel 自测使用的驱动版本为vertica-jdbc-12.0.3-0见 JdbcVerticaIT.java不同 Vertica 版本对应的驱动类可能不同请以实际环境为准。三、支持的数据源信息Datasource支持的版本DriverUrlMavenVertica不同依赖版本对应不同驱动类com.vertica.jdbc.Driverjdbc:vertica://localhost:5433/vertica见官方客户端驱动下载页Vertica 默认端口为5433URL 中可携带数据库名如/vertica。从源码看方言工厂通过url.startsWith(jdbc:vertica:)判断 URL 是否属于 Vertica见 VerticaDialectFactory.java。四、数据类型映射Vertica 与 SeaTunnel 数据类型映射关系如下表Vertica 数据类型SeaTunnel 数据类型BITBOOLEANTINYINT / 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.java 中实现几个值得注意的行为DECIMAL 精度保护当DECIMAL的精度precision大于 38 时映射为DECIMAL(38,18)并输出will probably cause value overflow告警日志DECIMAL UNSIGNED则映射为DECIMAL(precision 1, scale)为无符号数多保留一位整数位FLOAT/DOUBLE UNSIGNED 溢出告警映射到 FLOAT/DOUBLE 的同时打印溢出风险告警LONGTEXT 精度截断告警Vertica 中 LONGTEXT 最大精度为 536870911由于 SeaTunnel 类型系统限制精度会被置为 2147483647源码会打印对应告警VerticaTypeMapper.java不支持的类型直接报错GEOMETRY、UNKNOWN等类型会抛出convertToSeaTunnelTypeError转换异常VerticaTypeMapper.java任务会失败需要你在query中显式转换这些列如ST_AsText(geom)。五、Source 参数详解名称类型必填默认值说明urlString是-JDBC 连接 URL示例jdbc:vertica://localhost:5433/verticadriverString是-连接远端数据源使用的 JDBC 类名Vertica 固定为com.vertica.jdbc.DriveruserString否-连接实例的用户名passwordString否-连接实例的密码queryString是-查询语句connection_check_timeout_secInt否30用于校验连接的数据库操作的超时时间秒partition_columnString否-并行分区的列名仅支持数值类型主键列且只能配置一列partition_lower_boundBigDecimal否-扫描的partition_column最小值未设置时 SeaTunnel 会查询数据库获取 min 值partition_upper_boundBigDecimal否-扫描的partition_column最大值未设置时 SeaTunnel 会查询数据库获取 max 值partition_numInt否job parallelism分区数量仅支持正整数默认值为作业并行度fetch_sizeInt否0对返回大量对象的查询可通过配置行抓取大小减少数据库访问次数以提升性能0 表示使用 JDBC 默认值propertiesMap否-附加连接配置参数当 properties 与 URL 存在同名参数时优先级由驱动的具体实现决定例如 MySQL 中 properties 优先于 URLcommon-options-否-Source 插件公共参数详见 Source Common Options参数补充说明query必填项可以是全表查询select * from type_bin也可以是带投影的查询select id, name from type_bin。建议在 SQL 中把数据类型映射表中暂不支持的列做显式转换避免任务因类型转换失败而中断fetch_size默认 0 即使用 JDBC 驱动默认行为对 Vertica 这类列式数据库适当调大 fetch size 可减少网络往返批量读取大表时能显著改善吞吐partition_column仅支持数值类型列源码中分区器基于数值范围切分且只能配置一列文档描述其应为数值类型主键列实际使用中该列应具备可比较的数值语义如自增 idproperties用于传递驱动附加参数例如连接超时、SSL 等 Vertica 驱动支持的连接属性具体优先级以 Vertica 驱动实现为准。六、并行分区原理单并发 vs 并行Tips如果不设置partition_column任务将以单并发运行设置了partition_column则按任务并发度并行执行。从源码结构看JDBC Source 的分区能力由 JdbcSource.java、JdbcSourceSplitEnumerator.java 配合 FixedChunkSplitter.java 与 DynamicChunkSplitter.java 实现其工作方式可以概括为读取partition_column、partition_lower_bound、partition_upper_bound、partition_num配置若未显式设置上下界连接器会向数据库执行查询获取该列的实际 min / max 值将[lower_bound, upper_bound]数值区间按partition_num均匀切分为若干分片chunk每个分片生成一个带边界条件的查询分片分片由枚举器分发给各并行 reader 并行执行。要点显式指定partition_lower_bound/partition_upper_bound可以把扫描范围限定在目标数据区间内比全表扫描更高效也能避免对全表执行 min/max 聚合查询的开销partition_num决定切分份数未配置时默认取作业并行度未设置partition_column时退化为单分片、单并发执行。七、任务配置示例以下示例完整继承官方文档的三个典型场景均可直接复制修改后运行。7.1 简单模式单并发查询查询type_bin表前 16 条数据单并发读取全部字段并输出到控制台你也可以在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 user root password 123456 query select * from type_bin limit 16 } } transform { # 如需了解更多 transform 插件配置可参考项目 SQL transform 文档 } sink { Console {} }注意该示例没有配置partition_column因此即使env.parallelism 2Source 仍按单并发读取控制台 Sink 属于Console插件用于将结果打印到标准输出。7.2 并行模式按分区列分片读取通过partition_columnpartition_num将整表数据按分片并行读取source { Jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver connection_check_timeout_sec 100 user 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 } }该模式适合全表读取未指定上下界时SeaTunnel 会自动查询id列的 min/max 值并将[min, max]区间均匀切成 10 份分片并行扫描。7.3 并行边界模式限定上下界读取在并行基础上显式指定扫描边界只读取id在[1, 500]范围内的数据source { Jdbc { url jdbc:vertica://localhost:5433/vertica driver com.vertica.jdbc.Driver connection_check_timeout_sec 100 user 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 } }显式设置上下界可以避免对全表执行 min/max 聚合扫描范围更精确、效率更高上下界类型为 BigDecimal支持小数与超大数值。八、端到端测试验证仓库提供了 Vertica 的端到端测试用例可作为配置正确性的参考测试类JdbcVerticaIT.java 基于 Testcontainers 启动vertica/vertica-ce:latest社区版容器连接信息端口5433、数据库VMart、schemapublic、用户DBADMIN、驱动类com.vertica.jdbc.Driver验证流程在 Vertica 中创建e2e_table_source表并写入 100 条测试数据通过 Source 读取再经 Sink 写入e2e_table_sink表实际配置jdbc_vertica_source_and_sink.conf 展示了JdbcSource 与 Sink 配对使用的最小完整配置。如果你的环境中已有 Vertica 实例可以参照该测试配置url中数据库名、端口、用户密码按实际替换快速验证连通性。九、使用注意事项驱动版本匹配不同 Vertica 版本的驱动类可能不同请从官方下载与实例版本匹配的驱动 jar并放到plugins/Spark/Flink或lib/Zeta目录否则启动时会出现ClassNotFoundException: com.vertica.jdbc.Driver分区列约束partition_column仅支持数值类型、只能配置一列配置后按任务并发度并行执行不配置则单并发执行不支持的类型GEOMETRY、UNKNOWN类型目前不支持会直接抛类型转换异常请在query中预先转换DECIMAL 溢出风险精度大于 38 的 DECIMAL 会被截断为DECIMAL(38,18)源码会输出告警需结合实际数据评估精度损失连接校验超时connection_check_timeout_sec默认 30 秒网络波动或 Vertica 负载较高时可适当调大官方示例中使用 100避免连接校验误判失败同步与异步边界该连接器为批处理场景设计不支持流模式请确保job.mode BATCH。十、延伸阅读JDBC SinkVertica 写入端了解 Vertica 作为 Sink 的配置与 MERGE 写入语义VerticaDialect.java 中实现了基于MERGE INTO的 upsert 语句生成Source Common Optionsresult_table_name、parallelism等公共参数连接器 v2 特性说明batch、exactly-once、列投影等特性的语义定义源码参考Vertica 方言实现目录、JDBC Source 实现目录。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel JDBC SQL Server Source Connector 实战指南并行分片读取、类型映射与配置详解SeaTunnel JDBC SQL Server Source Connector 实战指南并行分片读取、类型映射与配置详解 本篇技术指南围绕 Apache数据工程大数据批处理流处理SeaTunnel Vertica 源连接器完全指南JDBC 读取、类型映射与并行分片实战SeaTunnel Vertica 源连接器完全指南JDBC 读取、类型映射与并行分片实战 Vertica 源连接器是 SeaTunnel 基于 JDBC 框数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Vertica Sink Connector 实战指南JDBC 写入、类型映射与 Exactly-Once 配置SeaTunnel Vertica Sink Connector 实战指南JDBC 写入、类型映射与 Exactly Once 配置 本文以仓库中 docs/数据工程大数据批处理流处理上一篇Sora2API POW验证终极指南本地计算与外部服务的完整对比下一篇20分钟打造企业级异常登录防护Fail2Ban自定义过滤器开发指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考