快速入门

Spark 快速入门

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

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

本指南通过 Spark 快速了解 Hudi 的核心功能。我们将使用 Spark Datasource API(Scala 和 Python)以及 Spark SQL,通过代码片段演示如何对 Hudi 表进行插入、更新、删除和查询操作。

环境准备

Hudi 支持 Spark 3.3 及以上版本。你可以按照此处的说明安装 Spark。

Spark 支持矩阵

Hudi支持的 Spark 版本Scala 版本Java 版本
1.2.x4.1.x, 4.0.x, 3.5.x(默认构建), 3.4.x, 3.3.x2.13(Spark 4.0/4.1)、2.12/2.13(Spark 3.5)、2.12(Spark 3.3-3.4)17+(Spark 4.0/4.1)、8+(Spark 3.x)
1.1.x4.0.x, 3.5.x(默认构建), 3.4.x, 3.3.x2.13(Spark 4.0)、2.12/2.13(Spark 3.5)、2.12(Spark 3.3-3.4)17+(Spark 4.0)、8+(Spark 3.x)
1.0.x3.5.x(默认构建), 3.4.x, 3.3.x2.12/2.13(Spark 3.5)、2.12(Spark 3.3-3.4)8+
0.15.x3.5.x(默认构建), 3.4.x, 3.3.x, 3.2.x, 3.1.x, 3.0.x2.12/2.13(Spark 3.5)、2.128+
0.14.x3.4.x(默认构建), 3.3.x, 3.2.x, 3.1.x, 3.0.x2.128+

note

默认构建的 Spark 版本表示我们构建 hudi-spark3-bundle 时所使用的版本。

Spark Shell/SQL

  • Scala
  • Python
  • Spark SQL

在解压后的目录中运行带有 Hudi 的 spark-shell:

# For Spark versions: 3.3 - 4.1
export SPARK_VERSION=3.5
export HUDI_VERSION=1.2.0
# For Scala versions: 2.12/2.13
export SCALA_VERSION=2.13

spark-shell --master "local[2]" \
  --packages org.apache.hudi:hudi-spark$SPARK_VERSION-bundle_$SCALA_VERSION:$HUDI_VERSION \
  --conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
  --conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog' \
  --conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension' \
  --conf 'spark.kryo.registrator=org.apache.spark.HoodieSparkKryoRegistrar'

注意: 你必须根据自身的环境需求调整 SPARK_VERSION 和 SCALA_VERSION 变量。

on Kryo serialization

建议用户设置此配置,以降低 Kryo 序列化的开销。

--conf 'spark.kryo.registrator=org.apache.spark.HoodieSparkKryoRegistrar'

设置项目

下面,我们执行导入操作,并设置表名及其对应的基础路径。

  • Scala
  • Python
  • Spark SQL
// spark-shell

val tableName = "trips_table"
val basePath = "file:///tmp/trips_table"

创建表

首先,让我们创建一个 Hudi 表。这里我们以分区表为例进行说明,但 Hudi 同样支持非分区表。

  • Scala
  • Python
  • Spark SQL
// scala
// First commit will auto-initialize the table, if it did not exist in the specified base path.

插入数据

  • Scala
  • Python
  • Spark SQL

生成一些新记录作为 DataFrame,并将该 DataFrame 写入 Hudi 表中。由于这是第一次写入,系统会自动创建该表。

// spark-shell
val columns = Seq("ts","uuid","rider","driver","fare","city")
val data =
  Seq((1695159649087L,"334e26e9-8355-45cc-97c6-c31daf0df330","rider-A","driver-K",19.10,"san_francisco"),
    (1695091554788L,"e96c4396-3fad-413a-a942-4cb36106d721","rider-C","driver-M",27.70 ,"san_francisco"),
    (1695046462179L,"9909a8b1-2d15-4d3d-8ec9-efc48c536a00","rider-D","driver-L",33.90 ,"san_francisco"),
    (1695516137016L,"e3cf430c-889d-4015-bc98-59bdce1e530c","rider-F","driver-P",34.15,"sao_paulo"    ),
    (1695115999911L,"c8abbe79-8d89-47ea-b4ce-4d224bae5bfa","rider-J","driver-T",17.85,"chennai"));

