写入表

SQL DML

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

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

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 子句有两种类型:

  1. INSERT * 子句要求源表的列与目标表的列相同。
  2. 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

在以下情况下,暂不支持部分更新:

  1. 目标表是引导(bootstrap)表时。
  2. 启用虚拟键时。
  3. 启用读时模式(schema on read)时。
  4. 源数据中包含枚举字段时。
  5. 启用了全局索引时,例如 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 表支持多种配置选项:

选项名称是否必填默认值说明
pathY不适用目标 hoodie 表的基础路径。如果路径不存在则会创建,否则 hudi 表期望能够成功初始化
table.typeNCOPY_ON_WRITE要写入的表类型。COPY_ON_WRITE(或)MERGE_ON_READ
write.operationNupsert本次写入应执行的写操作(支持 insert 或 upsert)
write.precombine.fieldN(无默认值)在实际写入之前用于对记录排序的字段。当两条记录具有相同的键值时,我们将选择排序字段值最大的那条,由 Object.compareTo(..) 确定。注意:此配置已弃用,请改用 ordering.fields
write.payload.classNOverwriteWithLatestAvroPayload.class使用的 Payload 类。如果你想在执行 upsert/insert 时自定义合并逻辑,可以覆盖此类。这将使该选项设置的任何值失效
write.insert.drop.duplicatesNfalse标志,指示插入时是否丢弃重复记录。默认情况下,插入会接受重复记录以获得额外的性能
write.ignore.failedNtrue标志,指示是否在检查点批次内忽略任何非异常错误(例如 writestatus 错误)。默认为 true(以牺牲数据完整性为代价优先保证流式处理的推进)
hoodie.datasource.write.recordkey.fieldNuuid记录键字段。其值将用作 HoodieKey 的 recordKey 组件。实际值将通过对字段值调用 .toString() 获得。嵌套字段可以使用点号表示法指定,例如:a.b.c
hoodie.datasource.write.keygenerator.classNSimpleAvroKeyGenerator.class键生成器类,其将从传入的记录中提取键
write.tasksN4执行实际写入的任务并行度,默认为 4
write.batch.size.MBN128刷新数据到底层文件系统的批次缓冲区大小,单位为 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 计划。

:::

评论

登录后参与评论

正在加载评论…