更新时间:2026-09-10 GMT+08:00
分享

Hudi Clustering操作说明

什么是Clustering

在Hudi数据湖中,随着数据持续写入,底层会产生大量小文件,导致查询引擎需要扫描更多文件、元数据管理开销增大,从而引起查询性能下降。同时,数据按写入顺序分布,缺乏基于查询模式的布局优化,导致数据局部性差,影响查询效率。

Clustering即数据布局优化服务,可重新组织数据以提高查询性能,也不会影响摄取速度。通过将小文件合并为大文件、按指定排序列重新组织数据分布,Clustering能够显著减少查询时需要扫描的文件数量,提升I/O效率。

典型应用场景:

  • 小文件过多:Hudi表经过多次upsert操作后产生大量小文件,查询性能下降,需要通过Clustering合并小文件。
  • 查询性能优化:查询条件涉及多个字段,通过Clustering按查询字段排序,提升数据局部性,减少全表扫描。
  • 数据布局优化:数据按时间或写入顺序分布,需要按业务字段重新组织以提升范围查询效率。

Clustering架构

Hudi通过其写入客户端API提供了不同的操作,如insert/upsert/bulk_insert来将数据写入Hudi表。为了能够在文件大小和入湖速度之间进行权衡,Hudi提供了一个hoodie.parquet.small.file.limit配置来设置最小文件大小。用户可以将该配置设置为“0”,以强制新数据写入新的文件组,或设置为更高的值以确保新数据被“填充”到现有小的文件组中,直到达到指定大小为止,但其会增加摄取延迟。

为能够支持快速摄取的同时不影响查询性能,引入了Clustering服务来重写数据以优化Hudi数据湖文件的布局。

Clustering服务可以异步或同步运行,Clustering会添加了一种新的REPLACE操作类型,该操作类型将在Hudi元数据时间轴中标记Clustering操作。

Clustering服务基于Hudi的MVCC设计,允许继续插入新数据,而Clustering操作在后台运行以重新格式化数据布局,从而确保并发读写者之间的快照隔离。

总体而言Clustering分为两个部分:

  • 调度Clustering:使用可插拔的Clustering策略创建Clustering计划。
    1. 识别符合Clustering条件的文件:根据所选的Clustering策略,调度逻辑将识别符合Clustering条件的文件。
    2. 根据特定条件对符合Clustering条件的文件进行分组。每个组的数据大小应为targetFileSize的倍数。分组是计划中定义的"策略"的一部分。此外还有一个选项可以限制组大小,以改善并行性并避免混排大量数据。
    3. 将Clustering计划以avro元数据格式保存到时间线。
  • 执行Clustering:使用执行策略处理计划以创建新文件并替换旧文件。
    1. 读取Clustering计划,并获得ClusteringGroups,其标记了需要进行Clustering的文件组。
    2. 对于每个组使用strategyParams实例化适当的策略类(例如:sortColumns),然后应用该策略重写数据。
    3. 创建一个REPLACE提交,并更新HoodieReplaceCommitMetadata中的元数据。

约束限制

  • Clustering的排序列不允许值存在null,是spark rdd的限制。
  • 当target.file.max.bytes的值较大时,启动Clustering执行需要提高--spark-memory,否则会导致executor内存溢出。
  • 当前clean不支持清理Clustering失败后的垃圾文件。
  • Clustering后可能出现新文件大小不等引起数据倾斜的情况。
  • cluster不支持和upsert并发。
  • 如果clustering处于inflight状态,该FileGroup下的文件不支持Update操作。
  • 如果存在未完成的Clustering计划,后续写入触发生成compaction调度计划时会报错失败,需要及时执行Clustering计划。

执行Clustering

  1. 同步执行Clustering配置。

    在写入数据时通过option方式配置Clustering参数,Clustering将在满足触发条件后自动执行。适用于写入过程中需要自动优化数据布局的场景。

    在写入时加上配置参数:

    option("hoodie.clustering.inline", "true").
    option("hoodie.clustering.inline.max.commits", "4").
    option("hoodie.clustering.plan.strategy.target.file.max.bytes", "1073741824").
    option("hoodie.clustering.plan.strategy.small.file.limit", "629145600").
    option("hoodie.clustering.plan.strategy.sort.columns", "column1,column2").
  2. 异步执行Clustering:

    异步执行Clustering需要先调度生成Clustering计划,再执行该计划。适用于不影响写入流程、需要在特定时间窗口执行Clustering的场景。

    • MRS 3.2.0及之后版本:

      通过Spark SQL命令来执行clustering,具体可以参考CLUSTERING章节。

    • MRS 3.1.2版本,执行如下命令:
      spark-submit --master yarn --class org.apache.hudi.utilities.HoodieClusteringJob /opt/client/Hudi/hudi/lib/hudi-utilities*.jar --schedule --base-path <table_path> --table-name <table_name> --props /tmp/clusteringjob.properties --spark-memory 1g
      spark-submit --master yarn --driver-memory 16G --executor-memory 12G --executor-cores 4 --num-executors 4 --class org.apache.hudi.utilities.HoodieClusteringJob /opt/client/Hudi/hudi/lib/hudi-utilities*.jar --base-path <table_path> --instant-time 20210605112954 --table-name <table_name> --props /tmp/clusteringjob.properties --spark-memory 12g
  3. 指定clustering的排序方式和排序列:

    当前clustering支持linear、z-order、hilbert三种排序方式,可以通过option方式或者set方式来设置。

    • linear:普通排序,默认排序,适合排序一个字段, 或者多个低级字段。
    • z-order和hilbert:多维排序,需要指定“hoodie.layout.optimize.strategy”为z-order或者hilbert。

      适合排序多个字段,例如查询条件中涉及到多个字段。推荐排序字段的个数2到4个。

      hilbert多维排序效果比z-order好,但是排序效率没z-order高。

验证Clustering执行结果:

  • 同步执行:写入数据达到hoodie.clustering.inline.max.commits指定的次数后,通过Hudi时间线查看是否生成了REPLACE提交。
  • 异步执行:在spark-submit执行完成后,查看Spark作业执行日志确认是否成功;通过Hudi时间线或Hudi CLI查看Clustering计划是否状态为completed。
  • 验证数据布局:通过查询Hudi表的文件列表,确认小文件是否已合并,文件大小是否接近target.file.max.bytes配置值。

常见问题

  • Clustering执行失败:请检查排序列是否存在null值(约束限制中说明排序列不允许null)、executor内存是否充足(target.file.max.bytes较大时需提高--spark-memory)、是否存在未完成的Clustering计划导致冲突。
  • Clustering后查询性能未提升:请确认排序列与查询条件匹配,排序字段个数建议2~4个,过多字段排序效果反而下降。可通过查看执行计划确认数据扫描范围是否减少。
  • Compaction调度报错:如果存在未完成的Clustering计划,后续compaction调度会报错失败,需先执行Clustering计划后再进行compaction。

相关文档

相关文档