表服务

聚簇

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

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

Apache Hudi 为大数据带来了流式处理能力,在提供新鲜数据的同时,其效率比传统批处理高出一个数量级。在数据湖/数据仓库中,关键的权衡之一在于摄取速度与查询性能之间。数据摄取通常偏好小文件,以提高并行度并让数据尽快可供查询。然而,大量小文件会严重降低查询性能。此外,在摄取过程中,数据通常按到达时间进行共置。但当经常被查询的数据共置在一起时,查询引擎的性能会更好。在大多数架构中,各个系统往往各自独立地添加优化措施以提升性能,但由于数据布局未得到优化,这些优化会遇到瓶颈。本文档介绍一种名为 clustering(聚类)的新型表服务 [RFC-19],它能够在不影响摄取速度的前提下,通过重组数据来提升查询性能。

压缩(Compaction)与聚类(Clustering)有何区别?

Hudi 被建模为一种日志结构存储引擎,其中保存着数据的多个版本。具体而言,Hudi 中的 Merge-on-Read 表采用列式格式的基文件与包含更新的行式增量日志相结合的方式来存储数据。压缩是将增量日志与基文件合并,以生成包含最新数据快照的最新文件切片的一种方式。压缩有助于控制查询性能(增量日志文件过大会导致查询端的合并时间变长)。而聚类则是一种数据布局优化技术。通过聚类,可以将小文件拼接成更大的文件。此外,还可以按键(sort key)对数据进行聚类,使查询能够利用数据局部性。

聚类架构

从宏观层面来看,Hudi 通过其写入客户端 API 提供 insert/upsert/bulk_insert 等不同操作,以便将数据写入 Hudi 表。为了能够在文件大小与摄取速度之间进行权衡,Hudi 提供了 hoodie.parquet.small.file.limit 参数,用于配置允许的最小文件大小。用户可以将小文件软限制配置为 0,以强制新数据写入新的一组文件组(filegroup);也可以将其设置为更大的值,以确保新数据被"填充"到现有文件中,直到达到该限制,但这会增加摄取延迟。

为了支持一种既能快速摄取又不牺牲查询性能的架构,我们引入了"聚类"服务,通过重写数据来优化 Hudi 数据湖的文件布局。

聚类表服务可以异步或同步运行,并新增一种名为 “REPLACE” 的动作类型,用于在 Hudi 元数据时间线中标记聚类操作。

总体上,聚类包含 2 个步骤

  1. 调度聚类:使用可插拔的聚类策略创建聚类计划。
  2. 执行聚类:使用执行策略处理该计划,以创建新文件并替换旧文件。

调度聚类

调度聚类需遵循以下步骤。

  1. 识别适合聚类的文件:根据所选的聚类策略,调度逻辑会识别出适合聚类的文件。
  2. 根据特定条件对适合聚类的文件进行分组。每组的数据大小预期应为 targetFileSize 的整数倍。分组是计划中定义的 strategy(策略)的一部分。此外,还可以对分组大小设置上限,以提高并行度并避免大量数据的 shuffle。
  3. 最后,聚类计划以 Avro 元数据格式保存到时间线中。

增量调度

Hudi 支持对聚类操作进行增量调度,这可以显著提升分区数量较多的表上的性能。增量调度不会在每次聚类调度运行时扫描所有分区,而只处理自上次完成的聚类操作以来发生变化的分区。

该功能默认通过 hoodie.table.services.incremental.enabled 启用。启用后,聚类调度将会:

  1. 通过分析上次完成的聚类与当前调度即时(instant)之间时间窗口内的提交元数据,识别出自上次完成的聚类操作以来被修改过的分区
  2. 纳入之前调度运行中标记为缺失的任何分区(例如,由于 IO 限制或分组大小限制)
  3. 仅扫描并处理这些增量分区
  4. 如果找不到上次完成的聚类即时(例如,由于归档),或在获取增量分区时发生异常,则回退为扫描所有分区

对于分区较多的表,此优化可以大幅降低调度开销,并提升作业的整体稳定性。