var inserts = spark.createDataFrame(data).toDF(columns:_*)
inserts.write.format("hudi").
  option("hoodie.datasource.write.partitionpath.field", "city").
  option("hoodie.table.name", tableName).
  mode("overwrite").
  save(basePath)

映射到 Hudi 写入操作

Hudi 提供了多种写入操作,包括批量写入和增量写入,可用于将数据写入 Hudi 表,它们具有不同的语义和性能表现。当未配置记录键(record key)时(参见下文的键),将选择 bulk_insert 作为写入操作,这与 Spark 的 Parquet 数据源的默认行为一致。

查询数据

Hudi 表可以被查询回 DataFrame 或 Spark SQL 中。

  • Scala
  • Python
  • Spark SQL
// spark-shell
val tripsDF = spark.read.format("hudi").load(basePath)
tripsDF.createOrReplaceTempView("trips_table")

spark.sql("SELECT uuid, fare, ts, rider, driver, city FROM  trips_table WHERE fare > 20.0").show()
spark.sql("SELECT _hoodie_commit_time, _hoodie_record_key, _hoodie_partition_path, rider, driver, fare FROM  trips_table").show()

更新数据

Hudi 表可以通过流式写入 DataFrame 或使用标准 UPDATE 语句进行更新。

  • Scala
  • Python
  • Spark SQL
// Lets read data from target Hudi table, modify fare column for rider-D and update it.
val updatesDf = spark.read.format("hudi").load(basePath).filter($"rider" === "rider-D").withColumn("fare", col("fare") * 10)

updatesDf.write.format("hudi").
  option("hoodie.datasource.write.operation", "upsert").
  option("hoodie.datasource.write.partitionpath.field", "city").
  option("hoodie.table.name", tableName).
  mode("append").
  save(basePath)

关键要求

使用 spark-datasource 进行更新时,只有当源 DataFrame 包含 Hudi 的元字段,或者配置了键字段时才可行。注意,保存模式现在是 Append。通常情况下,除非是首次创建表,否则始终使用追加模式。

再次查询数据将显示已更新的记录。每次写入操作会生成一个新的提交。对于给定的 _hoodie_record_key 值,比较前一次提交中 _hoodie_commit_time、fare 等字段的变化。

合并数据

  • Scala
  • Python
  • Spark SQL
// spark-shell
Feel free to use "upsert" operation as showed under "Update data" section. Or leverage MergeInto with Spark sql writes.

删除数据

删除操作会从表中移除指定的记录。例如,下面的代码片段会删除传入的 HoodieKey 所对应的记录。更多详情请参阅删除章节。

  • Scala
  • Python
  • Spark SQL
// spark-shell
// Lets  delete rider: rider-D
val deletesDF = spark.read.format("hudi").load(basePath).filter($"rider" === "rider-F")

deletesDF.write.format("hudi").
  option("hoodie.datasource.write.operation", "delete").
  option("hoodie.datasource.write.partitionpath.field", "city").
  option("hoodie.table.name", tableName).
  mode("append").
  save(basePath)

再次查询数据时,将不会显示已删除的记录。

关键要求

使用 spark-datasource 删除数据时,仅当源 DataFrame 包含 Hudi 的元字段,或者已配置键字段时才受支持。请注意,保存模式再次为 Append。

索引数据

Hudi 支持在列上建立索引以加速查询。可以使用 CREATE INDEX 语句在列上创建索引。

note

请注意,要创建二级索引,需满足以下条件:

  1. 表必须具有主键,且合并模式应为 COMMIT_TIME_ORDERING。
  2. 必须启用记录索引。可通过设置 hoodie.metadata.record.index.enable=true 并创建 record_index 来实现。请参见下方示例。
  • Scala
  • Spark SQL
