尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel HiveJdbc Source Connector 实战指南:基于 HiveServer2 JDBC 的批量数据读取

SeaTunnel HiveJdbc Source Connector 实战指南:基于 HiveServer2 JDBC 的批量数据读取 SeaTunnel HiveJdbc Source Connector 实战指南基于 HiveServer2 JDBC 的批量数据读取【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文聚焦 Apache SeaTunnel 的 HiveJdbc Source Connector即使用Jdbc插件配合 Hive dialect 读取 Hive 数据的方案它通过 HiveServer2 的标准 JDBC 接口org.apache.hive.jdbc.HiveDriver执行查询并拉取数据与直接读取 HDFS 文件的 Hive Source 形成互补。读完本文你将掌握 HiveJdbc 的驱动部署、参数配置、分区并行读取、类型映射与 Kerberos 认证的完整实战方案并理解其底层实现原理。一、连接器定位为什么需要 HiveJdbcSeaTunnel 生态中访问 Hive 数据存在两条技术路线理解二者的差异是选型的前提路线读取方式适用场景Hive Source直接读取 HDFS 上的表文件通过 metastore 解析SeaTunnel Worker 能直接访问 metastore 与 HDFSHiveJdbc Source本文通过 HiveServer2 的 JDBC 接口提交query拉取结果集Worker 无法直连 metastore/HDFS所有 I/O 委托给 HiveServer2HiveJdbc 的核心思路是把所有数据访问 I/O 全部委托给 HiveServer2SeaTunnel 侧只负责提交 SQL 和消费结果集。这意味着不需要在 SeaTunnel Worker 上配置 HDFS 客户端、metastore 地址或相关依赖读取逻辑完全由 Hive 侧HiveServer2 Hive 引擎决定天然支持 Hive 的 SQL 语义支持查询语句query可通过 SQL 实现列投影效果支持 Kerberos 认证适配企业级安全集群。二、支持范围与特性矩阵支持的 Hive 版本官方确认支持3.1.3 与 3.1.2其他版本需要自行测试验证。超时参数的支持边界socket_timeout_ms与connect_timeout_ms两个参数仅在Hive 3.2.0上完成过验证。对于更早的版本含 3.1.x尚未验证——这两个参数会被传递给 JDBC 驱动见下文源码分析但实际是否生效取决于目标 Hive 版本的驱动实现配置前需确认版本兼容性。支持的引擎SparkFlinkSeaTunnel Zeta特性矩阵Key Features特性支持情况batch批式✅ 支持stream流式❌ 不支持exactly-once精确一次❌ 不支持column projection列投影✅ 支持parallelism并行度✅ 支持support user-defined split用户自定义分片✅ 支持其中列投影通过query语句直接实现——只查询需要的字段即可减少传输数据量。三、数据源信息与驱动部署支持的数据源数据源支持的版本驱动URL 示例MavenHive不同依赖版本对应不同驱动类org.apache.hive.jdbc.HiveDriverjdbc:hive2://localhost:10000/defaultorg.apache.hive:hive-jdbc驱动部署Database Dependency根据运行引擎不同Hive JDBC 驱动 jar 的放置位置也不同Spark / Flink将hive-jdbc驱动 jar 放入${SEATUNNEL_HOME}/plugins/jdbc/lib/SeaTunnel Zeta将驱动 jar 放入${SEATUNNEL_HOME}/lib/。驱动由seatunnel.source.Jdbc connector-jdbc插件加载见仓库 plugin-mapping.properties 的映射关系。连接器通过 URL 前缀自动识别 Hive dialectjdbc:hive2://...因此配置中的插件名直接使用Jdbc即可无需单独注册 Hive 专用插件。四、数据类型映射附源码实现依据HiveJdbc 的类型映射在 HiveTypeMapper.java 中实现具体规则如下Hive 数据类型SeaTunnel 数据类型BOOLEANBOOLEANTINYINTSHORT源码实现为 BYTE见下方说明SMALLINTSHORTINT / INTEGERINTBIGINTLONGFLOATFLOATDOUBLE / DOUBLE PRECISIONDOUBLEDECIMAL(x,y) / NUMERIC(x,y)列宽 38DECIMAL(x,y)DECIMAL(x,y) / NUMERIC(x,y)列宽 ≥ 38DECIMAL(38,18)CHAR / VARCHAR / STRINGSTRINGDATEDATETIMESTAMPTIMESTAMPBINARY / ARRAY / INTERVAL / MAP / STRUCT / UNIONTYPE暂不支持源码级说明两处值得注意的细节DECIMAL 处理逻辑在 HiveTypeMapper.java 中DECIMAL/NUMERIC若带精度precision 0则原样映射为DecimalType(precision, scale)若未定义精度与标度则回退为DecimalType(38, 18)并输出 WARN 日志。因此官方文档表格中列宽 ≥ 38 时映射为 DECIMAL(38,18)对应的是后一种兜底路径。TINYINT 的映射差异文档表格将TINYINT归入 SHORT但源码中HIVE_TINYINT实际返回BasicType.BYTE_TYPEHiveTypeMapper.java。以源码实现为准TINYINT 在最新版本中被映射为 SeaTunnel 的 BYTE 类型。不支持类型的行为BINARY、ARRAY、INTERVAL、MAP、STRUCT、UNIONTYPE 等复杂类型会抛出convertToSeaTunnelTypeError异常任务直接失败而非静默丢数据建议在query中通过列投影规避这些字段。五、Source 参数详解HiveJdbc 的参数在源码中定义于 JdbcSourceOptions.java分区相关与 JdbcCommonOptions.java连接与 Kerberos 相关完整清单如下参数名类型是否必填默认值说明urlString是-JDBC 连接地址指向 HiveServer2 端点如jdbc:hive2://localhost:10000/defaultdriverString是-连接远程数据源的 JDBC 驱动类名Hive 固定为org.apache.hive.jdbc.HiveDriverusernameString否-连接实例的用户名passwordString否-连接实例的密码queryString是-查询语句HiveServer2 返回的结果集 schema 决定输出 schemaconnection_check_timeout_secInt否30等待连接校验操作完成的秒数socket_timeout_msInt否86400000从服务端读取数据的 Socket 超时毫秒0表示不超时仅在 Hive 3.2.0 验证过connect_timeout_msInt否86400000建立连接的连接超时毫秒0表示不超时仅在 Hive 3.2.0 验证过partition_columnString否-并行分区的列名仅支持数值型主键列partition_lower_boundBigDecimal否-扫描的partition_column最小值不设置时 SeaTunnel 会向数据库查询最小值partition_upper_boundBigDecimal否-扫描的partition_column最大值不设置时 SeaTunnel 会向数据库查询最大值partition_numInt否任务并行度分区数量仅支持正整数默认等于任务并行度fetch_sizeInt否0每次从数据库拉取的行数减少数据库命中次数以提升大结果集查询性能0表示使用 JDBC 驱动默认值use_kerberosBoolean否false是否启用 Kerberos 认证kerberos_principalString否-use_kerberos true时设置如test_userxxxkerberos_keytab_pathString否-use_kerberos true时设置 keytab 文件路径如/home/test/test_user.keytabkrb5_pathString否/etc/krb5.confuse_kerberos true时设置krb5.conf路径如/seatunnel/krb5.conf或保持默认common-options-否-Source 插件通用参数详见 Source Common Options关键参数底层原理connection_check_timeout_sec默认值 30源码 JdbcCommonOptions.java 中默认 30 秒用于连接有效性校验socket_timeout_ms/connect_timeout_ms默认 24 小时两个超时默认值均为1000 * 60 * 60 * 24JdbcCommonOptions.java。在 HiveJdbcConnectionProvider.java 中这两个值大于 0 时会被写入 JDBC 驱动的连接属性socketTimeout与connectTimeout最终由 Hive JDBC 驱动消费——这也是前文是否生效取决于 Hive 版本的由来partition_lower_bound/partition_upper_bound不设置时自动查询SeaTunnel 会主动向数据库查询列的最小/最大值来界定扫描范围。六、分区并行读取原理源码级HiveJdbc 的并行读取依赖partition_columnpartition_upper/lower_boundpartition_num三个参数。从源码结构看其实现位于 FixedChunkSplitter.java连接器将[partition_lower_bound, partition_upper_bound]数值区间按partition_num均匀切分为多个分片并通过JdbcNumericBetweenParametersProvider为每个分片生成带参数化WHERE column BETWEEN ? AND ?条件的查询FixedChunkSplitter.java每个分片对应一个JdbcSourceSplit。分片在 JdbcSourceSplitEnumerator 中生成并分发由 JdbcSourceReader.java 逐个消费Reader 从双端队列中取出 split打开JdbcInputFormat后循环nextRecord()收集SeaTunnelRow处理完一个 split 再取下一个每处理 50 个 split 打印一次进度日志直到所有分片处理完毕且收到noMoreSplit信号才调用signalNoMoreElement()结束任务。Tips并发与数据倾斜若未设置partition_columnSource 以单并发运行设置后按任务并行度并行。若分区列为BIGINT等大数值类型且数据严重倾斜建议设置parallelism 1以避免数据倾斜带来的长尾问题。七、Kerberos 认证源码级HiveJdbc 通过use_kerberos及相关参数支持 Kerberos 认证。其实现位于 HiveJdbcConnectionProvider.java当jdbcConfig.isUseKerberos()为 true 时走getConnectionWithKerberos()分支HiveJdbcConnectionProvider.java构造 HadoopConfiguration并设置hadoop.security.authentication kerberos调用HadoopLoginFactory.loginWithKerberos(configuration, krb5Path, principal, keytabPath, ...)完成 Kerberos 登录并在登录后的UserGroupInformation上下文中建立 JDBC 连接认证失败时抛出KERBEROS_AUTHENTICATION_FAILED错误码定义于 JdbcConnectorErrorCode.java。因此启用 Kerberos 时URL 中通常还要带上principalhive/_HOSTREALM服务主体参数见下文示例。八、任务配置示例以下四个示例完整继承自官方文档可直接在${SEATUNNEL_HOME}下按bin/seatunnel.sh --config conf方式运行。1. 简单示例单并发查询测试库type_bin表的前 16 行数据、输出全部字段到控制台也可以通过修改query指定要查询的字段实现列投影。# Defining the runtime environment env { parallelism 2 job.mode BATCH } source { Jdbc { url jdbc:hive2://localhost:10000/default driver org.apache.hive.jdbc.HiveDriver connection_check_timeout_sec 100 query select * from type_bin limit 16 } } transform { # 如需了解 transform 插件的完整配置请参阅官方 SQL transform 文档 } sink { Console {} }2. 并行示例全表分片读取按配置的切片字段并行读取整张表适用于全量读取场景source { Jdbc { url jdbc:hive2://localhost:10000/default driver org.apache.hive.jdbc.HiveDriver connection_check_timeout_sec 100 # Define query logic as required query select * from type_bin # 并行分片读取字段 partition_column id # 分片数量 partition_num 10 } }3. 并行边界示例显式上下界当数据值高度聚集时显式指定查询的上界与下界更高效source { Jdbc { url jdbc:hive2://localhost:10000/default driver org.apache.hive.jdbc.HiveDriver connection_check_timeout_sec 100 # Define query logic as required query select * from type_bin partition_column id # 读取起始边界 partition_lower_bound 1 # 读取结束边界 partition_upper_bound 500 partition_num 10 } }本例显式声明[1, 500]区间SeaTunnel 无需再向数据库查询 min/max可减少一次元数据开销。4. Kerberos 认证示例source { Jdbc { url jdbc:hive2://hive-server:10000/default;principalhive/_HOSTREALM driver org.apache.hive.jdbc.HiveDriver query select * from type_bin use_kerberos true kerberos_principal test_userREALM kerberos_keytab_path /home/test/test_user.keytab krb5_path /etc/krb5.conf } }九、使用建议与注意事项版本先行生产环境优先选用官方验证过的 Hive 3.1.2 / 3.1.3若使用 3.2.0可放心启用socket_timeout_ms与connect_timeout_ms否则需自行验证驱动行为。分区列选择partition_column仅支持数值型主键列且要求数据分布均匀出现数据倾斜时以parallelism 1兜底。复杂类型规避BINARY、ARRAY、MAP、STRUCT 等类型暂不支持应在query中排除避免任务直接失败。驱动部署位置Spark/Flink 与 Zeta 的驱动放置路径不同plugins/jdbc/lib/vslib/切换引擎时勿遗漏。变更追踪该连接器的历史变更记录可在 connector-jdbc Changelog 中查阅升级前建议核对版本差异。十、总结HiveJdbc Source 为无法直连 HDFS/metastore的 SeaTunnel Worker 提供了通过 HiveServer2 JDBC 接口读取 Hive 数据的标准方案它继承了 Jdbc 连接器成熟的参数体系超时控制、分区并行、Kerberos 认证类型映射逻辑清晰且可在源码 HiveTypeMapper.java 中逐项核对分区读取与认证流程在 FixedChunkSplitter.java 与 HiveJdbcConnectionProvider.java 中均有完整实现。按本文的配置示例与注意事项即可在 Spark、Flink 或 SeaTunnel Zeta 引擎上稳定落地 Hive 批量数据接入任务。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表