写入表

SQL DDL

师成师成· 更新于 2026-09-29· 阅读 79 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

本页介绍在各引擎中使用 SQL 创建和修改表的支持情况。

Spark SQL

创建表

你可以使用标准的 CREATE TABLE 语法创建表,该语法支持分区设置和传递表属性。

CREATE TABLE [IF NOT EXISTS] [db_name.]table_name
  [(col_name data_type [COMMENT col_comment], ...)]
  [COMMENT table_comment]
  [PARTITIONED BY (col_name, ...)]
  [ROW FORMAT row_format]
  [STORED AS file_format]
  [LOCATION path]
  [TBLPROPERTIES (property_name=property_value, ...)]
  [AS select_statement];

注意:

对于在本地运行本教程、且环境中已集成 Spark-Hive(HMS)的用户:如果你使用 default 数据库,或者在 DDL 语句中未提供 [LOCATION path],Spark 将返回 java.io.IOException: Mkdirs failed to create file:/user/hive/warehouse/hudi_table/.hoodie 错误。要避免此问题,你可以采用以下两种方式之一:

  1. 创建一个数据库,例如执行 CREATE DATABASE hudidb;,并在运行 DDL 语句之前使用它,例如执行 USE hudidb;。
  2. 或者在 DDL 语句中通过 LOCATION 关键字指定一个路径来持久化数据。

创建非分区表

创建非分区表与创建普通表一样简单。

-- create a Hudi table
CREATE TABLE IF NOT EXISTS hudi_table (
  id INT,
  name STRING,
  price DOUBLE
) USING hudi;

创建分区表

通过添加 partitioned by 子句即可创建分区表。分区有助于根据分区列将数据组织到多个文件夹中,同时还可以通过限制需要扫描的元数据、索引和数据量,来加快查询和索引查找的速度。

CREATE TABLE IF NOT EXISTS hudi_table_partitioned (
  id BIGINT,
  name STRING,
  dt STRING,
  hh STRING
) USING hudi
TBLPROPERTIES (
  type = 'cow'
)
PARTITIONED BY (dt, hh);

注意

你也可以通过提供逗号分隔的字段名来创建按多个字段分区的表。创建按多个字段分区的表时,请确保 PARTITIONED BY 子句中指定的列顺序与它们在 CREATE TABLE 模式中出现的顺序一致。例如,对于上面的表,分区字段应指定为 PARTITIONED BY (dt, hh)。

警告

请将分区列声明在 CREATE TABLE 列列表的最后。Spark 会将分区列移动到表模式的末尾,因此提前声明某一列会导致存储的列顺序与你编写的顺序不同。例如,声明 (id, name, price, dt, ts) 并搭配 PARTITIONED BY (dt) 时,表实际存储为 (id, name, price, ts, dt)。随后按位置执行 INSERT INTO ... SELECT 就会把值赋给错误的列:当错位后的类型不兼容时会报 CANNOT_SAFELY_CAST 错误而失败;当类型兼容时则会静默地把值写入错误的列。像 INSERT INTO tbl (id, name, price, dt, ts) SELECT ... 这样显式指定列名,也可以避免这种不匹配。

创建带有记录键和排序字段的表

如此处所述,表会使用记录键(record key)来跟踪表中的每条记录。在前面的示例中,Hudi 为每条新记录自动生成了一个高压缩率的键。如果你想使用已有字段作为键,可以设置 primaryKey 选项。通常,这还需要配置排序字段(通过 orderingFields 选项)来处理乱序数据以及传入写入中可能出现的相同键的重复记录。

注意

你可以根据实际需要为给定表选择多个字段作为主键。例如,primaryKey = 'id, name',这会将这两个字段物化为一个组合键,便于对表进行探索。

下面是同时使用这两个选项创建表的示例。通常,表示事件或事实发生时间的字段(例如订单创建时间、事件生成时间等)会被用作排序字段(通过 orderingFields)。当在表上执行查询时,Hudi 会依据该字段排序来解决同一条记录的多个版本。

CREATE TABLE IF NOT EXISTS hudi_table_keyed (
  id INT,
  name STRING,
  price DOUBLE,
  ts BIGINT
) USING hudi
TBLPROPERTIES (
  type = 'cow',
  primaryKey = 'id',
  orderingFields = 'ts'
);

使用合并模式创建表

Hudi 支持不同的记录合并模式,用于处理传入记录与现有记录的合并。要创建具有特定记录合并模式的表,可以设置 recordMergeMode 选项。

CREATE TABLE IF NOT EXISTS hudi_table_merge_mode (
  id INT,
  name STRING,
  ts LONG,
  price DOUBLE
) USING hudi
TBLPROPERTIES (
  type = 'mor',
  primaryKey = 'id',
  orderingFields = 'ts',
  recordMergeMode = 'EVENT_TIME_ORDERING'
)
LOCATION 'file:///tmp/hudi_table_merge_mode/';

使用 EVENT_TIME_ORDERING 时,事件时间较大(通过 orderingFields 指定)的记录会覆盖同一键上事件时间较小的记录,而不考虑事务的提交时间。用户可以设置 CUSTOM 模式来提供自己的合并逻辑。使用 CUSTOM 合并模式时,您可以提供一个实现了合并逻辑的自定义类。需要实现的接口在这里有详细说明。

CREATE TABLE IF NOT EXISTS hudi_table_merge_mode_custom (
  id INT,
  name STRING,
  ts LONG,
  price DOUBLE
) USING hudi
TBLPROPERTIES (
  type = 'mor',
  primaryKey = 'id',
  orderingFields = 'ts',
  recordMergeMode = 'CUSTOM',
  'hoodie.record.merge.strategy.id' = '<unique-uuid>'
)
LOCATION 'file:///tmp/hudi_table_merge_mode_custom/';

创建包含非结构化和半结构化列类型的表

Hudi 支持三种用于非结构化和半结构化数据的列类型,以及一种 Lance 基础文件格式。