执行聚类

  1. 读取聚类计划,获取标记了需要聚类的文件分组的 clusteringGroups。
  2. 对每个分组,我们使用 strategyParams(示例:sortColumns)实例化相应的策略类,并应用该策略来重写数据。
  3. 创建一个 “REPLACE” 提交,并在 HoodieReplaceCommitMetadata 中更新元数据。

Clustering Service 建立在 Hudi 基于 MVCC 的设计之上,使得写入器可以在后台执行 Clustering 操作、重新组织数据布局的同时,继续插入新数据,从而确保并发读写之间的快照隔离。

注意:Clustering 只能针对没有并发更新的表/分区进行调度。未来也会支持并发更新的使用场景。

Clustering 示例 图:展示通过 Clustering 提升查询性能

Clustering 使用场景

合并小文件

正如简介中提到的,流式数据摄入通常会在数据湖中产生较小的文件。而大量这样的小文件会导致较高的查询延迟。根据我们支持社区用户的经验,有不少用户使用 Hudi 就是为了其小文件处理能力。因此,你可以使用 Clustering 将大量这样的小文件合并成更大的文件。

合并小文件

按排序键进行 Clustering

数据湖中的另一个典型问题是数据到达时间与事件时间的矛盾。通常你按照到达时间写入数据,而这与查询谓词并不匹配。通过 Clustering,你可以基于查询谓词对数据排序并重写,这样数据跳过(data skipping)会非常高效,查询也能忽略大量不必要的数据扫描。

合并小文件

Clustering 策略

从宏观上看,Clustering 会基于可配置的策略生成计划,根据特定条件对符合条件的文件进行分组,然后执行该计划。如前所述,Clustering 计划及其执行都依赖于可配置的策略。这些策略大致可分为三类:Clustering 计划策略、执行策略和更新策略。

计划策略

该策略在创建 Clustering 计划时发挥作用。它用于决定哪些文件分组应当参与 Clustering,以及 Clustering 应当产生多少个输出文件分组。请注意,这些策略可以通过配置项 hoodie.clustering.plan.strategy.class 轻松插拔。

不同的计划策略如下所示:

基于大小的 Clustering 策略

该策略根据每组允许的最大大小来创建聚类分组。同时,它会将大于小文件限制的文件排除在聚类计划之外。根据写入客户端的不同,可用的策略有:SparkSizeBasedClusteringPlanStrategy、FlinkSizeBasedClusteringPlanStrategy 和 JavaSizeBasedClusteringPlanStrategy。此外,Hudi 提供了灵活性,可以包含或排除参与聚类的分区、调整文件大小限制以及最大输出分组数。更多详情请参阅 hoodie.clustering.plan.strategy.small.file.limit、hoodie.clustering.plan.strategy.max.num.groups、hoodie.clustering.plan.strategy.max.bytes.per.group 和 hoodie.clustering.plan.strategy.target.file.max.bytes。

配置名称 默认值 描述
hoodie.clustering.plan.strategy.partition.selected N/A (必填) 以逗号分隔的、需要运行聚类的分区列表

Config Param: PARTITION_SELECTED
Since Version: 0.11.0
hoodie.clustering.plan.strategy.partition.regex.pattern N/A (必填) 过滤与正则表达式模式匹配的聚类分区

Config Param: PARTITION_REGEX_PATTERN
Since Version: 0.11.0
hoodie.clustering.plan.partition.filter.mode NONE(可选) 创建聚类计划时使用的分区过滤模式。可选值:

