表服务

索引构建

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

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

Hudi 维护着一份可扩展的元数据,其中包含表的若干辅助数据。Hudi 的可插拔索引子系统依赖于元数据表。从用于高效定位记录的 files 索引,到用于数据跳过的 column_stats 索引,不同类型的索引都是元数据表的一部分。任何支持索引的数据系统,其一个根本性权衡在于如何平衡写入吞吐量与索引更新。一种粗暴的做法是在建索引期间阻塞写入。Hudi 支持通过 SQL、Datasource 以及异步索引来创建索引。然而,超大表的索引过程可能耗时数小时。这正是 Hudi 创新性的并发索引大显身手之处。

并发索引

Hudi 中的索引分两个阶段创建,并结合使用了乐观并发控制与多版本并发控制技术。两阶段方法确保其他写入者不会被阻塞。

  • 调度与规划:这是第一阶段,会调度一个索引计划并受锁保护。索引计划会考虑截至索引即时(instant)之前所有已完成的提交。
  • 执行:该阶段按照索引计划创建索引文件。阶段结束时,Hudi 会确保索引即时之后完成的提交使用已创建的索引计划来添加相应的索引元数据。该检查由元数据表锁保护,一旦失败,索引过程将被中止。

我们现在可以在 Hudi 中异步创建不同类型的索引和元数据,包括 bloom_filters、column_stats、partition_stats、record_index、secondary_index 和 expression_index。能够在不阻塞写入的情况下建索引,确保了写入性能不受影响,也无需额外的手动维护来添加或删除索引。同时,通过避免写入与建索引之间的竞争,还减少了资源浪费。

有关如何设置异步索引的更多细节,请参阅设置异步索引一节。想进一步了解异步索引功能的设计,请阅读这篇博客。

使用 SQL 创建索引

目前,二级索引、表达式索引和记录索引等索引可以通过 SQL 的 create index 命令创建。关于这些索引的更多信息,请参阅元数据章节。

note

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

  1. 表必须具有主键,且合并模式应为 COMMIT_TIME_ORDERING。
  2. 必须启用记录索引。可通过设置 hoodie.metadata.global.record.level.index.enable=true 并创建 record_index 来实现。请参考下面的示例。

示例

-- Create record index on primary key - uuid
CREATE INDEX record_index ON hudi_indexed_table (uuid);

-- Create secondary index on rider column.
CREATE INDEX idx_rider ON hudi_indexed_table (rider);

-- Create expression index by performing transformation on ts and driver column
-- The index is created on the transformed column. Here column stats index is created on ts column
-- and bloom filters index is created on driver column.
CREATE INDEX idx_column_ts ON hudi_indexed_table USING column_stats(ts) OPTIONS(expr='from_unixtime', format = 'yyyy-MM-dd');
CREATE INDEX idx_bloom_driver ON hudi_indexed_table USING bloom_filters(driver) OPTIONS(expr='identity');

有关使用 SQL 创建索引的更多信息,请参阅 SQL DDL

使用 Datasource 创建索引

bloom_filters、column_stats、partition_stats 和 record_index 等索引可以通过 Datasource 创建。下面列出了创建上述索引所需的各种配置。

-- [Required Configs] Partition stats
hoodie.metadata.index.partition.stats.enable=true
hoodie.metadata.index.column.stats.enable=true
-- [Optional Configs] - list of columns to index on. By default all columns are indexed
hoodie.metadata.index.column.stats.column.list=col1,col2,...

-- [Required Configs] Column stats
hoodie.metadata.index.column.stats.enable=true
-- [Optional Configs] - list of columns to index on. By default all columns are indexed
hoodie.metadata.index.column.stats.column.list=col1,col2,...

-- [Required Configs] Record Level Index (Global RLI — single record key unique across all partitions)
hoodie.metadata.global.record.level.index.enable=true

-- [Required Configs] Bloom filter Index
hoodie.metadata.index.bloom.filter.enable=true

以下示例展示了如何为使用 Datasource API 创建的表创建索引。