VECTOR(dim[, elementType])

固定维度的嵌入列。elementType 可以是 FLOAT(默认)、DOUBLE 或 INT8(别名 BYTE)之一:

元素类型存储类型
FLOAT(默认)ArrayType(FloatType)
DOUBLEArrayType(DoubleType)
INT8 / BYTEArrayType(ByteType)
CREATE TABLE products (
    product_id   STRING,
    name         STRING,
    embedding    VECTOR(768)          -- defaults to FLOAT
) USING hudi
TBLPROPERTIES (
    primaryKey = 'product_id',
    type = 'cow',
    hoodie.record.merger.impls = 'org.apache.hudi.DefaultSparkRecordMerger'
);

-- Other element types
-- embedding VECTOR(768, FLOAT)
-- embedding VECTOR(768, DOUBLE)
-- embedding VECTOR(256, INT8)

Hudi 的 SQL 解析器会将 VECTOR(128, FLOAT) 规范化为 VECTOR(128)。

VECTOR 列必须是顶层字段;不支持嵌套在 STRUCT、ARRAY 或 MAP 中。建表后,无法通过 schema evolution 修改维度和元素类型。

内部存储。 在 Parquet 中,VECTOR 列以 FIXED_LEN_BYTE_ARRAY 形式存储,并带有 hudi_type=VECTOR(dim[,elementType]) 字段元数据,以便 Spark 读取器将字节解码回类型化数组。在基于 Lance 的表中,VECTOR 以原生方式存储为 FixedSizeList<Float32/Float64, dim>;Lance 仅接受 FLOAT 和 DOUBLE 元素类型。

使用 hudi_vector_search 表值函数查询 VECTOR 列。

引擎限制。 Flink 无法解码 VECTOR 列(底层的 Parquet FIXED_LEN_BYTE_ARRAY 不会被转换回类型化数组);Flink 仍可以读取包含 VECTOR 列的表中的其他列。

BLOB

一种二进制列,有两种存储模式:

模式存储方式读取方式
INLINE原始字节直接嵌入表行直接读取;无需外部获取
OUT_OF_LINE指向外部文件中某段字节范围的指针通过 read_blob() 按需读取
CREATE TABLE media_assets (
    asset_id    STRING,
    file_name   STRING,
    mime_type   STRING,
    file_size   BIGINT,
    content     BLOB
) USING hudi
TBLPROPERTIES (
    primaryKey = 'asset_id',
    type = 'cow'
);

内部结构。 BLOB 列在内部表示为:

STRUCT<
  type:      STRING,        -- 'INLINE' or 'OUT_OF_LINE'
  data:      BINARY,        -- raw bytes for INLINE; null for OUT_OF_LINE
  reference: STRUCT<
    external_path: STRING,  -- file path for OUT_OF_LINE
    offset:        BIGINT,  -- byte offset; null = start of file
    length:        BIGINT,  -- byte length; null = read to end
    managed:       BOOLEAN  -- advisory flag for the (future) cleaner
  >
>

managed 目前仅具指导性:true 表示记录了 Hudi 负责所引用外部文件生命周期的意图(这样未来的清理器可以在没有任何行引用它时删除该文件),false 表示该文件由外部管理。清理器目前尚未使用此字段。

BLOB 列不会被纳入列统计信息索引。

VARIANT

一种半结构化(类似 JSON)的列,以两个二进制字段存储任意 JSON 兼容值(对象、数组、字符串、数字、布尔值、null):

字段说明
metadata对字段名称、类型和结构进行编码,以实现高效访问
value数据负载

在 Spark 4.0 及以上版本中,请声明原生类型:

CREATE TABLE events (
    event_id  STRING,
    payload   VARIANT,
    ts        BIGINT
) USING hudi
TBLPROPERTIES (
    primaryKey = 'event_id',
    preCombineField = 'ts'
);

从 Spark 3.x 读取 Spark 4.0+ 的 VARIANT 表?

若不使用下面的显式 DDL,读取时会抛出 HoodieSchemaException(由 UnsupportedOperationException: VARIANT type is only supported in Spark 4.0+ 引起)。参见从 Spark 3.x 读取。

从 Spark 3.x 读取 VARIANT(向后兼容)

Hudi 1.2.0 支持 Spark 3.4 和 3.5,但两者都无法构造原生的 VariantType。若要从 Spark 3.x 读取 Spark 4.0+ 的 VARIANT 表,请在现有位置创建一个外部 Hudi 表,并将 VARIANT 列声明为二进制 struct:

CREATE TABLE IF NOT EXISTS events_3x (
    event_id  STRING,
    payload   STRUCT<value: BINARY, metadata: BINARY>,
    ts        BIGINT
) USING hudi
LOCATION '/path/to/events'
TBLPROPERTIES (
    primaryKey      = 'event_id',
    preCombineField = 'ts',
    type            = 'cow'     -- or 'mor'
);

SELECT event_id, payload.value, payload.metadata FROM events_3x;

LOCATION 使表成为外部表(仅涉及目录元数据;DROP TABLE 不会删除数据)。primaryKey、preCombineField 和 type 必须与表的 .hoodie/hoodie.properties 保持一致;不一致会误导目录,但不会损坏数据。该映射在 COW 和 MOR 上均可用(日志记录类型 AVRO 和 SPARK 都支持)。

请将该映射视为只读

Spark SQL 允许对映射执行 DML,但写入会在 Hudi 的 Spark 3.x 写入路径中失败,且可能在部分工作(标记文件、日志块)已经完成之后才失败。这里只执行 SELECT;如果是由 Spark 4.0+ 的写入器生成该表,请将映射注册到单独的数据库中,以避免冲突。