-- 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);

INSERT INTO hudi_indexed_table
VALUES
(1695159649,'334e26e9-8355-45cc-97c6-c31daf0df330','rider-A','driver-K',19.10,'san_francisco'),
(1695091554,'e96c4396-3fad-413a-a942-4cb36106d721','rider-C','driver-M',27.70 ,'san_francisco'),
(1695046462,'9909a8b1-2d15-4d3d-8ec9-efc48c536a00','rider-D','driver-L',33.90 ,'san_francisco'),
(1695332066,'1dced545-862b-4ceb-8b43-d2a568f6616b','rider-E','driver-O',93.50,'san_francisco'),
(1695516137,'e3cf430c-889d-4015-bc98-59bdce1e530c','rider-F','driver-P',34.15,'sao_paulo'    ),
(1695376420,'7a84095f-737f-40bc-b62f-6b69664712d2','rider-G','driver-Q',43.40 ,'sao_paulo'    ),
(1695173887,'3eeb61f7-c2b0-4636-99bd-5d7a5a1d2c04','rider-I','driver-S',41.06 ,'chennai'      ),
(1695115999,'c8abbe79-8d89-47ea-b4ce-4d224bae5bfa','rider-J','driver-T',17.85,'chennai');

-- Setting the Lock Provider
SET hoodie.write.lock.provider = org.apache.hudi.client.transaction.lock.InProcessLockProvider;
-- 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';

-- To create secondary index, first create the record index
SET hoodie.metadata.record.index.enable=true;
CREATE INDEX record_index ON hudi_indexed_table (uuid);
-- Create secondary index on rider column
CREATE INDEX idx_rider ON hudi_indexed_table (rider);

-- 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';

-- Update a record and query the table based on indexed columns
UPDATE hudi_indexed_table SET rider = 'rider-B', driver = 'driver-N', ts = '1697516137' WHERE rider = 'rider-A';
-- Data skipping would be performed using column stat expression index
SELECT uuid, rider FROM hudi_indexed_table WHERE from_unixtime(ts, 'yyyy-MM-dd') = '2023-10-17';
-- Data skipping would be performed using bloom filter expression index
SELECT * FROM hudi_indexed_table WHERE driver = 'driver-N';
-- Data skipping would be performed using secondary index
SELECT * FROM hudi_indexed_table WHERE rider = 'rider-B';

-- Drop all the indexes
DROP INDEX record_index on hudi_indexed_table;
DROP INDEX secondary_index_idx_rider on hudi_indexed_table;
DROP INDEX expr_index_idx_bloom_driver on hudi_indexed_table;
DROP INDEX expr_index_idx_column_ts on hudi_indexed_table;
-- No indexes should show up for the table
SHOW INDEXES FROM hudi_indexed_table;

SET hoodie.metadata.record.index.enable=false;

时间旅行查询

Hudi 支持时间旅行查询,可用于查询历史某个时间点的表数据。支持三种时间戳格式,如下所示。

  • Scala
  • Python
  • Spark SQL
spark.read.format("hudi").
  option("as.of.instant", "20210728141108100").
  load(basePath)

spark.read.format("hudi").
  option("as.of.instant", "2021-07-28 14:11:08.200").
  load(basePath)

// It is equal to "as.of.instant = 2021-07-28 00:00:00"
spark.read.format("hudi").
  option("as.of.instant", "2021-07-28").
  load(basePath)

增量查询

Hudi 提供了一种独特的能力,可以在起始提交时间与结束提交时间之间获取一组发生变更的记录,并为每条记录提供截至结束提交时间的"最新状态"。默认情况下,Hudi 表已配置为支持增量查询,使用记录级别的元数据跟踪。

下面,我们从指定的起始时间开始获取变更,结束时间默认为表的最新提交时间。用户也可以通过 END_INSTANTTIME.key() 选项指定结束时间。

  • Scala
  • Python
  • Spark SQL