示例

import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.common.table.HoodieTableConfig._
import org.apache.hudi.config.HoodieWriteConfig._
import org.apache.hudi.keygen.constant.KeyGeneratorOptions._
import org.apache.hudi.common.model.HoodieRecord
import spark.implicits._

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

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).
  option("hoodie.write.record.merge.mode", "COMMIT_TIME_ORDERING").
  option(RECORDKEY_FIELD_OPT_KEY, "uuid").
  mode(Overwrite).
  save(basePath)

// Create record index and secondary index for the table
spark.sql(s"CREATE TABLE test_table_external USING hudi LOCATION '$basePath'")
spark.sql(s"SET hoodie.metadata.global.record.level.index.enable=true")
spark.sql(s"CREATE INDEX record_index ON test_table_external (uuid)")
spark.sql(s"CREATE INDEX idx_rider ON test_table_external (rider)")
spark.sql(s"SHOW INDEXES FROM hudi_indexed_table").show(false)
spark.sql(s"SELECT * FROM hudi_indexed_table WHERE rider = 'rider-E'").show(false)

设置异步索引

在本示例中,我们将通过 Hudi Streamer 持续写入数据,同时并行创建索引。示例中使用 HoodieIndexer 来创建索引,以便索引的调度(schedule)和执行(execute)阶段能够清晰可见。这些异步配置同样可以与基于 Datasource 和 SQL 的配置配合使用来创建索引。

首先,我们将生成一个持续写入的工作负载。在下面的示例中,我们将启动一个 Hudi Streamer,它会持续将数据从原始 parquet 文件写入 Hudi 表。我们使用了广泛可获取的 NY Taxi 数据集,其配置详情如下:

摄取写入配置

hoodie.datasource.write.recordkey.field=VendorID
hoodie.datasource.write.partitionpath.field=tpep_dropoff_datetime
hoodie.table.ordering.fields=tpep_dropoff_datetime
hoodie.streamer.source.dfs.root=/Users/home/path/to/data/parquet_files/
hoodie.streamer.schemaprovider.target.schema.file=/Users/home/path/to/schema/schema.avsc
hoodie.streamer.schemaprovider.source.schema.file=/Users/home/path/to/schema/schema.avsc
// set lock provider configs
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider
hoodie.write.lock.zookeeper.url=<zk_url>
hoodie.write.lock.zookeeper.port=<zk_port>
hoodie.write.lock.zookeeper.lock_key=<zk_key>
hoodie.write.lock.zookeeper.base_path=<zk_base_path>

运行 Hudi Streamer

spark-submit \
--jars "packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar,packaging/hudi-spark-bundle/target/hudi-spark3.5-bundle_2.12-1.2.0.jar" \
--class org.apache.hudi.utilities.streamer.HoodieStreamer `ls /Users/home/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar` \
--props `ls /Users/home/path/to/write/config.properties` \
--source-class org.apache.hudi.utilities.sources.ParquetDFSSource  --schemaprovider-class org.apache.hudi.utilities.schema.FilebasedSchemaProvider \
--source-ordering-field tpep_dropoff_datetime   \
--table-type COPY_ON_WRITE \
--target-base-path file:///tmp/hudi-ny-taxi/   \
--target-table ny_hudi_tbl  \
--op UPSERT  \
--continuous \
--source-limit 5000000 \
--min-sync-interval-seconds 60

Hudi 元数据表默认已启用,files index 也会自动创建。在 Hudi Streamer 以持续模式运行时,我们来安排 COLUMN_STATS 索引的索引构建任务。首先需要为索引器定义一个属性配置文件。

配置

如前所述,元数据索引是可插拔的。用户可以根据不断变化的业务需求,在任何时间点添加任意索引。下面列出了启用特定索引的一些配置。目前元数据表下可用的索引可以在此处查看,启用这些索引的配置也可在那里找到。完整的元数据配置集可以在此处查看。

note

启用元数据表并配置锁提供者(lock provider)是使用异步索引器的前提条件。请查看下方的示例配置。

