免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

SeaTunnel Hive Sink 连接器全解析:表写入、分区自动修复与 S3/OSS 部署实战

SeaTunnel Hive Sink 连接器全解析:表写入、分区自动修复与 S3/OSS 部署实战 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载Hive sink connector导读本文以 SeaTunnel 官方文档 Hive Sink 为核心骨架结合仓库内seatunnel-connectors-v2/connector-hive模块源码系统讲解 Hive Sink 连接器sink的完整能力它如何将上游数据写入 Hive 表、如何通过 2PC 提交实现exactly-once、如何在提交阶段自动向 Hive Metastore 注册分区以及如何在 HDFS、S3、OSS 等不同存储之上落地部署。读完本文你将掌握 Hive Sink 的全部配置项、单表与多表写入两种典型作业配置以及基于 AWS EMR / 阿里云 EMR 的 Hive-on-S3 与 Hive-on-OSS 实战步骤。连接器概述与运行前提Hive Sink 连接器用于将 SeaTunnel 作业中的数据写入 Hive 表。官方文档明确了两点运行前提如果运行在 Spark / Flink 集群上必须确保集群已经集成 Hive文档中测试过的 Hive 版本为 2.3.9如果使用 SeaTunnel EngineZeta需要把seatunnel-hadoop3-3.1.4-uber.jar、hive-exec-3.1.3.jar和libfb303-0.9.3.jar三个 jar 放入$SEATUNNEL_HOME/lib/目录否则运行时会因缺少 Hive Metastore Client 相关类而报错。核心能力Key features能力说明多表写入支持SupportMultiTableSink可配合多表 Source 使用${database_name}.${table_name}动态生成目标表名详见 connector-v2-features精确一次exactly-once默认使用 2PC两阶段提交提交保证精确一次语义详见 connector-v2-features文件格式text、csv、parquet、orc、json压缩编码lzo源码佐证HiveSink类实现了SeaTunnelSinkSeaTunnelRow, FileSinkState, FileCommitInfo, FileAggregatedCommitInfo与SupportMultiTableSink接口见 HiveSink.java。需要说明的是从当前仓库源码看底层HiveTableUtils.parseFileFormat严格按表存储的 InputFormat 类名只接受 TEXT / PARQUET / ORC 三种格式csv、json 类文本格式在实际使用中通常依托 text 格式通过field.delim、line.delim分隔符描述实现见 HiveTableUtils.java。工作原理从配置到提交的完整调用链理解源码有助于定位问题。HiveSink的初始化与执行流程如下解析表信息getTableInformation()通过HiveTableUtils.getTableInfo(readonlyConfig)连接 Hive Metastore 获取目标表的Table对象库名、表名、分区键、存储位置、SerDe 参数等见 HiveTableUtils.java。生成文件 Sink 配置generateFileSinkConfig()从表元数据推导出写入所需的一切参数——列字段自动取自表的普通列 分区列FILE_PATH直接使用tableInformation.getSd().getLocation()即 Hive 表在对象存储 / HDFS 上的数据目录文件名表达式固定为${transactionId}分区字段不写入数据文件IS_PARTITION_FIELD_WRITE_IN_FILEfalse。若表是 text 格式还会从 SerDe 参数中读取field.delim与line.delim作为字段、行分隔符见 HiveSink.java。识别底层存储createHadoopConf()根据hdfsLocation的前缀由StorageFactory.getStorageType(hdfsLocation)返回 HDFS / S3 / COS / OSS 对应的存储实现S3Storage、OSSStorage、COSStorage、HDFSStorage从而为不同 Schema 填入正确的 Hadoop 配置见 HiveSink.java。测试用例 StorageFactoryTest.java 验证了hdfs://、s3n://、s3://、s3a://、oss://、cosn://的识别映射。写入与提交HiveSinkWriter继承BaseFileSinkWriter按分区写文件见 HiveSinkWriter.java。提交阶段HiveSinkAggregatedCommitter在文件提交成功后通过HiveMetaStoreProxy.addPartitions()把本次写入产生的分区追加注册到 Hive Metastore分区已存在时仅打印警告实现自动分区修复见 HiveSinkAggregatedCommitter.java。注意只有文件提交成功errorCommitInfos为空后才注册分区注册失败会把对应提交信息标记为错误并返回重试。Metastore 客户端构造HiveMetaStoreProxy使用HiveMetaStoreClient将metastore_uri写入hive.metastore.uris并支持从hive.hadoop.conf-path目录加载hive-site.xml、从hive_site_path加载指定配置文件同时按配置依次尝试 Kerberos 登录、remote user 登录、匿名直连三种方式创建客户端见 HiveMetaStoreProxy.java。配置项详解官方文档给出的全部参数如下名称类型是否必填默认值table_namestring是-metastore_uristring是-compress_codecstring否nonehdfs_site_pathstring否-hive_site_pathstring否-hive.hadoop.confMap否-hive.hadoop.conf-pathstring否-krb5_pathstring否/etc/krb5.confkerberos_principalstring否-kerberos_keytab_pathstring否-abort_drop_partition_metadataboolean否见下文说明common-options-否-各参数含义如下table_name [string]目标 Hive 表名格式db1.table1。若 Source 为多表模式可用${database_name}.${table_name}作为模板作业运行时会被替换为 Source 端CatalogTable生成的实际库名与表名。metastore_uri [string]Hive Metastore 的 URI如thrift://namenode001:9083。hdfs_site_path [string]hdfs-site.xml的路径用于加载 NameNode 的 HA 配置。hive_site_path [string]hive-site.xml的路径。hive.hadoop.conf [map]直接以键值对形式注入 hadoop 配置可覆盖core-site.xml、hdfs-site.xml、hive-site.xml中的属性。源码中该选项类型为MapString, String默认空 Map见 HiveConfig.java。hive.hadoop.conf-path [string]指定core-site.xml、hdfs-site.xml、hive-site.xml文件的加载目录从源码看主要加载其中的hive-site.xml见 HiveMetaStoreProxy.java。krb5_path [string]krb5.conf的路径用于 Kerberos 认证。kerberos_principal [string]Kerberos 主体principal。kerberos_keytab_path [string]Kerberos keytab 文件路径。abort_drop_partition_metadata [boolean]决定 abort回滚操作时是否从 Hive Metastore 中删除分区元数据。需要特别说明官方文档表格标注默认值为true但当前仓库源码中该选项实际默认值为false见 HiveSinkOptions.java 与 HiveConfig.java。无论取值如何abort 时同步过程中生成的分区数据文件始终会被删除该参数只影响 Metastore 中的分区元数据是否一并清除回滚逻辑见 HiveSinkAggregatedCommitter.java。common-optionsSink 插件通用参数详见 Sink Common Options。工厂层面的必填/可选约束HiveSinkFactory.optionRule()将table_name、metastore_uri设为必填其余abort_drop_partition_metadata、kerberos_principal、kerberos_keytab_path、remote_user、hive.hadoop.conf、hive.hadoop.conf-path均为可选见 HiveSinkFactory.java。配置校验不通过时作业会在提交阶段直接报错。快速开始最小可用配置最简配置只需要两个必填参数Hive { table_name default.seatunnel_orc metastore_uri thrift://namenode001:9083 }示例一单表迁移含完整表结构 DDL假设上游有一张分区源表类型覆盖了 TINYINT 到 STRUCT 等 Hive 常见类型create table test_hive_source( test_tinyint TINYINT, test_smallint SMALLINT, test_int INT, test_bigint BIGINT, test_boolean BOOLEAN, test_float FLOAT, test_double DOUBLE, test_string STRING, test_binary BINARY, test_timestamp TIMESTAMP, test_decimal DECIMAL(8,2), test_char CHAR(64), test_varchar VARCHAR(64), test_date DATE, test_array ARRAYINT, test_map MAPSTRING, FLOAT, test_struct STRUCTstreet:STRING, city:STRING, state:STRING, zip:INT ) PARTITIONED BY (test_par1 STRING, test_par2 STRING);需要把数据从源表读取并写入另一张目标表create table test_hive_sink_text_simple( test_tinyint TINYINT, test_smallint SMALLINT, test_int INT, test_bigint BIGINT, test_boolean BOOLEAN, test_float FLOAT, test_double DOUBLE, test_string STRING, test_binary BINARY, test_timestamp TIMESTAMP, test_decimal DECIMAL(8,2), test_char CHAR(64), test_varchar VARCHAR(64), test_date DATE ) PARTITIONED BY (test_par1 STRING, test_par2 STRING);对应的作业配置文件如下Sink 侧以hive.hadoop.conf注入 S3 相关配置为例输出到 S3 桶s3a://mybucketenv { parallelism 3 job.nametest_hive_source_to_hive } source { Hive { table_name test_hive.test_hive_source metastore_uri thrift://ctyun7:9083 } } sink { # choose stdout output plugin to output data to console Hive { table_name test_hive.test_hive_sink_text_simple metastore_uri thrift://ctyun7:9083 hive.hadoop.conf { bucket s3a://mybucket fs.s3a.aws.credentials.providercom.amazonaws.auth.InstanceProfileCredentialsProvider } }注意目标表test_hive_sink_text_simple与源表在字段上略有差异不含 ARRAY/MAP/STRUCT 三列写入时会按目标 Hive 表的列集合进行列映射这正是 Sink 侧列定义自动取自表元数据的设计带来的灵活性。示例二多表写入${database_name}.${table_name}当存在多个源表需要迁移到对应的目标表时例如批量抽取多张分表create table test_1( ) PARTITIONED BY (xx); create table test_2( ) PARTITIONED BY (xx); ...作业配置中Source 使用tables_configs声明多张表Sink 使用${database_name}.${table_name}模板自动生成与源表一一对应的目标表名env { # You can set flink configuration here parallelism 3 job.nametest_hive_source_to_hive } source { Hive { tables_configs [ { table_name test_hive.test_1 metastore_uri thrift://ctyun6:9083 }, { table_name test_hive.test_2 metastore_uri thrift://ctyun7:9083 } ] } } sink { # choose stdout output plugin to output data to console Hive { table_name ${database_name}.${table_name} metastore_uri thrift://ctyun7:9083 } }该机制依赖 Hive Sink 实现的SupportMultiTableSink能力每个源表在运行时都会解析出独立的库名、表名并替换到模板中。Hive on S3 实战AWS EMR 场景以下步骤以 AWS EMR 环境为例把 Hive Sink 的数据写到 S3。Step 1为 Hive 插件创建 lib 目录mkdir -p ${SEATUNNEL_HOME}/plugins/Hive/libStep 2从 Maven 中心仓库下载所需 jar 到 lib 目录cd ${SEATUNNEL_HOME}/plugins/Hive/lib wget https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/2.6.5/hadoop-aws-2.6.5.jar wget https://repo1.maven.org/maven2/org/apache/hive/hive-exec/2.3.9/hive-exec-2.3.9.jarStep 3从 EMR 环境复制其余依赖 jarcp /usr/share/aws/emr/emrfs/lib/emrfs-hadoop-assembly-2.60.0.jar ${SEATUNNEL_HOME}/plugins/Hive/lib cp /usr/share/aws/emr/hadoop-state-pusher/lib/hadoop-common-3.3.6-amzn-1.jar ${SEATUNNEL_HOME}/plugins/Hive/lib cp /usr/share/aws/emr/hadoop-state-pusher/lib/javax.inject-1.jar ${SEATUNNEL_HOME}/plugins/Hive/lib cp /usr/share/aws/emr/hadoop-state-pusher/lib/aopalliance-1.0.jar ${SEATUNNEL_HOME}/plugins/Hive/libStep 4运行写入 S3 的作业env { parallelism 1 job.mode BATCH } source { FakeSource { schema { fields { pk_id bigint name string score int } primaryKey { name pk_id columnNames [pk_id] } } rows [ { kind INSERT fields [1, A, 100] }, { kind INSERT fields [2, B, 100] }, { kind INSERT fields [3, C, 100] } ] } } sink { Hive { table_name test_hive.test_hive_sink_on_s3 metastore_uri thrift://ip-192-168-0-202.cn-north-1.compute.internal:9083 hive.hadoop.conf-path /home/ec2-user/hadoop-conf hive.hadoop.conf { buckets3://ws-package fs.s3a.aws.credentials.providercom.amazonaws.auth.InstanceProfileCredentialsProvider } } }配置要点hive.hadoop.conf-path指向存放 Hadoop 配置文件含hive-site.xml的目录bucket指定 S3 桶地址fs.s3a.aws.credentials.provider使用 EMR 实例角色的临时凭证InstanceProfile。Hive on OSS 实战阿里云 EMR 场景Step 1创建 lib 目录mkdir -p ${SEATUNNEL_HOME}/plugins/Hive/libStep 2下载 hive-exec jarcd ${SEATUNNEL_HOME}/plugins/Hive/lib wget https://repo1.maven.org/maven2/org/apache/hive/hive-exec/2.3.9/hive-exec-2.3.9.jarStep 3复制 JindoSDK 相关 jar并删除冲突 jarcp -r /opt/apps/JINDOSDK/jindosdk-current/lib/jindo-*.jar ${SEATUNNEL_HOME}/plugins/Hive/lib rm -f ${SEATUNNEL_HOME}/lib/hadoop-aliyun-*.jar在只读仓库中rm命令请在实际 EMR 节点上按需执行本步骤描述的是官方文档给出的依赖清理做法。Step 4运行写入 OSS 的作业env { parallelism 1 job.mode BATCH } source { FakeSource { schema { fields { pk_id bigint name string score int } primaryKey { name pk_id columnNames [pk_id] } } rows [ { kind INSERT fields [1, A, 100] }, { kind INSERT fields [2, B, 100] }, { kind INSERT fields [3, C, 100] } ] } } sink { Hive { table_name test_hive.test_hive_sink_on_oss metastore_uri thrift://master-1-1.c-1009b01725b501f2.cn-wulanchabu.emr.aliyuncs.com:9083 hive.hadoop.conf-path /tmp/hadoop hive.hadoop.conf { bucketoss://emr-osshdfs.cn-wulanchabu.oss-dls.aliyuncs.com } } }OSS 场景的关键在于使用阿里云 EMR 自带的 JindoSDKjindo-*.jar驱动 OSS并移除与 Hadoop 自带 aliyun SDKhadoop-aliyun-*.jar的冲突依赖。版本演进记录Changelog2.2.0-beta2022-09-26新增 Hive Sink 连接器。2.3.0-beta2022-10-20Hive Sink 支持分区自动修复automatic partition repair即提交阶段自动把新写入的分区注册到 Metastore。2.3.02022-12-30修复写文件相关的三个缺陷——上游字段为 null 时抛出NullPointerException、Sink 列映射失败、从状态恢复 writer 时获取事务直接失败。后续版本支持 Kerberos 认证新增partition_dir_expression校验逻辑。常见注意事项依赖与版本官方文档测试的 Hive 版本为 2.3.9Zeta 引擎下必须在$SEATUNNEL_HOME/lib/放置seatunnel-hadoop3-3.1.4-uber.jar、hive-exec-3.1.3.jar、libfb303-0.9.3.jar。格式支持范围底层parseFileFormat目前只接受 text / parquet / orc 三种 InputFormat创建 Hive 目标表时请优先选择这三类存储格式lzo 压缩属于 text 体系的扩展能力。分区元数据与数据的一致性提交成功即注册分区、abort 时删除分区元数据可配置并始终删除同步产生的分区数据这是exactly-once语义在 Hive 场景下的具体体现若手工改动过文件注意与 Metastore 保持一致避免查询不到或重复分区。对象存储认证S3 与 OSS 场景必须同时正确配置bucket与对应的凭证/驱动依赖如 InstanceProfileCredentialsProvider、JindoSDK并确保hive.hadoop.conf-path目录下的hive-site.xml与作业内metastore_uri指向一致。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Hive Sink 连接器完全指南多表写入、精确一次与 S3/OSS 对象存储实践SeaTunnel Hive Sink 连接器完全指南多表写入、精确一次与 S3/OSS 对象存储实践 本文是 SeaTunnel Hive Sink 连接器数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Paimon Sink 连接器实战指南CDC 写入、自动建表与动态分桶SeaTunnel Paimon Sink 连接器实战指南CDC 写入、自动建表与动态分桶 本文以 Apache SeaTunnel 官方仓库 docs/zh数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Paimon Sink 连接器实战指南CDC 同步、自动建表、Schema Evolution 与多表写入全解析SeaTunnel Paimon Sink 连接器实战指南CDC 同步、自动建表、Schema Evolution 与多表写入全解析 导读 本文以 docs/数据集成ETL大数据批处理流处理变更数据捕获上一篇pbrt-v3与实时渲染对比离线渲染的优势与应用场景下一篇AutoGPT Platform 开发协作规范AGENTS.md 中的环境配置、分支策略与 Conventional Commits 实践创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表