- NONE:不过滤分区。聚类计划将包含所有存在聚类候选者的分区。
- RECENT_DAYS:此过滤器假定你的数据是按日期分区的。聚类计划将仅包含从 K 天前到 N 天前的分区,其中 K >= N。K 由 hoodie.clustering.plan.strategy.daybased.lookback.partitions 决定,N 由 hoodie.clustering.plan.strategy.daybased.skipfromlatest.partitions 决定。
- SELECTED_PARTITIONS:聚类计划将仅包含名称排序后落在闭区间 [hoodie.clustering.plan.strategy.cluster.begin.partition, hoodie.clustering.plan.strategy.cluster.end.partition] 内的分区路径。
- DAY_ROLLING:为了确定聚类计划中的分区,合格的分区将按升序排序。每个分区在该列表中都有一个索引 i。聚类计划将仅包含满足 i mod 24 = H 的分区,其中 H 是当前的小时数(从 0 到 23)。

Config Param: PLAN_PARTITION_FILTER_MODE_NAME
Since Version: 0.11.0

SparkSingleFileSortPlanStrategy

在此策略中,每个分区的聚簇分组的构建方式与 SparkSizeBasedClusteringPlanStrategy 相同。不同之处在于输出分组为 1 个,且文件组 ID 保持不变,而 SparkSizeBasedClusteringPlanStrategy 可以创建多个具有新 fileId 的文件组。

SparkConsistentBucketClusteringPlanStrategy

该策略专用于一致性桶索引(consistent bucket index)。它可用于扩展您的桶索引(从静态分区变为动态分区)。通常情况下,用户无需使用此策略。Hudi 内部使用它来动态扩展桶索引数据集的桶。

后两种策略仅适用于 Spark 引擎。

CommitBasedClusteringPlanStrategy

Hudi 1.2.0 引入了 org.apache.hudi.table.action.cluster.strategy.CommitBasedClusteringPlanStrategy,这是一种基于提交(commit)模式而非仅基于文件大小来调度聚簇的规划策略。它按照产生文件切片的提交对文件切片进行分组,从而更方便地对在特定时间窗口内写入的数据或满足特定提交条件的数据进行聚簇。

配置名称默认值描述
hoodie.clustering.plan.strategy.classSparkSizeBasedClusteringPlanStrategy设置为 org.apache.hudi.table.action.cluster.strategy.CommitBasedClusteringPlanStrategy 即可使用基于提交的规划。
hoodie.clustering.plan.strategy.earliest.commit.to.cluster(无)开始聚簇的最早提交时间(不包含该时间点)。仅考虑此时刻之后的提交。可用于对新数据进行增量聚簇,同时跳过已完成聚簇的历史数据。

SparkStreamCopyClusteringPlanStrategy

自 Hudi 1.2.0 起可用,org.apache.hudi.client.clustering.plan.strategy.SparkStreamCopyClusteringPlanStrategy 是一种仅适用于 Spark 的规划策略,它执行二进制文件拼接(字节级复制),而不是重新读取并重写记录。当目标仅仅是合并小文件且不需要排序时,这种方式可能显著更快。它与 org.apache.hudi.client.clustering.run.strategy.SparkStreamCopyClusteringExecutionStrategy 配合使用。

单分组聚簇控制

配置名称默认值描述
hoodie.clustering.plan.strategy.single.group.clustering.enabledtrue当只有一个文件组符合聚合条件时,是否仍然生成聚合计划。设置为 false,则在没有值得合并的内容时跳过聚合(即该分区已只有一个文件组)。

聚合计划生成中的文件切片排序顺序

从 1.2.0 开始,文件切片打包进聚合组的顺序变为可配置,从而可以更精细地控制哪些文件被归入同一组以及各组的填充方式。

配置名称默认值描述
hoodie.clustering.plan.strategy.file.slices.sort.bySIZE以逗号分隔的字段列表,用于在分区内将文件切片打包进聚合组时对文件切片进行排序。SIZE:按文件大小降序排序(最大在前)。INSTANT_TIME:按提交时间升序排序(最旧的文件在前)。示例:INSTANT_TIME,SIZE 表示先按提交时间排序,再按大小排序。

驱动端计划生成

配置名称默认值描述
hoodie.clustering.plan.generation.use.local.engine.contextfalse启用后,聚类分组的计算将在 driver(本地引擎上下文)上运行,而不是分布到各个 executor 上执行。当只有少数分区但包含大量文件时启用该配置,此时在 driver 本地进行计算比分配 executor 资源更为高效。