你能获得的是原始的 payload.metadata 和 payload.value 字节(采用开放的 Spark Variant 二进制规范),如有需要可在应用代码中解码。而在 Spark 3.x 上你无法获得:

  • parse_json()、variant_get()、cast(payload as STRING):仅 Spark 4.0+ 支持。
  • 将谓词下推到 VARIANT 字段。
  • Schema 自动解析:必须使用上述 DDL;仅调用 spark.read.format("hudi").load(path) 会失败。
引擎对 VARIANT 的支持
引擎行为
Spark 4.0+原生 VariantType,支持在 COW 和 MOR 上进行读取、写入与查询(使用 VARIANT 创建表,或通过 DataFrame 以 VariantType 写入)。
Spark 4.1与 Spark 4.0 相同。Spark 4.1 的 PushVariantIntoScan 会将 VARIANT 投影重写为提取操作组成的结构体;Hudi 能识别该形状,并将该列作为逻辑 VARIANT 返回。
Spark 3.x(3.4 / 3.5)不支持原生 VARIANT。向后兼容读取由 Spark 4.x 写入的表时,需要使用上述显式二进制结构体 DDL;仅能获得原始字节。
Flink < 2.1遇到 VARIANT 列时抛出 UnsupportedOperationException。
Flink ≥ 2.1将 VARIANT 表现为 ROW<metadata BYTES, value BYTES>。Flink 可以读取底层结构体,但无法将其解码为 variant 值。

Lance 支持的表不支持 VARIANT 列。

在用户层面仅暴露未拆分(unshredded)的 VARIANT:所有 Spark/Hudi 架构转换产生的都是未拆分的 VARIANT(包含 metadata 和 value 两个二进制字段)。引擎中虽存在拆分式 variant 的写入路径,但没有相应的 DDL、表属性或会话配置可用于启用它。

Lance 基础文件格式

可按表设置基础文件格式:

CREATE TABLE my_ai_table (
    id        STRING,
    embedding VECTOR(768),
    metadata  STRING
) USING hudi
TBLPROPERTIES (
    primaryKey = 'id',
    type = 'cow',
    hoodie.record.merger.impls = 'org.apache.hudi.DefaultSparkRecordMerger',
    hoodie.table.base.file.format = 'lance'
);

Lance 同样适用于 MOR 表:Lance 文件充当基础文件,而 Avro 日志文件记录增量变更。有关配置、依赖项及行为的说明,请参阅 存储布局 → Lance。

从外部位置创建表

Hudi 表通常由流式写入器创建,例如 streamer 工具,之后可能需要在其上执行一些 SQL 语句。你可以使用 location 语句创建外部表。

CREATE TABLE hudi_table_external
USING hudi
LOCATION 'file:///tmp/hudi_table/';
提示

除了分区列(如果存在)之外,你无需指定 schema 和任何属性。Hudi 可以自动识别 schema 和配置。

Create Table As Select (CTAS)

Hudi 支持 CTAS(Create table as select,即通过查询创建表),用于向 Hudi 表进行初始数据加载。为确保即使是大数据量的加载也能高效完成,CTAS 使用 批量插入(bulk insert)作为写入操作。

# create managed parquet table
CREATE TABLE parquet_table
USING parquet
LOCATION 'file:///tmp/parquet_dataset/';

# CTAS by loading data into Hudi table
CREATE TABLE hudi_table_ctas
USING hudi
TBLPROPERTIES (
  type = 'cow',
  orderingFields = 'ts'
)
PARTITIONED BY (dt)
AS SELECT * FROM parquet_table;

您也可以创建非分区表。

# create managed parquet table
CREATE TABLE parquet_table
USING parquet
LOCATION 'file:///tmp/parquet_dataset/';

# CTAS by loading data into Hudi table
CREATE TABLE hudi_table_ctas
USING hudi
TBLPROPERTIES (
  type = 'cow',
  orderingFields = 'ts'
)
AS SELECT * FROM parquet_table;

如果您倾向于显式设置记录键,可以通过在表属性中设置 primaryKey 配置来实现。

CREATE TABLE hudi_table_ctas
USING hudi
TBLPROPERTIES (
  type = 'cow',
  primaryKey = 'id'
)
PARTITIONED BY (dt)
AS
SELECT 1 AS id, 'a1' AS name, 10 AS price, 1000 AS dt;

你还可以使用 CTAS 在外部位置之间复制数据。

# create managed parquet table
CREATE TABLE parquet_table
USING parquet
LOCATION 'file:///tmp/parquet_dataset/*.parquet';

# CTAS by loading data into hudi table
CREATE TABLE hudi_table_ctas
USING hudi
LOCATION 'file:///tmp/hudi/hudi_tbl/'
TBLPROPERTIES (
  type = 'cow'
)
AS SELECT * FROM parquet_table;

创建索引

Hudi 支持在表上创建和删除不同类型的索引。有关不同索引类型的更多信息,请参阅多模态索引。可以使用 SQL 的 create index 命令创建二级索引、表达式索引和记录索引。

-- Create Index
CREATE INDEX [IF NOT EXISTS] index_name ON [TABLE] table_name
[USING index_type]
(column_name1 [OPTIONS(key1=value1, key2=value2, ...)], column_name2 [OPTIONS(key1=value1, key2=value2, ...)], ...)
[OPTIONS (key1=value1, key2=value2, ...)]

-- Record index syntax
CREATE INDEX indexName ON tableIdentifier (primaryKey1 [, primayKey2 ...]);

-- Secondary Index Syntax
CREATE INDEX indexName ON tableIdentifier (nonPrimaryKey);

-- Expression Index Syntax
CREATE INDEX indexName ON tableIdentifier USING column_stats(col) OPTIONS(expr='expr_val', format='format_val');
CREATE INDEX indexName ON tableIdentifier USING bloom_filters(col) OPTIONS(expr='expr_val');

-- Drop Index
DROP INDEX [IF EXISTS] index_name ON [TABLE] table_name
  • index_name 是要创建或删除的索引的名称。
  • table_name 是创建或删除索引所在的表的名称。
  • index_type 是要创建的索引类型。目前仅支持 column_stats 和 bloom_filters。如果省略 using .. 子句,则会创建一个二级记录索引。
  • column_name 是创建索引所基于的列的名称。