记录级索引配置项

Hudi 支持两种记录级索引(Record Level Index)实现,每种都有各自的启用开关和容量配置:

  • 全局 RLI — 记录键(record key)在整个表范围内(跨分区)唯一。
  • 分区 RLI — partition_path + record_key 在每个分区内唯一。
配置名称默认值说明
hoodie.metadata.global.record.level.index.enablefalse启用全局 RLI(记录级索引)。
hoodie.metadata.global.record.level.index.min.filegroup.count10全局 RLI 的最小文件组数量。
hoodie.metadata.global.record.level.index.max.filegroup.count10000全局 RLI 的最大文件组数量。
hoodie.metadata.record.level.index.enablefalse启用分区级 RLI。与上述全局 RLI 的开关相互独立。
hoodie.metadata.record.level.index.min.filegroup.count1分区级 RLI 的最小文件组数量。
hoodie.metadata.record.level.index.max.filegroup.count10分区级 RLI 的最大文件组数量。
hoodie.metadata.record.level.index.defer.initfalse启用后,会将 RLI 的初始化推迟到新表的第二次提交,使 Hudi 能够根据实际记录量来确定文件组大小。该设置同时适用于全局 RLI 和分区级 RLI。
hoodie.metadata.record.index.max.filegroup.size1073741824(1 GB)单个 RLI 文件组的最大字节数。文件组越大,压实所需的时间越长。
# ensure that async indexing is enabled
hoodie.metadata.index.async=true
# enable column_stats index config
hoodie.metadata.index.column.stats.enable=true
# set concurrency mode and lock configs as this is a multi-writer scenario
# check https://hudi.apache.org/docs/concurrency_control/ for differnt lock provider configs
hoodie.write.concurrency.mode=optimistic_concurrency_control
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider
hoodie.write.lock.zookeeper.url=<zk_url>
hoodie.write.lock.zookeeper.port=<zk_port>
hoodie.write.lock.zookeeper.lock_key=<zk_key>
hoodie.write.lock.zookeeper.base_path=<zk_base_path>

计划索引任务

现在,我们可以使用 HoodieIndexer 的 schedule 模式来计划索引任务,如下所示:

spark-submit \
--jars "packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar,packaging/hudi-spark-bundle/target/hudi-spark3.5-bundle_2.12-1.2.0.jar" \
--class org.apache.hudi.utilities.HoodieIndexer \
/Users/home/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar \
--props /Users/home/path/to/indexer.properties \
--mode schedule \
--base-path /tmp/hudi-ny-taxi \
--table-name ny_hudi_tbl \
--index-types COLUMN_STATS \
--parallelism 1 \
--spark-memory 1g

这会向 timeline 写入一个 indexing.requested instant。

执行索引

要执行索引,请以 execute 模式运行索引器,如下所示。

spark-submit \
--jars "packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar,packaging/hudi-spark-bundle/target/hudi-spark3.5-bundle_2.12-1.2.0.jar" \
--class org.apache.hudi.utilities.HoodieIndexer \
/Users/home/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar \
--props /Users/home/path/to/indexer.properties \
--mode execute \
--base-path /tmp/hudi-ny-taxi \
--table-name ny_hudi_tbl \
--index-types COLUMN_STATS \
--parallelism 1 \
--spark-memory 1g

我们也可以以 scheduleAndExecute 模式运行索引器,一次性完成上述两个步骤。分开执行则能让我们更好地掌控执行时机。

下面我们来看数据时间线。