执行策略

在规划阶段构建好聚类分组之后,Hudi 会针对每个分组应用执行策略,主要依据排序列和文件大小。策略可通过配置 hoodie.clustering.execution.strategy.class 指定。默认情况下,Hudi 会按指定列对计划中的文件组排序,同时满足配置的目标文件大小。

配置名称默认值描述

hoodie.clustering.execution.strategy.classorg.apache.hudi.client.clustering.run.strategy.SparkSortAndSizeExecutionStrategy (可选)用于提供策略类(RunClusteringStrategy 的子类)的配置,以定义聚类计划的执行方式。默认情况下,我们按指定列对计划中的文件组排序,同时满足配置的目标文件大小。

Config Param: EXECUTION_STRATEGY_CLASS_NAME
Since Version: 0.7.0

可用的策略如下:

  1. SPARK_SORT_AND_SIZE_EXECUTION_STRATEGY:使用 bulk_insert 从输入文件组重新写入数据。

  2. 将 hoodie.clustering.execution.strategy.class 设置为 org.apache.hudi.client.clustering.run.strategy.SparkSortAndSizeExecutionStrategy。

  3. hoodie.clustering.plan.strategy.sort.columns:聚类时用于对数据排序的列。该配置需与布局优化策略配合使用,具体取决于你的查询谓词。可以在该配置中设置以逗号分隔的待排序列列表。

  4. JAVA_SORT_AND_SIZE_EXECUTION_STRATEGY:与 SPARK_SORT_AND_SIZE_EXECUTION_STRATEGY 类似,适用于 Java 和 Flink 引擎。将 hoodie.clustering.execution.strategy.class 设置为 org.apache.hudi.client.clustering.run.strategy.JavaSortAndSizeExecutionStrategy。

  5. SPARK_CONSISTENT_BUCKET_EXECUTION_STRATEGY:顾名思义,该策略适用于动态扩展一致分桶索引,且仅适用于 Spark 引擎。将 hoodie.clustering.execution.strategy.class 设置为 org.apache.hudi.client.clustering.run.strategy.SparkConsistentBucketClusteringExecutionStrategy。

更新策略

目前,聚类只能调度到没有接收到任何并发更新的表或分区上。默认情况下,更新策略配置 hoodie.clustering.updates.strategy 被设置为 SparkRejectUpdateStrategy。如果在聚类过程中某个文件组发生了更新,该策略将拒绝这些更新并抛出异常。然而,在某些用例中,更新非常稀疏,并不会涉及大多数文件组。默认的简单拒绝更新的策略似乎并不合理。在这类用例中,用户可以将该配置设置为 SparkAllowUpdateStrategy。

以上我们讨论了关键的策略配置。所有与聚类相关的其他配置列于 聚类配置。在这份列表中,以下几个配置对内联(inline)或异步(async)聚类非常有用,下面通过示例加以说明。

内联聚类

内联聚类与常规的写入器同步执行,或者作为数据摄入管道的一部分执行。这意味着在聚类完成之前,下一轮摄入无法继续进行。通过内联聚类,Hudi 会在每次提交完成后调度并规划聚类操作,并在聚类计划生成后立即执行。这是最简单的部署模型,因为与运行多个异步 Spark 作业相比,它更易于管理。该模式在 Spark Datasource、Flink、Spark-SQL 以及 Hudi Streamer 的 sync-once 模式下均受支持。

对于这种部署模式,请启用并设置:hoodie.clustering.inline

若要选择聚类的触发频率,还需设置:hoodie.clustering.inline.max.commits。

内联聚类可以很方便地通过 Spark DataFrame 选项进行配置,示例如下:

import org.apache.hudi.QuickstartUtils._
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
val df =  //generate data frame
df.write.format("org.apache.hudi").
        options(getQuickstartWriteConfigs).
        option("hoodie.table.ordering.fields", "ts").
        option("hoodie.datasource.write.recordkey.field", "uuid").
        option("hoodie.datasource.write.partitionpath.field", "partitionpath").
        option("hoodie.table.name", "tableName").
        option("hoodie.parquet.small.file.limit", "0").
        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"). //optional, if sorting is needed as part of rewriting data
        mode(Append).
        save("dfs://location");