索引本身以及创建索引所在的列,都可以使用键值对形式的选项进行限定说明。

note

请注意,要创建二级索引:

  1. 表必须具有主键,且合并模式应为 COMMIT_TIME_ORDERING。
  2. 必须启用记录索引。这可以通过设置 hoodie.metadata.record.index.enable=true 并随后创建 record_index 来实现。请参考下面的示例。
  3. 复杂类型不支持二级索引。

示例

-- Create a table with primary key
CREATE TABLE hudi_indexed_table (
    ts BIGINT,
    uuid STRING,
    rider STRING,
    driver STRING,
    fare DOUBLE,
    city STRING
) USING HUDI
options(
    primaryKey ='uuid',
    hoodie.write.record.merge.mode = "COMMIT_TIME_ORDERING"
)
PARTITIONED BY (city);

-- Add some data.
INSERT INTO hudi_indexed_table
VALUES
 ...

-- Create bloom filter expression index on driver column
CREATE INDEX idx_bloom_driver ON hudi_indexed_table USING bloom_filters(driver) OPTIONS(expr='identity');
-- It would show bloom filter expression index
SHOW INDEXES FROM hudi_indexed_table;
-- Query on driver column would prune the data using the idx_bloom_driver index
SELECT uuid, rider FROM hudi_indexed_table WHERE driver = 'driver-S';

-- Create column stat expression index on ts column
CREATE INDEX idx_column_ts ON hudi_indexed_table USING column_stats(ts) OPTIONS(expr='from_unixtime', format = 'yyyy-MM-dd');
-- Shows both expression indexes
SHOW INDEXES FROM hudi_indexed_table;
-- Query on ts column would prune the data using the idx_column_ts index
SELECT * FROM hudi_indexed_table WHERE from_unixtime(ts, 'yyyy-MM-dd') = '2023-09-24';

-- Create secondary index on rider column
CREATE INDEX record_index ON hudi_indexed_table (uuid);
CREATE INDEX idx_rider ON hudi_indexed_table (rider);
SET hoodie.metadata.record.index.enable=true;
-- Expression index and secondary index should show up
SHOW INDEXES FROM hudi_indexed_table;
-- Query on rider column would leverage the secondary index idx_rider
SELECT * FROM hudi_indexed_table WHERE rider = 'rider-E';

创建表达式索引

表达式索引是对某一列的函数表达式建立的索引,是 Hudi 多模态索引子系统新增的能力。通过在列的表达式上创建 column_stats 索引,表达式索引可用于实现表的逻辑分区。例如,对时间戳字段提取日期的表达式索引,可以有效地实现基于日期的分区,即使物理布局不同,也能为查询带来相同的收益。

-- Create an expression index on the column `ts` (unix epoch) of the table `hudi_table` using the function `from_unixtime` with the format `yyyy-MM-dd`
CREATE INDEX IF NOT EXISTS ts_datestr ON hudi_table
  USING column_stats(ts)
  OPTIONS(expr='from_unixtime', format='yyyy-MM-dd');
-- Create an expression index on the column `ts` (timestamp in yyyy-MM-dd HH:mm:ss) of the table `hudi_table` using the function `hour`
CREATE INDEX ts_hour ON hudi_table
  USING column_stats(ts)
  options(expr='hour');

note

  1. 表达式索引目前只能通过 Spark 引擎使用 SQL 创建,尚不支持 Spark DataSource API。
  2. 表达式索引目前尚不支持复杂类型。
  3. 表达式索引支持一元表达式以及部分二元表达式。更多详情请参阅 SQL DDL 文档。

创建表达式索引时必须提供 expr 选项,且其取值应为有效的 Spark SQL 函数。有关上述函数的语法,请参阅 Spark SQL 文档,并据此提供相应选项。例如,from_unixtime 函数需要 format 选项。

以下列出了一些受支持且较为常用的函数。

  • identity
  • from_unixtime
  • date_format
  • to_date
  • to_timestamp
  • year
  • month
  • day
  • hour
  • lower
  • upper
  • substring
  • regexp_extract
  • regexp_replace
  • concat
  • length

请注意,目前仅支持以单个列作为输入的函数,且不支持 UDF。

创建并使用表达式索引的完整示例

CREATE TABLE hudi_table_expr_index (
    ts STRING,
    uuid STRING,
    rider STRING,
    driver STRING,
    fare DOUBLE,
    city STRING
) USING HUDI
tblproperties (primaryKey = 'uuid')
PARTITIONED BY (city)
location 'file:///tmp/hudi_table_expr_index';

-- Query with hour function filter but no index yet --
spark-sql> SELECT city, fare, rider, driver FROM  hudi_table_expr_index WHERE  city NOT IN ('chennai') AND hour(ts) > 12;
san_francisco	93.5	rider-E	driver-O
san_francisco	33.9	rider-D	driver-L
sao_paulo	43.4	rider-G	driver-Q
Time taken: 0.208 seconds, Fetched 3 row(s)