// spark-shell
spark.read.format("hudi").load(basePath).createOrReplaceTempView("trips_table")

val commits = spark.sql("SELECT DISTINCT(_hoodie_commit_time) AS commitTime FROM  trips_table ORDER BY commitTime").map(k => k.getString(0)).take(50)
val beginTime = commits(commits.length - 2) // commit time we are interested in

// incrementally query data
val tripsIncrementalDF = spark.read.format("hudi").
  option("hoodie.datasource.query.type", "incremental").
  option("hoodie.datasource.read.begin.instanttime", 0).
  load(basePath)
tripsIncrementalDF.createOrReplaceTempView("trips_incremental")

spark.sql("SELECT `_hoodie_commit_time`, fare, rider, driver, uuid, ts FROM  trips_incremental WHERE fare > 20.0").show()

变更数据捕获查询

Hudi 还提供了一等公民级别的变更数据捕获(Change Data Capture,CDC)查询支持。CDC 查询适用于需要在给定的提交时间范围内获取所有变更及其记录前后镜像的应用场景。

  • Scala
  • Python
  • Spark SQL
// spark-shell
// Lets first insert data to a new table with cdc enabled.
val columns = Seq("ts","uuid","rider","driver","fare","city")
val data =
  Seq((1695158649187L,"334e26e9-8355-45cc-97c6-c31daf0df330","rider-A","driver-K",19.10,"san_francisco"),
    (1695091544288L,"e96c4396-3fad-413a-a942-4cb36106d721","rider-B","driver-L",27.70 ,"san_paulo"),
    (1695046452379L,"9909a8b1-2d15-4d3d-8ec9-efc48c536a00","rider-C","driver-M",33.90 ,"san_francisco"),
    (1695332056404L,"1dced545-862b-4ceb-8b43-d2a568f6616b","rider-D","driver-N",93.50,"chennai"));
var df = spark.createDataFrame(data).toDF(columns:_*)

// Insert data
df.write.format("hudi").
  option("hoodie.datasource.write.partitionpath.field", "city").
  option("hoodie.table.cdc.enabled", "true").
  option("hoodie.table.name", tableName).
  mode("overwrite").
  save(basePath)

// Update fare for riders: rider-A and rider-B
val updatesDf = spark.read.format("hudi").load(basePath).filter($"rider" === "rider-A" || $"rider" === "rider-B").withColumn("fare", col("fare") * 10)

updatesDf.write.format("hudi").
  option("hoodie.datasource.write.operation", "upsert").
  option("hoodie.datasource.write.partitionpath.field", "city").
  option("hoodie.table.cdc.enabled", "true").
  option("hoodie.table.name", tableName).
  mode("append").
  save(basePath)


// Query CDC data
spark.read.option("hoodie.datasource.read.begin.instanttime", 0).
  option("hoodie.datasource.query.type", "incremental").
  option("hoodie.datasource.query.incremental.format", "cdc").
  format("hudi").load(basePath).show(false)

关键要求

请注意,CDC 查询目前仅支持 Copy-on-Write 表。

表类型

到目前为止的示例展示的是 Hudi 支持的两种表类型之一——Copy-on-Write(COW,写时复制)表。Hudi 还支持一种更高级的写优化表类型 Merge-on-Read(MOR,读时合并)表,它能以更灵活的方式平衡读写性能。更多详情请参阅表类型。

上述任意示例都可以在 Merge-on-Read 表上运行,只需在创建表时将表类型改为 MOR 即可,如下所示。

  • Scala
  • Python
  • Spark SQL
// spark-shell
inserts.write.format("hudi").
  ...
  option("hoodie.datasource.write.table.type", "MERGE_ON_READ").
  ...

键

Hudi 还允许用户指定记录键(record key),用于在 Hudi 表中唯一标识一条记录。这对于以一致的方式支持索引和聚类等特性至关重要,这些特性可分别加快更新和查询的速度。关于键的其他好处,这里有详细说明。为此,Hudi 支持多种内置的键生成器,可以轻松地为给定表生成记录键。如果用户没有配置键,Hudi 会自动生成记录键,这些键具有很高的压缩率。

  • Scala
  • Python
  • Spark SQL