异步 Clustering

异步 clustering 在后台运行 clustering 表服务,不会阻塞常规的写入(ingestion)流程。异步 clustering 进程有三种不同的部署方式:

  • 同一进程内异步执行:在这种部署模式下,Hudi 会在每次提交完成后、作为写入流水线的一部分来调度和规划 clustering 操作。同时,Hudi 会在同一作业内另起一个线程来执行 clustering 表服务。Spark Streaming、Flink 以及处于连续模式的 Hudi Streamer 均支持这种方式。对于该部署模式,请启用 hoodie.clustering.async.enabled 和 hoodie.clustering.async.max.commits​。
  • 由独立进程异步调度与执行:在这种部署模式下,应用会作为写入流水线的一部分将数据写入 Hudi 表;另有一个独立的 clustering 作业来调度、规划并执行 clustering 操作。通过为 clustering 操作运行独立的作业,可以重新平衡 Hudi 对计算资源的使用:写入所需的计算资源更少,从而使写入延迟更加稳定,同时为 clustering 进程预留一组独立的计算资源。请为所有作业(写入作业和表服务作业)之间的并发控制配置锁提供者(lock provider)。一般而言,当存在两个不同的作业或两个不同的进程时,就需要配置锁提供者。所有写入方式均支持这种部署模式。对于该部署模式,写入端不应配置任何 clustering 相关参数。
  • 内联调度、异步执行:在这种部署模式下,一个作业负责应用的数据写入和 clustering 调度;另一个作业负责执行 clustering 计划。所支持的写入方式(见下文)不会被数据写入阻塞。如果启用了元数据表,则不需要锁提供者。但是,如果启用了元数据表,请确保所有作业都已配置锁提供者以进行并发控制。所有写入方式均支持该部署选项。对于该部署模式,请启用 hoodie.clustering.schedule.inline 和 hoodie.clustering.async.enabled。

Hudi 支持多写入者(multi-writers),它在多个表服务之间提供快照隔离,从而允许写入者在 clustering 于后台运行的同时继续进行写入。

配置名称 默认值 描述
hoodie.clustering.async.enabled false(可选) 允许在表上发生插入操作时异步运行 clustering 服务。
Config Param: ASYNC_CLUSTERING_ENABLE
Since Version: 0.7.0
hoodie.clustering.async.max.commits 4(可选) 控制异步 clustering 频率的配置。
Config Param: ASYNC_CLUSTERING_MAX_COMMITS
Since Version: 0.9.0

设置异步 Clustering

用户可以借助 HoodieClusteringJob 来设置两步式异步 Clustering。

HoodieClusteringJob

通过指定 scheduleAndExecute 模式,可以在同一个步骤中同时完成调度和 Clustering。可以使用 -mode 或 -m 选项指定相应的模式。共有三种模式:

  1. schedule:生成 Clustering 计划。该模式会返回一个 instant,可作为参数传入 execute 模式。
  2. execute:在指定的 instant 上执行 Clustering 计划。如果未指定 instant-time,HoodieClusteringJob 将在 Hudi 时间线上最早的 instant 上执行。
  3. scheduleAndExecute:先生成 Clustering 计划,然后立即执行该计划。

可用选项

除了基础的模式选项外,HoodieClusteringJob 还支持以下重试和超时选项(在 scheduleAndExecute 模式下生效):

选项名称简写默认值描述
--retry-last-failed-job-rcfalse设置为 true 时,会检查并回滚上一次失败的 Clustering 计划并执行它,而不是直接规划新的 Clustering 作业。这对于从之前的失败中恢复非常有用。
--job-max-processing-time-ms-jt0认定 Clustering 作业失败前的最大处理时间(毫秒)。若超过该时间作业仍未完成,Hudi 将视为作业失败并重新启动它(需与 --retry-last-failed-job 搭配使用)。值为 0 或负数时将禁用超时检查。

