SQL DML
Spark SQL
SparkSQL 提供了多种数据操纵语言(DML)操作,用于与 Hudi 表进行交互。通过这些操作,你可以向 Hudi 表中插入、更新、合并和删除数据。下面让我们逐一了解这些操作。
关于如何使用 SQL 创建 Hudi 表,请参阅 SQL DDL。
INSERT INTO
你可以使用 INSERT INTO 语句通过 Spark SQL 向 Hudi 表中添加数据。以下是一些示例:
INSERT INTO <table>
SELECT <columns> FROM <source>;::: info
INSERT INTO 语句不支持演化表结构。请使用 DDL(例如 ALTER TABLE)或数据源写入(df.write.format("hudi")....save(basePath))来演化表结构。
:::
弃用说明
从 0.14.0 起,hoodie.sql.bulk.insert.enable 和 hoodie.sql.insert.mode 已被弃用。用户应改用 hoodie.spark.sql.insert.into.operation。若需通过 INSERT INTO 处理重复数据,请参阅 插入去重策略配置。
示例:
-- Insert into a copy-on-write (COW) Hudi table
INSERT INTO hudi_cow_nonpcf_tbl SELECT 1, 'a1', 20;
-- Insert into a merge-on-read (MOR) Hudi table
INSERT INTO hudi_mor_tbl SELECT 1, 'a1', 20, 1000;
-- Insert into a COW Hudi table with static partition
INSERT INTO hudi_cow_pt_tbl PARTITION(dt = '2021-12-09', hh='11') SELECT 2, 'a2', 1000;
-- Insert into a COW Hudi table with dynamic partition
INSERT INTO hudi_cow_pt_tbl PARTITION(dt, hh) SELECT 1 AS id, 'a1' AS name, 1000 AS ts, '2021-12-09' AS dt, '10' AS hh;插入 VECTOR 列
VECTOR(dim[, elementType]) 列接受一个 ARRAY,其长度和元素类型必须与声明的维度和元素类型相匹配:
INSERT INTO products VALUES (
'prod_001', 'Running Shoes', 'Lightweight trail runner',
ARRAY(0.123, -0.456, 0.789, /* ... 768 floats ... */)
);查询向量的元素类型必须与语料库嵌入的元素类型完全一致(不允许隐式类型转换)。
插入 BLOB 列
BLOB 列在内部是一个结构体,包含 type(INLINE/OUT_OF_LINE)、data 以及一个 reference 指针结构体。可通过 named_struct 构造该结构体。
内联(Inline) —— 字节直接嵌入到行中:
INSERT INTO media_assets VALUES (
'asset_001', 'logo.png', 'image/png', 45230,
named_struct(
'type', 'INLINE',
'data', /* binary literal or column reference */,
'reference', CAST(NULL AS STRUCT<external_path: STRING, offset: BIGINT, length: BIGINT, managed: BOOLEAN>)
)
);非内联(Out-of-line) — 指向外部文件中某段字节范围的指针:
INSERT INTO media_assets VALUES (
'asset_002', 'video.mp4', 'video/mp4', 1073741824,
named_struct(
'type', 'OUT_OF_LINE',
'data', CAST(NULL AS BINARY),
'reference', named_struct(
'external_path', 's3://my-bucket/media/container_001.bin',
'offset', 8388608,
'length', 1073741824,
'managed', false
)
)
);OUT_OF_LINE 存储模式的一种常见做法是:将多个对象打包到单个二进制容器文件中,并让每个 BLOB 行将 (external_path, offset, length) 三元组存入该容器:
container_001.bin
├── [offset=0, len=45230] → logo.png
├── [offset=45230, len=89012] → photo.jpg
├── [offset=134242, len=1073741824] → video.mp4
└── ...插入 VARIANT 列
VARIANT 类型的列可以接受任意与 JSON 兼容的值。在 Spark 4.0 及以上版本中,可使用 parse_json 构建 VARIANT 值:
INSERT INTO events VALUES
('evt_001', parse_json('{"action": "click", "x": 120, "y": 450}'), 1000),
('evt_002', parse_json('{"action": "purchase", "items": ["sku_a", "sku_b"], "total": 59.99}'), 1001),
('evt_003', parse_json('"simple string value"'), 1002);写操作的映射关系
Hudi 通过 hoodie.spark.sql.insert.into.operation 配置,为 INSERT INTO 语句提供了灵活选择底层写操作的能力。可选的值包括 "bulk_insert"(大批量插入)、"insert"(带小文件管理)以及 "upsert"(带去重/合并)。如果未设置排序字段,则默认选择 "insert";对于已通过 orderingFields 设置了排序字段的表,默认选择 "upsert" 作为写操作。
Insert Overwrite
INSERT OVERWRITE 语句用于替换 Hudi 表中的现有数据。
INSERT OVERWRITE <table>
SELECT <columns> FROM <source>;所有受 INSERT OVERWRITE 语句影响的现有分区都将被源数据替换。以下是一些示例:
-- Overwrite non-partitioned table
INSERT OVERWRITE hudi_mor_tbl SELECT 99, 'a99', 20.0, 900;
INSERT OVERWRITE hudi_cow_nonpcf_tbl SELECT 99, 'a99', 20.0;
-- Overwrite partitioned table with dynamic partition
INSERT OVERWRITE TABLE hudi_cow_pt_tbl SELECT 10, 'a10', 1100, '2021-12-09', '10';
-- Overwrite partitioned table with static partition
INSERT OVERWRITE hudi_cow_pt_tbl PARTITION(dt = '2021-12-09', hh='12') SELECT 13, 'a13', 1100;更新
你可以使用 UPDATE 语句直接修改 Hudi 表中已有的数据。
UPDATE tableIdentifier SET column = EXPRESSION(,column = EXPRESSION) [ WHERE boolExpression]以下是示例:
-- Update data in a Hudi table
UPDATE hudi_mor_tbl SET price = price * 2, ts = 1111 WHERE id = 1;
-- Update data in a partitioned Hudi table
UPDATE hudi_cow_pt_tbl SET name = 'a1_1', ts = 1001 WHERE id = 1;
-- update using non-PK field
update hudi_cow_pt_tbl set ts = 1001 where name = 'a1';info
UPDATE 操作需要指定排序字段(通过 orderingFields)。
Merge Into
MERGE INTO 语句允许你对源数据执行更复杂的更新和合并操作。MERGE INTO 语句与 UPDATE 语句类似,但它允许你为匹配和未匹配的记录指定不同的操作。
MERGE INTO tableIdentifier AS target_alias
USING (sub_query | tableIdentifier) AS source_alias
ON <merge_condition>
[ WHEN MATCHED [ AND <condition> ] THEN <matched_action> ]
[ WHEN NOT MATCHED [ AND <condition> ] THEN <not_matched_action> ]
<merge_condition> =A equal bool condition
<matched_action> =
DELETE |
UPDATE SET * |
UPDATE SET column1 = expression1 [, column2 = expression2 ...]
<not_matched_action> =
INSERT * |
INSERT (column1 [, column2 ...]) VALUES (value1 [, value2 ...])info
MERGE INTO 语句不支持演进表结构。请使用 DDL(例如 ALTER TABLE)或 Datasource 写入(df.write.format("hudi")....save(basePath))来演进表结构。
WHEN NOT MATCHED 子句指定当值不匹配时要执行的操作。INSERT 子句有两种类型:
INSERT *子句要求源表的列与目标表的列相同。INSERT (column1 [, column2 ...]) VALUES (value1 [, value2 ...])子句不要求指定目标表的所有列。对于未指定的目标列,插入NULL值。
对于配置了用户自定义主键的 Hudi 表,MERGE INTO 中的连接条件以及 UPDATE/INSERT INTO 子句应包含该表的主键。
对于 Hudi 自动生成主键的表,MERGE INTO 中的连接条件可以基于任意数据列。
如果 hoodie.record.merge.mode 设置为 EVENT_TIME_ORDERING,则需要通过 orderingFields 设置排序字段,并且 UPDATE/INSERT 子句中必须包含对应的值。
系统强制要求:如果目标表有主键和分区键列,则源表对应的列必须遵循相同的数据类型。此外,如果目标表配置了 hoodie.record.merge.mode = EVENT_TIME_ORDERING(即目标表应具有有效的排序字段配置),则源表对应的排序字段也必须具有相同的数据类型。
以下为示例
-- source table using hudi for testing merging into non-partitioned table
create table merge_source (id int, name string, price double, ts bigint) using hudi
tblproperties (primaryKey = 'id', orderingFields = 'ts');
insert into merge_source values (1, "old_a1", 22.22, 900), (2, "new_a2", 33.33, 2000), (3, "new_a3", 44.44, 2000);
merge into hudi_mor_tbl as target
using merge_source as source
on target.id = source.id
when matched then update set *
when not matched then insert *
;
-- source table using parquet for testing merging into partitioned table
create table merge_source2 (id int, name string, flag string, dt string, hh string) using parquet;
insert into merge_source2 values (1, "new_a1", 'update', '2021-12-09', '10'), (2, "new_a2", 'delete', '2021-12-09', '11'), (3, "new_a3", 'insert', '2021-12-09', '12');
MERGE into hudi_cow_pt_tbl as target
using (
select id, name, '1000' as ts, flag, dt, hh from merge_source2
) source
on target.id = source.id
when matched and flag != 'delete' then
update set id = source.id, name = source.name, ts = source.ts, dt = source.dt, hh = source.hh
when matched and flag = 'delete' then delete
when not matched then
insert (id, name, ts, dt, hh) values(source.id, source.name, source.ts, source.dt, source.hh)
;关键要求
对于配置了用户自定义主键的 Hudi 表,Merge Into 中的连接条件应包含该表的主键。对于由 Hudi 自动生成主键的表,MIT 中的连接条件可以基于任意数据列。
使用部分更新的 Merge Into
部分更新只写入被更新的列,而不是完整的变更记录。当你拥有包含数百列的宽表(机器学习特征存储的典型场景),且每次只有少数几列被更新时,这种方式非常有用。它不仅可以降低写放大,还有助于减少查询延迟。上面的 MERGE INTO 语句可以改为使用部分更新,如下所示。
-- Create a Merge-on-Read table
CREATE TABLE tableName (
id INT,
name STRING,
price DOUBLE,
_ts LONG,
description STRING
) USING hudi
TBLPROPERTIES (
type = 'mor',
primaryKey = 'id',
orderingFields = '_ts'
)
LOCATION '/location/to/basePath';
-- Insert values into the table
INSERT INTO tableName VALUES
(1, 'a1', 10, 1000, 'a1: desc1'),
(2, 'a2', 20, 1200, 'a2: desc2'),
(3, 'a3', 30, 1250, 'a3: desc3');
-- Perform partial updates using a MERGE INTO statement
MERGE INTO tableName t0
USING (
SELECT 1 AS id, 'a1' AS name, 12 AS price, 1001 AS ts
UNION ALL
SELECT 3 AS id, 'a3' AS name, 25 AS price, 1260 AS ts
) s0
ON t0.id = s0.id
WHEN MATCHED THEN UPDATE SET
price = s0.price,
_ts = s0.ts;
SELECT id, name, price, _ts, description FROM tableName;注意,这里我们没有使用 UPDATE SET *,而只是更新 price 和 _ts 列。
note
在以下情况下,暂不支持部分更新:
- 目标表是引导(bootstrap)表时。
- 启用虚拟键时。
- 启用读时模式(schema on read)时。
- 源数据中包含枚举字段时。
- 启用了全局索引时,例如 GLOBAL_SIMPLE。
对于用户自定义了主键的 Hudi 表,MERGE INTO 中的连接条件以及 UPDATE/INSERT INTO 子句应包含该表的主键。
Delete From
你可以使用 DELETE FROM 语句从 Hudi 表中删除数据。
DELETE FROM tableIdentifier [ WHERE boolExpression ]以下示例
-- Delete data from a Hudi table
DELETE FROM hudi_cow_nonpcf_tbl WHERE uuid = 1;
-- Delete data from a MOR Hudi table based on a condition
DELETE FROM hudi_mor_tbl WHERE id % 2 = 0;
-- Delete data using a non-primary key field
DELETE FROM hudi_cow_pt_tbl WHERE name = 'a1';数据跳过与索引
DML 操作可以通过利用列统计信息进行数据跳过,以及通过索引来减少扫描的数据量,从而得到加速。例如,下面的写法通过使用记录级索引,加快了 Hudi 表上 DELETE 操作的速度。
SET hoodie.metadata.record.index.enable=true;
DELETE from hudi_table where uuid = 'c8abbe79-8d89-47ea-b4ce-4d224bae5bfa';这些 DML 操作为你提供了使用 Spark SQL 管理表格的强大工具。你可以通过各种配置选项来控制这些操作的行为,具体说明见文档。
Flink SQL
Flink SQL 提供了多种数据操纵语言(DML)操作,用于与 Hudi 表进行交互。这些操作允许你向 Hudi 表中插入、更新和删除数据。下面我们逐一来了解它们。
Insert Into
你可以使用 INSERT INTO 语句通过 Flink SQL 将数据写入 Hudi 表。以下是一些示例:
INSERT INTO <table>
SELECT <columns> FROM <source>;示例:
-- Insert into a Hudi table
INSERT INTO hudi_table SELECT 1, 'a1', 20;如果 write.operation 为 upsert,则 INSERT INTO 语句不仅会插入新记录,还会更新具有相同记录键的现有行。
-- Insert into a Hudi table in upsert mode
INSERT INTO hudi_table/*+ OPTIONS('write.operation'='upsert')*/ SELECT 1, 'a1', 20;更新
使用 Flink SQL,你可以通过 update 命令来更新 Hudi 表。以下是一些示例:
UPDATE tableIdentifier SET column = EXPRESSION(,column = EXPRESSION) [ WHERE boolExpression]UPDATE hudi_table SET price = price * 2, ts = 1111 WHERE id = 1;关键要求
更新查询仅在批处理执行模式下有效。
Delete From
使用 Flink SQL,你可以通过 delete 命令从 Hudi 表中删除行。以下是一些示例:
DELETE FROM tableIdentifier [ WHERE boolExpression ]DELETE FROM hudi_table WHERE price < 100;关键要求
删除查询仅在批处理执行模式下有效。
查找连接
查找连接通常用于从外部系统查询数据来丰富(enrich)某个表。该连接要求其中一个表具有处理时间属性,而另一个表由查找源连接器(lookup source connector)支撑。
CREATE TABLE datagen_source(
id int,
name STRING,
proctime as PROCTIME()
) WITH (
'connector' = 'datagen',
'rows-per-second'='1',
'number-of-rows' = '2',
'fields.id.kind'='sequence',
'fields.id.start'='1',
'fields.id.end'='2'
);
SELECT o.id,o.name,b.id as id2
FROM datagen_source AS o
JOIN hudi_table/*+ OPTIONS('lookup.join.cache.ttl'= '2 day') */ FOR SYSTEM_TIME AS OF o.proctime AS b on o.id = b.id;设置写入器/读取器配置
使用 Flink SQL,你还可以在查询的同时设置写入器/读取器配置。
INSERT INTO hudi_table/*+ OPTIONS('${hoodie.config.key1}'='${hoodie.config.value1}')*/INSERT INTO hudi_table/*+ OPTIONS('hoodie.keep.max.commits'='true')*/Flink SQL 实战
hudi-flink 模块为 Hudi 的 Source 和 Sink 定义了 Flink SQL 连接器。Sink 表支持多种配置选项:
| 选项名称 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|
| path | Y | 不适用 | 目标 hoodie 表的基础路径。如果路径不存在则会创建,否则 hudi 表期望能够成功初始化 |
| table.type | N | COPY_ON_WRITE | 要写入的表类型。COPY_ON_WRITE(或)MERGE_ON_READ |
| write.operation | N | upsert | 本次写入应执行的写操作(支持 insert 或 upsert) |
| write.precombine.field | N | (无默认值) | 在实际写入之前用于对记录排序的字段。当两条记录具有相同的键值时,我们将选择排序字段值最大的那条,由 Object.compareTo(..) 确定。注意:此配置已弃用,请改用 ordering.fields |
| write.payload.class | N | OverwriteWithLatestAvroPayload.class | 使用的 Payload 类。如果你想在执行 upsert/insert 时自定义合并逻辑,可以覆盖此类。这将使该选项设置的任何值失效 |
| write.insert.drop.duplicates | N | false | 标志,指示插入时是否丢弃重复记录。默认情况下,插入会接受重复记录以获得额外的性能 |
| write.ignore.failed | N | true | 标志,指示是否在检查点批次内忽略任何非异常错误(例如 writestatus 错误)。默认为 true(以牺牲数据完整性为代价优先保证流式处理的推进) |
| hoodie.datasource.write.recordkey.field | N | uuid | 记录键字段。其值将用作 HoodieKey 的 recordKey 组件。实际值将通过对字段值调用 .toString() 获得。嵌套字段可以使用点号表示法指定,例如:a.b.c |
| hoodie.datasource.write.keygenerator.class | N | SimpleAvroKeyGenerator.class | 键生成器类,其将从传入的记录中提取键 |
| write.tasks | N | 4 | 执行实际写入的任务并行度,默认为 4 |
| write.batch.size.MB | N | 128 | 刷新数据到底层文件系统的批次缓冲区大小,单位为 MB |
如果表类型为 MERGE_ON_READ,你还可以通过选项指定异步合并(compaction)策略:
| 选项名称 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|
| compaction.async.enabled | 否 | true | 异步合并,默认对 MOR 表开启 |
| compaction.trigger.strategy | 否 | num_commits | 触发合并的策略,可选值:'num_commits':当达到 N 次增量提交(delta commits)时触发合并;'time_elapsed':当距上次合并经过的时间超过 N 秒时触发合并;'num_and_time':当 NUM_COMMITS 和 TIME_ELAPSED 同时满足时触发合并;'num_or_time':当 NUM_COMMITS 或 TIME_ELAPSED 满足其中之一时触发合并。默认为 'num_commits' |
| compaction.delta_commits | 否 | 5 | 触发合并所需的最大增量提交次数,默认 5 次提交 |
| compaction.delta_seconds | 否 | 3600 | 触发合并所需的最大增量时间(秒),默认 1 小时 |
你可以使用 SQL INSERT INTO 语句写入数据:
INSERT INTO hudi_table select ... from ...;注意:INSERT OVERWRITE 尚未支持,但已在路线图中。
非阻塞并发控制(实验性)
Hudi Flink 支持一种新的非阻塞并发控制模式,多个写入任务可以并发执行而互不阻塞。关于该模式的更多信息,可参阅并发控制文档。下面让我们看看它的实际应用。
在下面的示例中,我们有两条流式摄取管道并发更新同一张表。其中一条管道负责 compaction 和 cleaning 表服务,另一条管道仅用于数据摄取。
为了提交数据集,需要启用 checkpoint。以下是 flink-conf.yaml 的一个配置示例:
-- set the interval as 30 seconds
execution.checkpointing.interval: 30000
state.backend: rocksdb-- 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' = '${work_path}/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' = '${work_path}/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;从上面的示例可以看出,我们有两条流水线,每条流水线包含多个任务,它们并发地写入同一个表。要使用新的并发模式,只需将 hoodie.write.concurrency.mode 设置为 NON_BLOCKING_CONCURRENCY_CONTROL 即可。write.tasks 选项用于指定写入该表时使用的写任务数量。compaction.schedule.enabled、compaction.async.enabled 和 clean.async.enabled 选项用于为第二条流水线禁用压缩(compaction)和清理(clean)服务,这样做是为了确保同一张表的压缩和清理服务不会被执行两次。
一致性哈希索引(实验性)
自 0.13.0 版本起,我们引入了一致性哈希分桶索引(Consistent Hashing Bucket Index)。这是 Hudi 中可用的三种分桶索引变体之一。一致性哈希分桶索引为写入方提供了数据分桶的动态伸缩能力。你可以在该 RFC 中了解此功能的设计。在 0.13.X 版本中,一致性哈希索引仅支持 Spark 引擎;而从 0.14.0 版本开始,该索引也支持 Flink 引擎。
要使用此功能,请将 index.type 选项配置为 BUCKET,并将 hoodie.index.bucket.engine 设置为 CONSISTENT_HASHING。启用一致性哈希索引时,在写入方启用聚类调度非常重要。在此过程中,在聚类尚未完成期间,写入方会对新旧数据分桶执行双重写入。虽然双重写入不会影响正确性,但仍强烈建议尽快执行聚类。
在下面的示例中,我们将创建一个 datagen 数据源,并使用一致性分桶索引将数据流式写入 Hudi 表。为了提交数据集,需要启用 checkpoint(检查点)。以下是 flink-conf.yaml 的示例配置:
-- set the interval as 30 seconds
execution.checkpointing.interval: 30000
state.backend: rocksdb-- 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'
);
-- Create the hudi table with consistent bucket index
CREATE TABLE t1(
uuid VARCHAR(20) PRIMARY KEY NOT ENFORCED,
name VARCHAR(10),
age INT,
ts TIMESTAMP(3),
`partition` VARCHAR(20)
)
PARTITIONED BY (`partition`)
WITH (
'connector'='hudi',
'path' = '${work_path}/hudi-demo/hudiT',
'table.type' = 'MERGE_ON_READ',
'index.type' = 'BUCKET',
'clustering.schedule.enabled'='true',
'hoodie.index.bucket.engine'='CONSISTENT_HASHING',
'hoodie.clustering.plan.strategy.class'='org.apache.hudi.client.clustering.plan.strategy.FlinkConsistentBucketClusteringPlanStrategy',
'hoodie.clustering.execution.strategy.class'='org.apache.hudi.client.clustering.run.strategy.SparkConsistentBucketClusteringExecutionStrategy',
'hoodie.bucket.index.num.buckets'='8',
'hoodie.bucket.index.max.num.buckets'='128',
'hoodie.bucket.index.min.num.buckets'='8',
'hoodie.bucket.index.split.threshold'='1.5',
'write.tasks'='2'
);
-- submit the pipelines
insert into t1 select * from sourceT;
select * from t1 limit 20;:::caution
一致性哈希索引自 0.14.0 版本 起支持 Flink 引擎,截至 0.14.0 使用时存在以下限制:
- 该索引仅支持 MOR 表。即使使用 Spark 引擎,也存在此限制。
- 在启用元数据表的情况下无法使用。即使使用 Spark 引擎,也存在此限制。
- Flink 引擎下一致性哈希索引尚不支持 bulk insert,请在 bulk insert 管道中使用简单桶索引或 Spark 引擎。
- Flink 引擎生成的 resize 计划仅支持合并小文件组,尚不支持文件拆分。
- resize 计划应通过离线 Spark 作业执行。Flink 引擎尚不支持执行 resize 计划。
:::
评论
登录后参与评论
KnowForge