// spark-shell
inserts.write.format("hudi").
...
option("hoodie.datasource.write.recordkey.field", "uuid").
...

定义记录键的影响

为 Hudi 表配置键会对表产生新的影响。如果由用户设置记录键,upsert 将被选为写入操作。此外,如果配置了记录键,建议同时指定排序字段,以便正确处理源数据中存在多个相同键记录的情况。详见下文。

合并模式

Hudi 还允许用户指定排序字段,用于对同一记录的多个版本进行排序并解决冲突。这一点在将数据库 CDC 日志应用到 Hudi 表等场景中非常重要,因为在源数据中,由于上游的重复更新,同一条记录可能会出现多次。Hudi 也利用该机制支持数据乱序到达表的场景,在这些场景中,记录可能需要按照与提交时间不同的顺序进行解析。例如,使用 created_at 时间戳字段作为排序字段,可以防止记录的旧版本覆盖新版本或被查询到,即使它们是在更晚的提交时间写入表中的。这是使 Hudi 成为处理流式数据最佳选择的关键特性之一。

为实现不同的合并语义,Hudi 支持合并模式。基于提交时间和事件时间的合并模式开箱即用。用户也可以定义自己的自定义合并策略,参见此处。

  • Scala
  • Python
  • Spark SQL
// spark-shell
updatesDf.write.format("hudi").
  ...
  option("hoodie.table.ordering.fields", "ts").
  ...

接下来可以做什么?

你也可以自行构建 Hudi,并在使用快速入门时加上 --jars <spark bundle jar 的路径>(另请参阅使用 Scala 2.12 构建)以获取更多信息。如果你想寻找将现有数据迁移到 Hudi 的方法,请参阅迁移指南。

Spark SQL 参考

关于 Spark SQL 的高级用法,请参阅 Spark SQL DDL 和 Spark SQL DML 参考指南。关于 ALTER TABLE 命令,请查看此处。存储过程借助 Hudi SparkSQL 提供了许多强大能力,可用于监控、管理和运维 Hudi 表,请查看此处。

流式工作负载

Hudi 为流式数据提供了业界领先的性能和功能。

Hudi Streamer - Hudi 提供了一款增量摄取/ETL 工具 Hudi Streamer,支持以流式方式从各种不同来源将数据摄取到 Hudi 中,并内置了自动检查点、通过 schema 提供者进行 schema 校验、转换支持、自动表服务等强大功能。

Structured Streaming - Hudi 同样支持 Spark Structured Streaming 的读取和写入。更多信息请见此处。

了解更多关于Hudi 中的数据建模的信息,以及执行批量写入和流式写入的不同方式。

Docker 化演示

尽管我们已经展示了核心能力,但 Hudi 还支持更多高级功能,可以帮助你在 Hive、Flink、Spark、Presto、Trino 等多种查询引擎上快速搭建并运行事务型数据湖。我们整理了一个演示视频,在一个基于 Docker 的环境中展示了所有这些功能,所有依赖系统均在本地运行。我们建议你按照此处的步骤复现相同的环境并亲自运行该演示,以获得切身体验。

交互式 Notebooks

另外,你也可以通过交互式笔记本(Notebook)配置,使用 Docker 在本地启动所需的服务。与上文完整的 Docker 化演示环境不同,这是一个更为轻量级的环境,专为快速动手实验而设计。它仅 provisioning 必要的本地组件——包括 Spark、Hive 以及一个本地的 S3 兼容存储——并通过 Docker Compose 无缝打包在一起。你只需要一份克隆的 Hudi 仓库,以及系统上已安装的 Docker(已在 macOS 上测试通过)。具体的分步配置说明和更多详情,请参阅 Notebooks 页面。

评论

登录后参与评论

正在加载评论…