note

这些重试选项仅在使用 --mode scheduleAndExecute 时生效。--retry-last-failed-job 选项需要将 --job-max-processing-time-ms 设置为正值,才能检测到处于停滞状态的 in-flight instant。

过期 Clustering Instant 的自动清理

当一个聚合作业(clustering job)已被调度但从未成功执行时(例如由于 driver 故障),处于 inflight 状态的 replacecommit 实例会阻塞后续的聚合作业运行。Hudi 1.2.0 新增了对这类过期聚合作用实例的自动过期机制,作为上述手动重试选项的补充。

note

过期聚合作用计划的清理需要设置 hoodie.clean.failed.writes.policy=LAZY。在 LAZY 清理模式下,失败写入的回滚(由下一次写入触发)也会一并回滚已过期的聚合作用实例。

配置名称默认值描述
hoodie.clustering.enable.expirationsfalse启用后,在 LAZY 清理模式下,失败写入的回滚会同时回滚心跳已过期的聚合作用 replacecommit 实例。聚合作用作业会在调度前记录心跳,以便其他写入方检测到过期的尝试。
hoodie.clustering.expiration.threshold.mins60聚合作用实例只有在其创建时间达到该分钟数后才会被视为过期。该配置作为保护措施,避免回滚仍在进行中的聚合作用尝试。

请注意,若要在原写入方仍在运行时执行此作业,请启用多写入(multi-writing):

hoodie.write.concurrency.mode=optimistic_concurrency_control
hoodie.write.lock.provider=org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider

以下是一个用于配置 HoodieClusteringJob 的 spark-submit 命令示例:

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.HoodieClusteringJob \
/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar \
--props /path/to/config/clusteringjob.properties \
--mode scheduleAndExecute \
--base-path /path/to/hudi_table/basePath \
--table-name hudi_table_schedule_clustering \
--spark-memory 1g

clusteringjob.properties 配置文件示例:

hoodie.clustering.async.enabled=true
hoodie.clustering.async.max.commits=4
hoodie.clustering.plan.strategy.target.file.max.bytes=1073741824
hoodie.clustering.plan.strategy.small.file.limit=629145600
hoodie.clustering.execution.strategy.class=org.apache.hudi.client.clustering.run.strategy.SparkSortAndSizeExecutionStrategy
hoodie.clustering.plan.strategy.sort.columns=column1,column2

Hudi Streamer

这就引出了 Hudi 用户最喜爱的工具。现在,我们可以通过 Hudi Streamer 触发异步聚类。只需将 hoodie.clustering.async.enabled 配置设置为 true,并在属性文件中指定其他聚类相关配置——该属性文件的位置可在启动 Hudi Streamer 时通过 —props 参数传入(与 HoodieClusteringJob 的情况类似)。

下面是一个用于配置 Hudi Streamer 的 spark-submit 命令示例:

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 \
/path/to/hudi-utilities-slim-bundle/target/hudi-utilities-slim-bundle_2.12-1.2.0.jar \
--props /path/to/config/clustering_kafka.properties \
--schemaprovider-class org.apache.hudi.utilities.schema.SchemaRegistryProvider \
--source-class org.apache.hudi.utilities.sources.AvroKafkaSource \
--source-ordering-field impresssiontime \
--table-type COPY_ON_WRITE \
--target-base-path /path/to/hudi_table/basePath \
--target-table impressions_cow_cluster \
--op INSERT \
--hoodie-conf hoodie.clustering.async.enabled=true \
--continuous

Spark Structured Streaming

我们也可以通过 Spark Structured Streaming sink 启用异步聚类,如下所示。

val commonOpts = Map(
   "hoodie.insert.shuffle.parallelism" -> "4",
   "hoodie.upsert.shuffle.parallelism" -> "4",
   "hoodie.datasource.write.recordkey.field" -> "_row_key",
   "hoodie.datasource.write.partitionpath.field" -> "partition",
   "hoodie.table.ordering.fields" -> "timestamp",
   "hoodie.table.name" -> "hoodie_test"
)