ls -lrt /tmp/hudi-ny-taxi/.hoodie
total 1816
-rw-r--r--  1 sagars  wheel       0 Apr 14 19:53 20220414195327683.commit.requested
-rw-r--r--  1 sagars  wheel  153423 Apr 14 19:54 20220414195327683.inflight
-rw-r--r--  1 sagars  wheel  207061 Apr 14 19:54 20220414195327683.commit
-rw-r--r--  1 sagars  wheel       0 Apr 14 19:54 20220414195423420.commit.requested
-rw-r--r--  1 sagars  wheel     659 Apr 14 19:54 20220414195437837.indexing.requested
-rw-r--r--  1 sagars  wheel  323950 Apr 14 19:54 20220414195423420.inflight
-rw-r--r--  1 sagars  wheel       0 Apr 14 19:55 20220414195437837.indexing.inflight
-rw-r--r--  1 sagars  wheel  222920 Apr 14 19:55 20220414195423420.commit
-rw-r--r--  1 sagars  wheel     734 Apr 14 19:55 hoodie.properties
-rw-r--r--  1 sagars  wheel     979 Apr 14 19:55 20220414195437837.indexing

在数据时间线中,我们可以看到索引操作是在一次提交完成(20220414195327683.commit)之后、另一次提交被请求(20220414195423420.commit.requested)之后才被调度的。因此它会选取 20220414195327683 作为基准即时(base instant)。索引操作处于进行中状态,同时还有一个进行中的写入者。如果解析索引器的日志,就会发现它在为基准即时完成索引之后,确实追上了即时 20220414195423420。

22/04/14 19:55:22 INFO HoodieTableMetaClient: Finished Loading Table of type MERGE_ON_READ(version=1, baseFileFormat=HFILE) from /tmp/hudi-ny-taxi/.hoodie/metadata
22/04/14 19:55:22 INFO RunIndexActionExecutor: Starting Index Building with base instant: 20220414195327683
22/04/14 19:55:22 INFO HoodieBackedTableMetadataWriter: Creating a new metadata index for partition 'column_stats' under path /tmp/hudi-ny-taxi/.hoodie/metadata upto instant 20220414195327683
...
...
22/04/14 19:55:38 INFO RunIndexActionExecutor: Total remaining instants to index: 1
22/04/14 19:55:38 INFO HoodieTableMetaClient: Loading HoodieTableMetaClient from /tmp/hudi-ny-taxi/.hoodie/metadata
22/04/14 19:55:38 INFO HoodieTableConfig: Loading table properties from /tmp/hudi-ny-taxi/.hoodie/metadata/.hoodie/hoodie.properties
22/04/14 19:55:38 INFO HoodieTableMetaClient: Finished Loading Table of type MERGE_ON_READ(version=1, baseFileFormat=HFILE) from /tmp/hudi-ny-taxi/.hoodie/metadata
22/04/14 19:55:38 INFO HoodieActiveTimeline: Loaded instants upto : Option{val=[20220414195423420__deltacommit__COMPLETED]}
22/04/14 19:55:38 INFO RunIndexActionExecutor: Starting index catchup task
...

删除索引

要删除索引,只需以 dropindex 模式运行索引即可。请注意,从 Hudi 1.2.0 开始,当某个索引在写入配置中被禁用时,Hudi 默认会自动删除其元数据表分区;参见 hoodie.metadata.auto.delete.partitions 以控制此行为。

spark-submit \
--jars "packaging/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar,packaging/hudi-spark-bundle/target/hudi-spark3.5-bundle_2.12-1.2.0.jar" \
--class org.apache.hudi.utilities.HoodieIndexer \
/Users/home/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar \
--props /Users/home/path/to/indexer.properties \
--mode dropindex \
--base-path /tmp/hudi-ny-taxi \
--table-name ny_hudi_tbl \
--index-types COLUMN_STATS \
--parallelism 1 \
--spark-memory 2g

注意事项

异步索引功能仍在不断完善中。在运行索引器时,从部署角度需要注意以下几点:

  • 只要启用了元数据表,文件索引就会默认创建。
  • 一次只为一个元数据分区(或一种索引类型)触发索引。
  • 如果通过异步索引启用了某个索引,那么务必在对应常规写入 writer 的配置中也启用该索引。否则,元数据 writer 会认为该索引已被禁用,进而清理相应的元数据分区。

这些限制中的一部分将在后续版本中移除。请关注该 GitHub issue以了解此功能的最新进展。

视频

评论

登录后参与评论

正在加载评论…