spark-sql> EXPLAIN COST SELECT city, fare, rider, driver FROM  hudi_table_expr_index WHERE  city NOT IN ('chennai') AND hour(ts) > 12;
== Optimized Logical Plan ==
Project [city#3465, fare#3464, rider#3462, driver#3463], Statistics(sizeInBytes=899.5 KiB)
+- Filter ((isnotnull(city#3465) AND isnotnull(ts#3460)) AND (NOT (city#3465 = chennai) AND (hour(cast(ts#3460 as timestamp), Some(Asia/Kolkata)) > 12))), Statistics(sizeInBytes=2.5 MiB)
   +- Relation default.hudi_table_expr_index[_hoodie_commit_time#3455,_hoodie_commit_seqno#3456,_hoodie_record_key#3457,_hoodie_partition_path#3458,_hoodie_file_name#3459,ts#3460,uuid#3461,rider#3462,driver#3463,fare#3464,city#3465] parquet, Statistics(sizeInBytes=2.5 MiB)

== Physical Plan ==
*(1) Project [city#3465, fare#3464, rider#3462, driver#3463]
+- *(1) Filter (isnotnull(ts#3460) AND (hour(cast(ts#3460 as timestamp), Some(Asia/Kolkata)) > 12))
   +- *(1) ColumnarToRow
      +- FileScan parquet default.hudi_table_expr_index[ts#3460,rider#3462,driver#3463,fare#3464,city#3465] Batched: true, DataFilters: [isnotnull(ts#3460), (hour(cast(ts#3460 as timestamp), Some(Asia/Kolkata)) > 12)], Format: Parquet, Location: HoodieFileIndex(1 paths)[file:/tmp/hudi_table_expr_index], PartitionFilters: [isnotnull(city#3465), NOT (city#3465 = chennai)], PushedFilters: [IsNotNull(ts)], ReadSchema: struct<ts:string,rider:string,driver:string,fare:double>


-- create the expression index --
CREATE INDEX ts_hour ON hudi_table_expr_index USING column_stats(ts) options(expr='hour');

-- query after creating the index --
spark-sql> SELECT city, fare, rider, driver FROM  hudi_table_expr_index WHERE  city NOT IN ('chennai') AND hour(ts) > 12;
san_francisco	93.5	rider-E	driver-O
san_francisco	33.9	rider-D	driver-L
sao_paulo	43.4	rider-G	driver-Q
Time taken: 0.202 seconds, Fetched 3 row(s)
spark-sql> EXPLAIN COST SELECT city, fare, rider, driver FROM  hudi_table_expr_index WHERE  city NOT IN ('chennai') AND hour(ts) > 12;
== Optimized Logical Plan ==
Project [city#2970, fare#2969, rider#2967, driver#2968], Statistics(sizeInBytes=449.8 KiB)
+- Filter ((isnotnull(city#2970) AND isnotnull(ts#2965)) AND (NOT (city#2970 = chennai) AND (hour(cast(ts#2965 as timestamp), Some(Asia/Kolkata)) > 12))), Statistics(sizeInBytes=1278.3 KiB)
   +- Relation default.hudi_table_expr_index[_hoodie_commit_time#2960,_hoodie_commit_seqno#2961,_hoodie_record_key#2962,_hoodie_partition_path#2963,_hoodie_file_name#2964,ts#2965,uuid#2966,rider#2967,driver#2968,fare#2969,city#2970] parquet, Statistics(sizeInBytes=1278.3 KiB)

== Physical Plan ==
*(1) Project [city#2970, fare#2969, rider#2967, driver#2968]
+- *(1) Filter (isnotnull(ts#2965) AND (hour(cast(ts#2965 as timestamp), Some(Asia/Kolkata)) > 12))
   +- *(1) ColumnarToRow
      +- FileScan parquet default.hudi_table_expr_index[ts#2965,rider#2967,driver#2968,fare#2969,city#2970] Batched: true, DataFilters: [isnotnull(ts#2965), (hour(cast(ts#2965 as timestamp), Some(Asia/Kolkata)) > 12)], Format: Parquet, Location: HoodieFileIndex(1 paths)[file:/tmp/hudi_table_expr_index], PartitionFilters: [isnotnull(city#2970), NOT (city#2970 = chennai)], PushedFilters: [IsNotNull(ts)], ReadSchema: struct<ts:string,rider:string,driver:string,fare:double>

创建分区统计索引

分区统计索引与列统计类似,它跟踪表中各列的 min、max、null、count 等统计信息,可用于查询规划。关键区别在于:column_stats 跟踪的是文件级别的统计信息,而 partition_stats 索引跟踪的是存储分区路径级别的聚合统计信息,从而在查询规划和执行期间更高效地跳过整个文件夹路径。

要启用分区统计索引,只需在建表选项中设置 hoodie.metadata.index.partition.stats.enable = 'true' 即可。

note

  1. partition_stats 索引需要在启用 column_stats 索引的前提下才能使用,两者相辅相成。
  2. partition_stats 索引不会自动为所有列创建,用户必须指定需要为其创建分区统计索引的列列表。
  3. column_stats 和 partition_stats 索引目前尚不支持复杂类型。

创建二级索引

二级索引是构建在表中任意列上的记录级索引。它能够高效支持具有相同二级列值的多条记录,并且构建在基于表记录键(record key)的现有记录级索引之上。二级索引属于基于哈希的索引,通过对键空间进行哈希分片实现可横向扩展的写入性能,同时借助基于行的文件格式实现快速查找。

接下来我们看一个创建包含多个索引的表的示例,以及查询如何利用这些索引同时进行分区裁剪和数据跳过。

DROP TABLE IF EXISTS hudi_table;
-- Let us create a table with multiple partition fields, and enable record index and partition stats index
CREATE TABLE hudi_table (
    ts BIGINT,
    id STRING,
    rider STRING,
    driver STRING,
    fare DOUBLE,
    city STRING,
    state STRING
) USING hudi
 OPTIONS(
    primaryKey ='id',
    hoodie.metadata.record.index.enable = 'true', -- enable record index
    hoodie.metadata.index.partition.stats.enable = 'true', -- enable partition stats index
    hoodie.metadata.index.column.stats.enable = 'true', -- enable column stats
    hoodie.metadata.index.column.stats.column.list = 'rider', -- create column stats index on rider column
    hoodie.write.record.merge.mode = "COMMIT_TIME_ORDERING" -- enable commit time ordering, required for secondary index
)
PARTITIONED BY (city, state)
LOCATION 'file:///tmp/hudi_test_table';

INSERT INTO hudi_table VALUES (1695159649,'trip1','rider-A','driver-K',19.10,'san_francisco','california');
INSERT INTO hudi_table VALUES (1695091554,'trip2','rider-C','driver-M',27.70,'sunnyvale','california');
INSERT INTO hudi_table VALUES (1695332066,'trip3','rider-E','driver-O',93.50,'austin','texas');
INSERT INTO hudi_table VALUES (1695516137,'trip4','rider-F','driver-P',34.15,'houston','texas');

-- simple partition predicate --
select * from hudi_table where city = 'sunnyvale';
20240710215107477	20240710215107477_0_0	trip2	city=sunnyvale/state=california	1dcb14a9-bc4a-4eac-aab5-015f2254b7ec-0_0-40-75_20240710215107477.parquet	1695091554	trip2	rider-C	driver-M	27.7	sunnyvale	california
Time taken: 0.58 seconds, Fetched 1 row(s)

-- simple partition predicate on other partition field --
select * from hudi_table where state = 'texas';
20240710215119846	20240710215119846_0_0	trip4	city=houston/state=texas	08c6ed2c-a87b-4798-8f70-6d8b16cb1932-0_0-74-133_20240710215119846.parquet	1695516137	trip4	rider-F	driver-P	34.15	houston	texas
20240710215110584	20240710215110584_0_0	trip3	city=austin/state=texas	0ab2243c-cc08-4da3-8302-4ce0b4c47a08-0_0-57-104_20240710215110584.parquet	1695332066	trip3	rider-E	driver-O	93.5	austin	texas
Time taken: 0.124 seconds, Fetched 2 row(s)

-- predicate on a column for which partition stats are present --
select id, rider, city, state from hudi_table where rider > 'rider-D';
trip4	rider-F	houston	texas
trip3	rider-E	austin	texas
Time taken: 0.703 seconds, Fetched 2 row(s)

-- record key predicate --
SELECT id, rider, driver FROM hudi_table WHERE id = 'trip1';
trip1	rider-A	driver-K
Time taken: 0.368 seconds, Fetched 1 row(s)

-- create secondary index on driver --
CREATE INDEX driver_idx ON hudi_table (driver);

-- secondary key predicate --
SELECT id, driver, city, state FROM hudi_table WHERE driver IN ('driver-K', 'driver-M');
trip1	driver-K	san_francisco	california
trip2	driver-M	sunnyvale	california
Time taken: 0.83 seconds, Fetched 2 row(s)

创建布隆过滤器索引

布隆过滤器索引会针对被索引的列或列表达式,为每个文件存储一个布隆过滤器。在跳过不包含某个高基数列值(例如 UUID)的文件时,这种索引非常有效。

-- Create a bloom filter index on the column derived from expression `lower(rider)` of the table `hudi_table`
CREATE INDEX idx_bloom_rider ON hudi_indexed_table USING bloom_filters(rider) OPTIONS(expr='lower');

设置 Hudi 配置

你可以通过多种方式为给定的 Hudi 表传递配置。

使用 set 命令

你可以使用 set 命令设置 Hudi 的任何写入配置。这将应用于整个 Spark 会话中的所有操作。

set hoodie.insert.shuffle.parallelism = 100;
set hoodie.upsert.shuffle.parallelism = 100;
set hoodie.delete.shuffle.parallelism = 100;

使用表属性

在创建表时,您还可以配置表选项。该配置仅对当前表生效,并会覆盖任何 SET 命令的值。

CREATE TABLE IF NOT EXISTS tableName (
  colName1 colType1,
  colName2 colType2,
  ...
) USING hudi
TBLPROPERTIES (
  primaryKey = '${colName1}',
  type = 'cow',
  ${hoodie.config.key1} = '${hoodie.config.value1}',
  ${hoodie.config.key2} = '${hoodie.config.value2}',
  ....
);

e.g.
CREATE TABLE IF NOT EXISTS hudi_table (
  id BIGINT,
  name STRING,
  price DOUBLE
) USING hudi
TBLPROPERTIES (
  primaryKey = 'id',
  type = 'cow',
  hoodie.cleaner.fileversions.retained = '20',
  hoodie.keep.max.commits = '20'
);

表属性

用户在创建表时可以设置表属性。下面讨论几个重要的表属性。

参数名称默认值描述
typecow要创建的表类型。type = 'cow' 创建写时复制(COPY-ON-WRITE)表,而 type = 'mor' 创建读时合并(MERGE-ON-READ)表。与 hoodie.datasource.write.table.type 相同。更多详情请参见此处
primaryKeyuuid表的主键字段名称,多个字段之间用逗号分隔。与 hoodie.datasource.write.recordkey.field 相同。如果忽略该配置,Hudi 将自动生成主键。如果显式设置,主键生成将遵循用户的配置。
orderingFields表的排序字段。用于在记录的多个版本之间确定最终版本。通常会使用 event time(事件时间)或类似的列作为排序依据。Hudi 能够利用排序字段的值处理乱序数据。

注意

primaryKey、orderingFields、type 以及其他属性区分大小写。

为并发写入者传入锁提供者

当使用 OCC 和 NBCC(非阻塞并发控制)并发模式时,Hudi 需要一个锁提供者来支持并发写入者或异步表服务。对于 NBCC 模式,加锁仅用于写入时间线中的提交元数据文件,写入操作按完成时间进行串行化。用户同样可以通过 TBLPROPERTIES 传入这些表属性。下面是一个基于 Zookeeper 的配置示例。

-- Properties to use Lock configurations to support Multi Writers
TBLPROPERTIES(
  hoodie.write.lock.zookeeper.url = "zookeeper",
  hoodie.write.lock.zookeeper.port = "2181",
  hoodie.write.lock.zookeeper.lock_key = "tableName",
  hoodie.write.lock.provider = "org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider",
  hoodie.write.concurrency.mode = "optimistic_concurrency_control",
  hoodie.write.lock.zookeeper.base_path = "/tableName"
)

为表启用列统计 / 记录级索引

Hudi 提供了利用表的丰富元数据和索引来加速 DML 操作与查询的能力。例如,可以通过下列表属性启用列统计信息的收集以实现快速的数据跳过,或者使用记录级索引来进行快速更新或点查。

更多内容请参阅元数据配置。

TBLPROPERTIES(
  'hoodie.metadata.index.column.stats.enable' = 'true'
  'hoodie.metadata.record.index.enable' = 'true'
)

Spark Alter Table

语法

-- Alter table name
ALTER TABLE oldTableName RENAME TO newTableName;

-- Alter table add columns
ALTER TABLE tableIdentifier ADD COLUMNS(colAndType [, colAndType]);

示例

--rename to:
ALTER TABLE hudi_table RENAME TO hudi_table_renamed;

--add column:
ALTER TABLE hudi_table ADD COLUMNS(remark STRING);

修改表属性

语法

-- alter table ... set|unset
ALTER TABLE tableIdentifier SET|UNSET TBLPROPERTIES (table_property = 'property_value');

示例

ALTER TABLE hudi_table SET TBLPROPERTIES (hoodie.keep.max.commits = '10');
ALTER TABLE hudi_table SET TBLPROPERTIES ("note" = "don't drop this table");

ALTER TABLE hudi_table UNSET TBLPROPERTIES IF EXISTS (hoodie.keep.max.commits);
ALTER TABLE hudi_table UNSET TBLPROPERTIES IF EXISTS ('note');

注意

目前,尝试更改列类型可能会抛出错误 ALTER TABLE CHANGE COLUMN is not supported for changing column colName with oldColType to colName with newColType.,这是由于存在一个未修复的 SPARK 问题所致。

修改配置选项

你也可以通过 ALTER TABLE SET SERDEPROPERTIES 来修改表的写入配置

语法

-- alter table ... set|unset
ALTER TABLE tableName SET SERDEPROPERTIES ('property' = 'property_value');

示例

 ALTER TABLE hudi_table SET SERDEPROPERTIES ('key1' = 'value1');

显示与删除分区

语法

-- Show partitions
SHOW PARTITIONS tableIdentifier;

-- Drop partition
ALTER TABLE tableIdentifier DROP PARTITION ( partition_col_name = partition_col_val [ , ... ] );

示例

--Show partition:
SHOW PARTITIONS hudi_table;

--Drop partition:
ALTER TABLE hudi_table DROP PARTITION (dt='2021-12-09', hh='10');

斜杠分隔的日期分区与 SHOW PARTITIONS

当表以 hoodie.datasource.write.slash.separated.date.partitioning=true 写入时,物理目录布局使用 yyyy/MM/dd 路径。SHOW PARTITIONS 能够正确处理这一点:它以标准的 col=yyyy-MM-dd 显示格式返回分区值,为了便于阅读,会将 / 分隔符规范化回 -。有关配置斜杠分隔分区的详细信息,请参阅Key Generation。

显示与删除索引

语法

-- Show Indexes
SHOW INDEXES FROM tableIdentifier;

-- Drop partition
DROP INDEX indexIdentifier ON tableIdentifier;

示例

-- Show indexes
SHOW INDEXES FROM hudi_indexed_table;

-- Drop Index
DROP INDEX record_index ON hudi_indexed_table;

显示建表语句

语法

SHOW CREATE TABLE tableIdentifier;

示例

SHOW CREATE TABLE hudi_table;

注意事项

Hudi 目前在使用 Spark SQL 创建/修改表时存在以下限制。

  • 当使用 AWS Glue Data Catalog 作为 hive metastore 时,不支持 ALTER TABLE ... RENAME TO ...,因为 Glue 本身不支持重命名表。
  • 为了便于使用,通过 Spark SQL 创建的新 Hudi 表默认会设置 hoodie.datasource.write.hive_style_partitioning=true。可以通过表属性覆盖该设置。

Flink SQL

创建 Catalog

Catalog 用于管理 SQL 表,如果 Catalog 持久化保存表定义,则表可以在会话之间共享。对于 hms 模式,Catalog 还会补充 hive 同步选项。

示例

CREATE CATALOG hoodie_catalog
  WITH (
    'type'='hudi',
    'catalog.path' = '${catalog default root path}',
    'hive.conf.dir' = '${directory where hive-site.xml is located}',
    'mode'='hms' -- supports 'dfs' mode that uses the DFS backend for table DDLs persistence
  );

选项

选项名称是否必需默认值说明
catalog.path是--目录(catalog)表存储的默认路径,该路径用于自动推断表路径,默认表路径为:${catalog.path}/${db_name}/${table_name}
default-database否default默认数据库名称
hive.conf.dir否--hive-site.xml 所在的目录,仅在 hms 模式下有效
mode否dfs支持 hms 模式,使用 HMS 持久化表选项
table.external否false是否创建外部表,仅在 hms 模式下有效

创建表

你可以使用标准的 FLINK SQL CREATE TABLE 语法创建表,该语法支持分区,并可通过 WITH 子句传入 Flink 选项。

CREATE TABLE [IF NOT EXISTS] [catalog_name.][db_name.]table_name
  (
    { <physical_column_definition>
    [ <table_constraint> ][ , ...n]
  )
  [COMMENT table_comment]
  [PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
  WITH (key1=val1, key2=val2, ...)

创建非分区表

创建非分区表与创建普通表一样简单。

-- create a Hudi table
CREATE TABLE hudi_table(
  id BIGINT,
  name STRING,
  price DOUBLE
)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'MERGE_ON_READ'
);

创建分区表

以下是一个创建 Flink 分区表示例。

CREATE TABLE hudi_table(
  id BIGINT,
  name STRING,
  dt STRING,
  hh STRING
)
PARTITIONED BY (`dt`)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'MERGE_ON_READ'
);

创建带有记录键和排序字段的表

以下示例展示了如何像 Spark 一样,创建一个具有记录键和排序字段的 Flink 表。

CREATE TABLE hudi_table(
  id BIGINT PRIMARY KEY NOT ENFORCED,
  name STRING,
  price DOUBLE,
  ts BIGINT
)
PARTITIONED BY (`dt`)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'MERGE_ON_READ',
'ordering.fields' = 'ts'
);

创建无主键的仅追加表

Hudi 1.2.0 支持在不定义 PRIMARY KEY 的情况下创建 Flink 表,以满足纯追加(append)写入的场景。在此模式下,需将 write.operation 设置为 insert;Hudi 不会强制保证记录级唯一性,record key 和排序字段均为可选项。

-- Append-only table: no PRIMARY KEY required
CREATE TABLE hudi_append_table (
  id      BIGINT,
  name    STRING,
  ts      BIGINT,
  city    STRING
)
PARTITIONED BY (`city`)
WITH (
  'connector'       = 'hudi',
  'path'            = 'file:///tmp/hudi_append_table',
  'table.type'      = 'COPY_ON_WRITE',
  'write.operation' = 'insert'
);

INSERT INTO hudi_append_table VALUES (1, 'Alice', 1695159649, 'sf'), (2, 'Bob', 1695091554, 'ny');

注意

没有主键时,Hudi 会使用自动生成的记录键,并且不会执行去重或 upsert 合并。这等同于 bulk_insert 语义,非常适合日志/事件摄取管道,其中每一行传入数据都应按原样追加。如果 write.operation 的值不是 insert 且未定义 PRIMARY KEY,Hudi 将在建表时报错 "Primary key definition is missing"。

在非阻塞并发控制模式下创建表

以下是在非阻塞并发控制模式下创建 Flink 表的示例。

-- This is a datagen source that can generate records continuously
CREATE TABLE sourceT (
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts TIMESTAMP(3),
  `partition` AS 'par1'
) WITH (
  'connector' = 'datagen',
  'rows-per-second' = '200'
);

-- pipeline1: by default, enable the compaction and cleaning services
CREATE TABLE t1 (
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts TIMESTAMP(3),
  `partition` VARCHAR(20)
) WITH (
  'connector' = 'hudi',
  'path' = '/tmp/hudi-demo/t1',
  'table.type' = 'MERGE_ON_READ',
  'index.type' = 'BUCKET',
  'hoodie.write.concurrency.mode' = 'NON_BLOCKING_CONCURRENCY_CONTROL',
  'write.tasks' = '2'
);

-- pipeline2: disable the compaction and cleaning services manually
CREATE TABLE t1_2 (
  uuid VARCHAR(20),
  name VARCHAR(10),
  age INT,
  ts TIMESTAMP(3),
  `partition` VARCHAR(20)
) WITH (
  'connector' = 'hudi',
  'path' = '/tmp/hudi-demo/t1',
  'table.type' = 'MERGE_ON_READ',
  'index.type' = 'BUCKET',
  'hoodie.write.concurrency.mode' = 'NON_BLOCKING_CONCURRENCY_CONTROL',
  'write.tasks' = '2',
  'compaction.schedule.enabled' = 'false',
  'compaction.async.enabled' = 'false',
  'clean.async.enabled' = 'false'
);

-- Submit the pipelines
INSERT INTO t1
SELECT * FROM sourceT;

INSERT INTO t1_2
SELECT * FROM sourceT;

SELECT * FROM t1 LIMIT 20;

修改表

ALTER TABLE tableA RENAME TO tableB;

设置 Hudi 配置

使用表选项

你可以在创建表时通过表选项来配置 hoodie 相关配置。Flink 专用的 hoodie 配置可以参考此处。这些配置将应用于该表的所有操作。

CREATE TABLE IF NOT EXISTS tableName (
  colName1 colType1 PRIMARY KEY NOT ENFORCED,
  colName2 colType2,
  ...
)
WITH (
  'connector' = 'hudi',
  'path' = '${path}',
  ${hoodie.config.key1} = '${hoodie.config.value1}',
  ${hoodie.config.key2} = '${hoodie.config.value2}',
  ....
);

e.g.
CREATE TABLE hudi_table(
  id BIGINT PRIMARY KEY NOT ENFORCED,
  name STRING,
  price DOUBLE,
  ts BIGINT
)
PARTITIONED BY (`dt`)
WITH (
'connector' = 'hudi',
'path' = 'file:///tmp/hudi_table',
'table.type' = 'MERGE_ON_READ',
'ordering.fields' = 'ts',
'hoodie.cleaner.fileversions.retained' = '20',
'hoodie.keep.max.commits' = '20',
'hoodie.datasource.write.hive_style_partitioning' = 'true'
);

支持的类型

SparkHudi备注
booleanboolean
byteint
shortint
integerint
longlong
datedate
timestamptimestamp
floatfloat
doubledouble
stringstring
decimaldecimal
binarybytes
arrayarray
mapmap
structstruct
ArrayType(<elementType>),并带有字段元数据 hudi_type=VECTOR(dim[, elementType])VECTOR固定维度的嵌入列。元素类型为 FloatType、DoubleType 或 ByteType(INT8)。参见创建包含非结构化和半结构化列类型的表。
StructType(type STRING, data BINARY, reference STRUCT<…>),并带有字段元数据 hudi_type=BLOBBLOB二进制列,支持 INLINE / OUT_OF_LINE 存储方式。参见创建包含非结构化和半结构化列类型的表。
VariantType(Spark 4.0+)或 StructType(metadata BINARY NOT NULL, value BINARY NOT NULL),并带有字段元数据 hudi_type=VARIANT(Spark 3.x)VARIANT半结构化(类 JSON)列(未展开)。参见创建包含非结构化和半结构化列类型的表。
char不支持
varchar不支持
numeric不支持
null不支持
object不支持

评论

登录后参与评论

正在加载评论…