Apache Hive Metastore
Hive Metastore 是 Apache Hive 提供的一个由关系型数据库支撑的服务,它充当数据仓库或数据湖的目录(catalog)。它可以存储关于表的所有元数据,例如分区、列、列类型等。你也可以将 Hudi 表的元数据同步到 Hive metastore 中。这样不仅可以使用 Hive 查询 Hudi 表,还能借助 Presto、Trino 等交互式查询引擎进行查询。本文将介绍将 Hudi 表同步到 Hive metastore 的不同方式。
Spark 数据源示例
前提条件:正确设置 Hive metastore,并将 hive-site.xml 放置在 $SPARK_HOME/conf 下,使 Spark 安装指向该 Hive metastore。
假设:
- hiveserver2 运行在 10000 端口
- metastore 运行在 9083 端口
然后启动一个以 Hudi spark bundle jar 为依赖的 spark-shell(参见快速开始示例)。
我们可以运行以下脚本来创建一个示例 Hudi 表,并将其同步到 Hive。
// spark-shell
import org.apache.hudi.QuickstartUtils._
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
import org.apache.spark.sql.types._
import org.apache.spark.sql.Row
val databaseName = "my_db"
val tableName = "hudi_cow"
val basePath = "/user/hive/warehouse/hudi_cow"
val schema = StructType(Array(
StructField("rowId", StringType,true),
StructField("partitionId", StringType,true),
StructField("orderingField", LongType,true),
StructField("name", StringType,true),
StructField("versionId", StringType,true),
StructField("toBeDeletedStr", StringType,true),
StructField("intToLong", IntegerType,true),
StructField("longToInt", LongType,true)
))
val data0 = Seq(Row("row_1", "2021/01/01",0L,"bob","v_0","toBeDel0",0,1000000L),
Row("row_2", "2021/01/01",0L,"john","v_0","toBeDel0",0,1000000L),
Row("row_3", "2021/01/02",0L,"tom","v_0","toBeDel0",0,1000000L))
var dfFromData0 = spark.createDataFrame(data0,schema)
dfFromData0.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.table.ordering.fields", "orderingField").
option("hoodie.datasource.write.recordkey.field", "rowId").
option("hoodie.datasource.write.partitionpath.field", "partitionId").
option("hoodie.database.name", databaseName).
option("hoodie.table.name", tableName).
option("hoodie.datasource.write.table.type", "COPY_ON_WRITE").
option("hoodie.datasource.write.operation", "upsert").
option("hoodie.datasource.write.hive_style_partitioning","true").
option("hoodie.datasource.meta.sync.enable", "true").
option("hoodie.datasource.hive_sync.mode", "hms").
option("hoodie.datasource.hive_sync.metastore.uris", "thrift://hive-metastore:9083").
mode(Overwrite).
save(basePath)如果更希望使用 JDBC 而非 HMS 同步模式,请省略 hoodie.datasource.hive_sync.metastore.uris,并改为配置以下项:
hoodie.datasource.hive_sync.mode=jdbc
hoodie.datasource.hive_sync.jdbcurl=<e.g., jdbc:hive2://hiveserver:10000>
hoodie.datasource.hive_sync.username=<username>
hoodie.datasource.hive_sync.password=<password>使用 HiveQL 查询
beeline -u jdbc:hive2://hiveserver:10000/my_db
--hiveconf hive.input.format=org.apache.hadoop.hive.ql.io.HiveInputFormat
--hiveconf hive.stats.autogather=false
Beeline version 1.2.1.spark2 by Apache Hive
0: jdbc:hive2://hiveserver:10000> show tables;
+-----------+--+
| tab_name |
+-----------+--+
| hudi_cow |
+-----------+--+
1 row selected (0.531 seconds)
0: jdbc:hive2://hiveserver:10000> select * from hudi_cow limit 1;
+-------------------------------+--------------------------------+------------------------------+----------------------------------+----------------------------------------------------------------------------+-----------------+-------------------+----------------+---------------------+--------------------------+---------------------+---------------------+-----------------------+--+
| hudi_cow._hoodie_commit_time | hudi_cow._hoodie_commit_seqno | hudi_cow._hoodie_record_key | hudi_cow._hoodie_partition_path | hudi_cow._hoodie_file_name | hudi_cow.rowid | hudi_cow.orderingfield | hudi_cow.name | hudi_cow.versionid | hudi_cow.tobedeletedstr | hudi_cow.inttolong | hudi_cow.longtoint | hudi_cow.partitionid |
+-------------------------------+--------------------------------+------------------------------+----------------------------------+----------------------------------------------------------------------------+-----------------+-------------------+----------------+---------------------+--------------------------+---------------------+---------------------+-----------------------+--+
| 20220120090023631 | 20220120090023631_1_2 | row_1 | partitionId=2021/01/01 | 0bf9b822-928f-4a57-950a-6a5450319c83-0_1-24-314_20220120090023631.parquet | row_1 | 0 | bob | v_0 | toBeDel0 | 0 | 1000000 | 2021/01/01 |
+-------------------------------+--------------------------------+------------------------------+----------------------------------+----------------------------------------------------------------------------+-----------------+-------------------+----------------+---------------------+--------------------------+---------------------+---------------------+-----------------------+--+
1 row selected (5.475 seconds)
0: jdbc:hive2://hiveserver:10000>正确使用分区值提取器
同步到 Hive Metastore 时,分区值是通过 hoodie.datasource.hive_sync.partition_value_extractor 进行提取的。在 0.12 版本之前,该值默认设置为 org.apache.hudi.hive.SlashEncodedDayPartitionValueExtractor,用户通常需要手动覆盖此设置。从 0.12 版本开始,默认值改为更通用的 org.apache.hudi.hive.MultiPartKeysValueExtractor,该提取器使用 / 作为分隔符来提取分区值。
如果使用 TimestampBasedKeyGenerator 等键生成器,分区值可能呈现为 yyyy/MM/dd 的形式。通常不希望将分区值提取为多个部分,例如 [yyyy, MM, dd]。此时,用户可以设置 org.apache.hudi.hive.SinglePartPartitionValueExtractor,以将分区值提取为 yyyy-MM-dd 格式。
当表未分区时,应设置 org.apache.hudi.hive.NonPartitionedExtractor。该值会根据分区字段配置自动推断,因此用户可能无需手动设置。类似地,如果表使用了 Hive 风格分区,则 org.apache.hudi.hive.HiveStylePartitionValueExtractor 会被自动推断并设置。
Hive 同步工具
使用 DataSource 写入器或 Hudi Streamer 写入数据时,支持将表的最新 schema 同步到 Hive Metastore,使查询能够识别新增的列和分区。如果更倾向于从命令行或独立的 JVM 中运行同步操作,Hudi 提供了一个 HiveSyncTool,在构建 hudi-hive 模块后,可以按如下方式调用。以下展示了如何将上述由 DataSource Writer 写入的表同步到 Hive Metastore。
cd hudi-hive
./run_sync_tool.sh --jdbc-url jdbc:hive2://hiveserver:10000 --user hive --pass hive --partitioned-by partition --base-path <basePath> --database default --table <tableName>从 Hudi 0.5.1 版本起,merge-on-read 表的读优化版本默认会加上 '\_ro' 后缀。为了兼容旧版 Hudi,提供了一个可选的 HiveSyncConfig 参数 --skip-ro-suffix,可在需要时关闭 '\_ro' 后缀的添加。使用以下命令可查看其他 Hive 同步选项:
cd hudi-hive
./run_sync_tool.sh
[hudi-hive]$ ./run_sync_tool.sh --helpHive 同步配置
请查看 HiveSyncConfig 中可以传递给 run_sync_tool 的参数。其中,以下参数为必填参数:
@Parameter(names = {"--database"}, description = "name of the target database in Hive", required = true);
@Parameter(names = {"--table"}, description = "name of the target table in Hive", required = true);
@Parameter(names = {"--base-path"}, description = "Basepath of Hudi table to sync", required = true);最常用的 hive 同步配置所对应的数据源选项如下:
:::note
下表中 (N/A) 表示未设置默认值。
:::
| HiveSyncConfig | DataSourceWriteOption | 默认值 | 说明 |
|---|---|---|---|
| --database | hoodie.datasource.hive_sync.database | default | Hive 元数据存储中目标数据库的名称 |
| --table | hoodie.datasource.hive_sync.table | (不适用) | Hive 中目标表的名称。若未指定,则根据 Hudi 表配置中的表名推断得出。 |
| --user | hoodie.datasource.hive_sync.username | hive | Hive 元数据存储的用户名 |
| --pass | hoodie.datasource.hive_sync.password | hive | Hive 元数据存储的密码 |
| --jdbc-url | hoodie.datasource.hive_sync.jdbcurl | jdbc:hive2://localhost:10000 | 使用 jdbc 模式进行同步时的 Hive Server 地址 |
| --sync-mode | hoodie.datasource.hive_sync.mode | (不适用) | Hive 操作所选择的模式。有效值为 hms、jdbc 和 hiveql。详见下文。 |
| --partitioned-by | hoodie.datasource.hive_sync.partition_fields | (不适用) | 表中以逗号分隔的列名,用于确定 Hive 分区。 |
| --partition-value-extractor | hoodie.datasource.hive_sync.partition_extractor_class | org.apache.hudi.hive.MultiPartKeysValueExtractor | 实现 PartitionValueExtractor 以提取分区值的类。会根据指定的分区字段自动推断。 |
同步模式
HiveSyncTool 支持三种模式,即 HMS、HIVEQL、JDBC,用于连接 Hive metastore 服务器。这些模式只是针对 Hive 执行 DDL 的三种不同方式。在这些模式中,JDBC 或 HMS 优于 HIVEQL,后者主要用于执行 DML 而非 DDL。
note
所有这些模式都假定已经配置了 Hive metastore,并在 hive-site.xml 配置文件中设置了相应的属性。此外,如果你使用 spark-shell/spark-sql 将 Hudi 表同步到 Hive,还需要将 hive-site.xml 文件放置在 <SPARK_HOME>/conf 目录下。
HMS
HMS 模式使用 Hive metastore 客户端,通过 thrift API 直接同步 Hudi 表。要使用此模式,请向 run_sync_tool 传入 --sync-mode=hms 并设置 --use-jdbc=false。另外,如果使用远程 metastore,则需要在 hive-site.xml 配置文件中设置 hive.metastore.uris。否则,该工具默认假定 metastore 运行在本地的 9083 端口上。
JDBC
此模式使用 JDBC 规范连接到 Hive metastore。
@Parameter(names = {"--jdbc-url"}, description = "Hive jdbc connect url");HIVEQL
HQL 是 Hive 自有的 SQL 方言。此模式仅使用 Hive QL 的 driver 将 DDL 作为 HQL 命令执行。要使用此模式,请向 run_sync_tool 传入 --sync-mode=hiveql,并将 --use-jdbc 设置为 false。
Flink 环境准备
安装
现在你可以 git clone Hudi master 分支来测试 Flink hive sync。第一步是安装 Hudi,以获取 hudi-flink1.1x-bundle-0.x.x.jar。hudi-flink-bundle 模块的 pom.xml 默认将与 hive 相关的 scope 设置为 provided。如果你要使用 hive sync,则需要在打包时使用 flink-bundle-shade-hive profile。执行以下命令进行安装:
# Maven install command
mvn install -DskipTests -Drat.skip=true -Pflink-bundle-shade-hive2
# For hive1, you need to use profile -Pflink-bundle-shade-hive1
# For hive3, you need to use profile -Pflink-bundle-shade-hive3注意
Hive 1.x 只能将元数据同步到 Hive,但目前无法使用 Hive 进行查询。如果需要查询,可以使用 Spark 查询 Hive 表。
注意
如果使用 hive profile,需要将该 profile 中的 Hive 版本修改为你的 Hive 集群版本(只需修改该 profile 中的 Hive 版本)。此
pom.xml的位置为packaging/hudi-flink-bundle/pom.xml,对应的 profile 位于该文件的底部。
Hive 环境
- 将
hudi-hadoop-mr-bundle导入 Hive。在 Hive 的根目录下创建auxlib/文件夹,并将hudi-hadoop-mr-bundle-0.x.x-SNAPSHOT.jar移动到auxlib中。hudi-hadoop-mr-bundle-0.x.x-SNAPSHOT.jar位于packaging/hudi-hadoop-mr-bundle/target目录下。 - 当 Flink SQL Client 远程连接 Hive Metastore 时,需要启用
hive metastore和hiveserver2服务,并正确设置端口号。启动服务的命令如下:
# Enable hive metastore and hiveserver2
nohup ./bin/hive --service metastore &
nohup ./bin/hive --service hiveserver2 &
# While modifying the jar package under auxlib, you need to restart the service.同步模板
Flink hive 同步现在支持两种 hive 同步模式:hms 和 jdbc。hms 模式只需要配置 metastore uri。对于 jdbc 模式,需要同时配置 JDBC 属性和 metastore uri。选项模板如下:
-- hms mode template
CREATE TABLE t1(
uuid VARCHAR(20),
name VARCHAR(10),
age INT,
ts TIMESTAMP(3),
`partition` VARCHAR(20)
)
PARTITIONED BY (`partition`)
WITH (
'connector' = 'hudi',
'path' = '${db_path}/t1',
'table.type' = 'COPY_ON_WRITE', -- If MERGE_ON_READ, hive query will not have output until the parquet file is generated
'hive_sync.enable' = 'true', -- Required. To enable hive synchronization
'hive_sync.mode' = 'hms', -- Required. Setting hive sync mode to hms, default hms. (Before 0.13, the default sync mode was jdbc.)
'hive_sync.metastore.uris' = 'thrift://${ip}:9083' -- Required. The port need set on hive-site.xml
);
-- jdbc mode template
CREATE TABLE t1(
uuid VARCHAR(20),
name VARCHAR(10),
age INT,
ts TIMESTAMP(3),
`partition` VARCHAR(20)
)
PARTITIONED BY (`partition`)
WITH (
'connector' = 'hudi',
'path' = '${db_path}/t1',
'table.type' = 'COPY_ON_WRITE', --If MERGE_ON_READ, hive query will not have output until the parquet file is generated
'hive_sync.enable' = 'true', -- Required. To enable hive synchronization
'hive_sync.mode' = 'jdbc', -- Required. Setting hive sync mode to jdbc, default hms. (Before 0.13, the default sync mode was jdbc.)
'hive_sync.metastore.uris' = 'thrift://${ip}:9083', -- Required. The port need set on hive-site.xml
'hive_sync.jdbc_url'='jdbc:hive2://${ip}:10000', -- required, hiveServer port
'hive_sync.table'='${table_name}', -- required, hive table name
'hive_sync.db'='${db_name}', -- required, hive database name
'hive_sync.username'='${user_name}', -- required, JDBC username
'hive_sync.password'='${password}' -- required, JDBC password
);查询
使用 Hive beeline 进行查询时,需要输入以下设置:
set hive.input.format = org.apache.hudi.hadoop.hive.HoodieCombineHiveInputFormat;Spark Catalog 元数据存储客户端
当在已经启用 Hive 支持的 Spark 环境中运行 Hudi 时(例如使用 spark.sql.catalogImplementation=hive 的 SparkSQL),标准的 IMetaStoreClient 初始化可能会与 Spark 自身的 Hive 类加载器产生冲突。通过设置
hoodie.datasource.hive_sync.use_spark_catalog=true(默认值:false)使 Hudi 使用 SparkCatalogMetaStoreClient —— 一个 Spark 原生的 IMetaStoreClient 实现 —— 而不是自行创建一个。这可以避免 Hive-on-Spark 环境中的类加载器冲突。需要一个启用了 Hive 支持的 SparkSession 处于活跃状态。
通过 JDBC 回退支持 HMS 4.x
HMS 4.x 修改了若干 Thrift API 方法签名(例如 get_table → get_table_req),这使得基于标准 Thrift 的 HMS 客户端不再兼容。Hudi 1.2.0 增加了自动回退机制:当某次 Thrift 元数据调用的异常原因链中出现 TApplicationException 时,Hudi 会翻转内部的 thriftIncompatible 标志,并将该次同步运行的后续操作全部改走 JDBC 路径。
要求: 同步模式必须为 jdbc,并且配置了有效的 JDBC URL,以便回退客户端可用:
hoodie.datasource.hive_sync.mode=jdbc
hoodie.datasource.hive_sync.jdbcurl=jdbc:hive2://hiveserver:10000
hoodie.datasource.hive_sync.username=<username>
hoodie.datasource.hive_sync.password=<password>如果表使用 mode=hms 或 mode=hiveql 与 HMS 4.x 同步,Hudi 会记录 "Thrift API incompatible with HMS but no JDBC fallback available. Consider using mode=jdbc with a valid jdbcUrl." 并抛出原始异常——不会进行任何自动恢复。
检测范围。 该标志是 HoodieHiveSyncClient 实例级别的,而非全局的,并且只能从 false 变为 true(永远不会重置)。实际上,这意味着每次同步运行的第一次 Thrift 调用会探测一次,该次运行的其余调用都使用 JDBC 回退。下一次同步运行会重新开始一次全新的探测。
JDBC 连接失败会单独暴露。 使用 mode=jdbc 时,Hudi 会在同步客户端构建时就立即建立 JDBC 连接——早于任何 Thrift 调用的尝试。因此,无效的 JDBC URL、缺失的驱动或错误的凭据会在启动阶段失败,报出 HoodieHiveSyncException: Failed to create HiveMetaStoreClient,并在堆栈跟踪中以底层 JDBC 异常作为 cause。这属于配置错误路径,而非 HMS API 不匹配,其行为与 mode=jdbc 连接任意版本的 HMS 完全一致。
VECTOR、BLOB 和 VARIANT 的元数据存储映射
当列类型(VECTOR、BLOB、VARIANT)被同步到外部目录时,Hudi 会将其映射为该目录的原生二进制/结构体类型:
| 类型 | Hive | BigQuery |
|---|---|---|
VECTOR | BINARY | BYTES |
BLOB | STRUCT<type:STRING, data:BINARY, reference:STRUCT<external_path:STRING, offset:BIGINT, length:BIGINT, managed:BOOLEAN>> | 等价的 STRUCT 字段 |
VARIANT | STRUCT<metadata:BINARY, value:BINARY> | 包含 metadata 和 value 字段的 STRUCT(BYTES) |
对于 VECTOR,原始的 VECTOR(dim, elementType) 维度和元素类型元信息会保留在 TBLPROPERTIES/表描述中,从而保证表在经过元数据存储往返之后仍能被 Spark 正确重建。
原生支持 VARIANT 的引擎(Spark 4.0+、Flink 2.1+)会直接依据 Parquet 的 VARIANT 标注读取表数据,而不经过 Hive/BigQuery 的元数据存储表示。read_blob() SQL 函数仅适用于 Spark。对 BLOB 列执行 Hive 和 BigQuery 查询时,会直接读取底层的 struct。
写入器版本表属性
Hudi 1.2.0 的同步会在每次同步时,将表属性 hudi_writer_version(设置为最近一次同步该表的 Hudi 版本)写入 Hive 元存储条目中。这样,工具和元存储管理员就能识别出某个表是由哪个 Hudi 版本写入的。
若要向元存储发出 TOUCH 事件以实现分区级变更跟踪(例如用于下游目录通知),请设置:
hoodie.meta.sync.touch.partitions.enabled=true默认值为 false。启用后,同步操作中每个被修改的分区都会触发一次 TOUCH 事件。
评论
登录后参与评论
KnowForge