模式演进
模式演进是数据管理的重要方面,Hudi 开箱即用地支持写入时的模式演进,并对读取时的模式演进提供实验性支持。本页将介绍 Hudi 中的模式演进支持。
写入时的模式演进
Hudi 开箱即用地支持向后兼容的模式演进场景,例如添加可空字段或提升字段的数据类型。
info
我们建议尽可能采用这种方式。这是一种实用且高效的模式演进方法,已在 Uber、Walmart 和 LinkedIn 等公司的大规模数据湖中得到验证。Confluent 等厂商也在流式数据领域大规模实现了这种方式。鉴于流式数据的持续特性,无法定义一个可能与之前模式不兼容的模式变更(例如重命名列)。
此外,演进后的模式可以在 Presto 和 Spark SQL 等高性能引擎中直接查询,无需额外的列 ID 转换或类型调和开销。下表总结了与不同 Hudi 表类型兼容的模式变更。
传入的模式会自动根据表模式补充缺失的列,并以 null 值填充。为此,我们需要启用配置 hoodie.write.set.null.for.missing.columns,否则数据管道将会失败。
| 架构变更 | COW | MOR | 备注 |
|---|---|---|---|
| 在根层级末尾新增一个可为空的列 | 是 | 是 | 是表示使用演化后的架构写入成功,且随后的读取能够成功读取整个数据集。 |
| 向内部 struct 中新增一个可为空的列(位于末尾) | 是 | 是 | |
| 新增一个带默认值的复杂类型字段(map 和 array) | 是 | 是 | |
| 新增一个可为空的列并调整字段顺序 | 否 | 否 | 写入成功,但如果使用演化后的架构进行的写入只更新了部分基础文件而未更新全部文件,则读取会失败。目前,Hudi 不维护跨基础文件的架构变更历史注册表。不过,如果 upsert 触及了所有基础文件,读取将会成功。 |
新增一个自定义的可为空的 Hudi 元数据列,例如 _hoodie_meta_col | 是 | 是 | |
| 对根层级的字段进行数据类型提升 | 是 | 是 | |
| 对嵌套字段进行数据类型提升 | 是 | 是 | |
| 对复杂类型(map 或 array 的值)进行数据类型提升 | 是 | 是 | |
| 在根层级末尾新增一个不可为空的列 | 否 | 否 | 对于使用 Spark 数据源的 MOR 表,写入成功但读取失败。作为变通方案,可以将该字段设为可为空。 |
| 向内部 struct 中新增一个不可为空的列(位于末尾) | 否 | 否 | |
| 对根层级的字段进行数据类型降级 | 否 | 否 | |
| 对嵌套字段进行数据类型降级 | 否 | 否 | |
| 对复杂类型(map 或 array 的值)进行数据类型降级 | 否 | 否 |
类型提升
下表展示了当传入的列类型发生变化时,表的最终 schema 会是什么(X 表示不允许):
| 传入 Schema ↓ \ 表 Schema → | int | long | float | double | string | bytes |
|---|---|---|---|---|---|---|
| int | int | long | float | double | string | X |
| long | long | long | float | double | string | X |
| float | float | float | float | double | string | X |
| double | double | double | double | double | string | X |
| string | string | string | string | string | string | bytes |
| bytes | X | X | X | X | string | bytes |
读时 Schema 演化
在很多场景中,人们希望能够更灵活地演化 schema。例如:
- 可以新增、删除、修改和移动列(包括嵌套列)。
- 重命名列(包括嵌套列)。
- 对 Array 类型的嵌套列进行新增、删除或相关操作。
Hudi 实验性地支持在写入时允许不向后兼容的 schema 演化场景,并在读取时进行解析。要启用此功能,需要在写入方配置(Datasource)或表属性(SQL)中设置 hoodie.schema.on.read.enable=true。
note
需要 Hudi 版本 > 0.11 以及 Spark 版本 > 3.1.x 和 3.2.1。对于 Spark 3.2.1 及以上版本,还必须设置 spark.sql.catalog.spark_catalog。如果启用了读时 schema,则无法再次禁用,因为表已经接受了这样的 schema 变更。
新增列
-- add columns
ALTER TABLE tableName ADD COLUMNS(col_spec[, col_spec ...])列规范由五个相邻的字段组成。
| 参数 | 说明 |
|---|---|
| col_name | 新列的名称。若要向嵌套 map 类型的列成员 map<string, struct<n: string, a: int> 中添加子列 col1,请将此字段设置为 member.value.col1 |
| col_type | 新列的类型。 |
| nullable | 新列是否允许 null 值。(可选) |
| comment | 新列的注释。(可选) |
| col_position | 新列被添加的位置。该值可以是 FIRST 或 AFTER origin_col。如果设置为 FIRST,新列将被添加到表的第一列之前;如果设置为 AFTER origin_col,新列将被添加到原有列之后。FIRST 仅在向嵌套列添加新子列时使用,不能用于顶层列。AFTER 的使用则没有限制。 |
示例
ALTER TABLE h0 ADD COLUMNS(ext0 string);
ALTER TABLE h0 ADD COLUMNS(new_col int not null comment 'add new column' AFTER col1);
ALTER TABLE complex_table ADD COLUMNS(col_struct.col_name string comment 'add new column to a struct col' AFTER col_from_col_struct);修改列
语法
-- alter table ... alter column
ALTER TABLE tableName ALTER [COLUMN] col_old_name TYPE column_type [COMMENT] col_comment[FIRST|AFTER] column_name参数说明
| 参数 | 说明 |
|---|---|
| tableName | 表名。 |
| col_old_name | 要修改的列的名称。 |
| column_type | 目标列的类型。 |
| col_comment | 对被修改列的可选注释。 |
| column_name | 被修改列的新位置。例如,AFTER column_name 表示目标列放置在 column_name 之后。 |
示例
--- Changing the column type
ALTER TABLE table1 ALTER COLUMN a.b.c TYPE bigint
--- Altering other attributes
ALTER TABLE table1 ALTER COLUMN a.b.c COMMENT 'new comment'
ALTER TABLE table1 ALTER COLUMN a.b.c FIRST
ALTER TABLE table1 ALTER COLUMN a.b.c AFTER x
ALTER TABLE table1 ALTER COLUMN a.b.c DROP NOT NULL列类型变更
| 源\目标 | long | float | double | string | decimal | date | int |
|---|---|---|---|---|---|---|---|
| int | Y | Y | Y | Y | Y | N | Y |
| long | Y | Y | Y | Y | Y | N | N |
| float | N | Y | Y | Y | Y | N | N |
| double | N | N | Y | Y | Y | N | N |
| decimal | N | N | N | Y | Y | N | N |
| string | N | N | N | Y | Y | Y | N |
| date | N | N | N | Y | N | Y | N |
删除列
语法
-- alter table ... drop columns
ALTER TABLE tableName DROP COLUMN|COLUMNS cols示例
ALTER TABLE table1 DROP COLUMN a.b.c
ALTER TABLE table1 DROP COLUMNS a.b.c, x, y重命名列
语法
-- alter table ... rename column
ALTER TABLE tableName RENAME COLUMN old_columnName TO new_columnName示例
ALTER TABLE table1 RENAME COLUMN a.b.c TO x禁用 Hive 元存储列类型兼容性检查
对于在 Hive 元存储中注册的表,元存储会逐个位置将新的列列表与现有列列表进行比较,如果某一位置上的类型不兼容,就会拒绝执行 ALTER TABLE。因此,修改现有列的类型或位置的模式变更——例如使用 FIRST 或 AFTER 添加列——可能会失败并报出如下错误:
The following columns have types incompatible with the existing columns in their respective positions
将 hive.metastore.disallow.incompatible.col.type.changes 设置为 false 即可跳过此检查。该检查由元存储执行,因此属性设置在哪里取决于引擎所连接的元存储。
当 Spark 运行其自带的嵌入式元存储(未配置 hive.metastore.uris)时,元存储与 Spark 共享同一个 JVM,因此在启动作业时,需要使用 Spark 的 spark.hadoop. 前缀来传递该属性:
--conf 'spark.hadoop.hive.metastore.disallow.incompatible.col.type.changes=false'当引擎与远程 Hive 元存储服务通信时,该属性必须在服务端生效。可以将其设置在元存储的 hive-site.xml 中并重启服务:
<property>
<name>hive.metastore.disallow.incompatible.col.type.changes</name>
<value>false</value>
</property>或者在单个 Beeline 会话中使用 metaconf: 前缀覆盖该值,这会将该值推送到此连接的元数据存储中:
set metaconf:hive.metastore.disallow.incompatible.col.type.changes=false;:::note
只有 Hive 的 SQL 处理器会识别 metaconf: 前缀,而且该覆盖设置在 HiveServer2 上生效。Spark SQL 能接受同样的语句并回显出该值,但会将其作为普通的会话属性存储,根本不会传递到元存储(metastore)。旧版 Hive CLI 确实会把值发送到元存储,但并非通过执行 ALTER 的那个连接。在这两种情况下,请改用上面提到的 hive-site.xml 配置项。
Schema Evolution 实战
下面我们通过一个示例来演示 Hudi 的 schema 演进支持。在下面的例子中,我们将新增一个 string 类型的字段,并把某个字段的数据类型从 int 改为 long。
scala> :paste
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 tableName = "hudi_trips_cow"
val basePath = "file:///tmp/hudi_trips_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("intToLong", IntegerType,true)
))
val data1 = Seq(Row("row_1", "part_0", 0L, "bob", "v_0", 0),
Row("row_2", "part_0", 0L, "john", "v_0", 0),
Row("row_3", "part_0", 0L, "tom", "v_0", 0))
var dfFromData1 = spark.createDataFrame(data1, schema)
dfFromData1.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.index.type","SIMPLE").
option("hoodie.table.name", tableName).
mode(Overwrite).
save(basePath)
var tripsSnapshotDF1 = spark.read.format("hudi").load(basePath + "/*/*")
tripsSnapshotDF1.createOrReplaceTempView("hudi_trips_snapshot")
ctrl+D
scala> spark.sql("desc hudi_trips_snapshot").show()
+--------------------+---------+-------+
| col_name|data_type|comment|
+--------------------+---------+-------+
| _hoodie_commit_time| string| null|
|_hoodie_commit_seqno| string| null|
| _hoodie_record_key| string| null|
|_hoodie_partition...| string| null|
| _hoodie_file_name| string| null|
| rowId| string| null|
| partitionId| string| null|
| orderingField| bigint| null|
| name| string| null|
| versionId| string| null|
| intToLong| int| null|
+--------------------+---------+-------+
scala> spark.sql("select rowId, partitionId, orderingField, name, versionId, intToLong from hudi_trips_snapshot").show()
+-----+-----------+-------------+----+---------+---------+
|rowId|partitionId|orderingField|name|versionId|intToLong|
+-----+-----------+-------------+----+---------+---------+
|row_3| part_0| 0| tom| v_0| 0|
|row_2| part_0| 0|john| v_0| 0|
|row_1| part_0| 0| bob| v_0| 0|
+-----+-----------+-------------+----+---------+---------+
// In the new schema, we are going to add a String field and
// change the datatype `intToLong` field from int to long.
scala> :paste
val newSchema = StructType( Array(
StructField("rowId", StringType,true),
StructField("partitionId", StringType,true),
StructField("orderingField", LongType,true),
StructField("name", StringType,true),
StructField("versionId", StringType,true),
StructField("intToLong", LongType,true),
StructField("newField", StringType,true)
))
val data2 = Seq(Row("row_2", "part_0", 5L, "john", "v_3", 3L, "newField_1"),
Row("row_5", "part_0", 5L, "maroon", "v_2", 2L, "newField_1"),
Row("row_9", "part_0", 5L, "michael", "v_2", 2L, "newField_1"))
var dfFromData2 = spark.createDataFrame(data2, newSchema)
dfFromData2.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.index.type","SIMPLE").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)
var tripsSnapshotDF2 = spark.read.format("hudi").load(basePath + "/*/*")
tripsSnapshotDF2.createOrReplaceTempView("hudi_trips_snapshot")
Ctrl + D
scala> spark.sql("desc hudi_trips_snapshot").show()
+--------------------+---------+-------+
| col_name|data_type|comment|
+--------------------+---------+-------+
| _hoodie_commit_time| string| null|
|_hoodie_commit_seqno| string| null|
| _hoodie_record_key| string| null|
|_hoodie_partition...| string| null|
| _hoodie_file_name| string| null|
| rowId| string| null|
| partitionId| string| null|
| orderingField| bigint| null|
| name| string| null|
| versionId| string| null|
| intToLong| bigint| null|
| newField| string| null|
+--------------------+---------+-------+
scala> spark.sql("select rowId, partitionId, orderingField, name, versionId, intToLong, newField from hudi_trips_snapshot").show()
+-----+-----------+-------------+-------+---------+---------+----------+
|rowId|partitionId|orderingField| name|versionId|intToLong| newField|
+-----+-----------+-------------+-------+---------+---------+----------+
|row_3| part_0| 0| tom| v_0| 0| null|
|row_2| part_0| 5| john| v_3| 3|newField_1|
|row_1| part_0| 0| bob| v_0| 0| null|
|row_5| part_0| 5| maroon| v_2| 2|newField_1|
|row_9| part_0| 5|michael| v_2| 2|newField_1|
+-----+-----------+-------------+-------+---------+---------+----------+视频
评论
登录后参与评论
KnowForge