def getAsyncClusteringOpts(isAsyncClustering: String,
                           clusteringNumCommit: String,
                           executionStrategy: String):Map[String, String] = {
   commonOpts + (DataSourceWriteOptions.ASYNC_CLUSTERING_ENABLE.key -> isAsyncClustering,
           HoodieClusteringConfig.ASYNC_CLUSTERING_MAX_COMMITS.key -> clusteringNumCommit,
           HoodieClusteringConfig.EXECUTION_STRATEGY_CLASS_NAME.key -> executionStrategy
   )
}

def initStreamingWriteFuture(hudiOptions: Map[String, String]): Future[Unit] = {
   val streamingInput = // define the source of streaming
   Future {
      println("streaming starting")
      streamingInput
              .writeStream
              .format("org.apache.hudi")
              .options(hudiOptions)
              .option("checkpointLocation", basePath + "/checkpoint")
              .mode(Append)
              .start()
              .awaitTermination(10000)
      println("streaming ends")
   }
}

def structuredStreamingWithClustering(): Unit = {
   val df = //generate data frame
   val hudiOptions = getClusteringOpts("true", "1", "org.apache.hudi.client.clustering.run.strategy.SparkSortAndSizeExecutionStrategy")
   val f1 = initStreamingWriteFuture(hudiOptions)
   Await.result(f1, Duration.Inf)
}

Flink 离线聚类

Flink 的离线聚类需要在命令行中以 Flink 作业的形式提交。程序入口位于 hudi-flink-bundle.jar 中的 org.apache.hudi.sink.clustering.HoodieFlinkClusteringJob。

# Command line
./bin/flink run -c org.apache.hudi.sink.clustering.HoodieFlinkClusteringJob lib/hudi-flink-bundle.jar --path hdfs://xxx:9000/table

选项

选项名称默认值描述
--pathn/a **(必填)**目标表在 Hudi 中的存储路径
--schedulefalse(可选)是否执行调度聚类计划的操作。在写入过程仍在进行时开启该参数存在数据丢失的风险。因此,开启该参数前必须确保当前没有写入任务正在向该表写入数据
--servicefalse(可选)是否启动一个监控服务,按配置的时间间隔检查并调度新的聚类任务。
--min-clustering-interval-seconds600(s)(可选)服务模式下的检查间隔,默认为 10 分钟。
--retry0(可选)聚类操作的重试次数。仅在单次运行模式(非服务模式)下有效。默认为 0(不重试)。
--retry-last-failed-jobfalse(可选)当进行中的 instant 超过最大处理时间时,检查并重试上一次失败的聚类作业。仅在单次运行模式下有效。需要将 --job-max-processing-time-ms 设置为正值。
--job-max-processing-time-ms0(可选)将聚类作业判定为失败前的最大处理时间(毫秒)。与 --retry-last-failed-job 搭配使用。默认值 0 表示不进行超时检查。

注意

重试相关选项(--retry、--retry-last-failed-job、--job-max-processing-time-ms)仅在单次运行模式下生效,在服务模式下无效。服务模式通过其持续监控循环来实现隐式的重试语义。如果启用了 --retry-last-failed-job,但 --job-max-processing-time-ms 未设置为正值,系统将记录一条警告。

Java 客户端

Clustering 同样支持通过 Java 客户端执行。开箱即用的规划策略为 org.apache.hudi.client.clustering.plan.strategy.JavaSizeBasedClusteringPlanStrategy,执行策略为 org.apache.hudi.client.clustering.run.strategy.JavaSortAndSizeExecutionStrategy。请注意,目前 Java 执行策略仅支持线性排序。

博客

Apache Hudi Z-Order 和希尔伯特空间填充曲线 Hudi Z-Order 和希尔伯特空间填充曲线

视频

评论

登录后参与评论

正在加载评论…