- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
本篇文章聚焦于 Apache Hudi 仓库中hudi-spark-datasource模块,该模块是 Hudi 与 Apache Spark 集成的核心枢纽,为 Spark SQL 与 DataFrame API 提供读写 Hudi 表的 DataSource 能力。通过阅读本文,你将完整掌握该模块的子模块组成与分层设计思路、不同 Spark 版本对应的构建 Profile 与语言/JDK 要求、以及 Parquet 压缩编码随 Spark 版本变化的默认行为与潜在风险,并了解其关键能力(SQL 扩展、存储过程、时间旅行、增量查询、索引、流式写入与 CDC)在仓库源码中的落点。
模块定位:Hudi 与 Spark 的集成枢纽
hudi-spark-datasource是整个 Hudi 仓库中面向 Spark 引擎的一站式集成模块。根据 hudi-spark-datasource/README.md 的说明,该模块为使用 Spark SQL 和 DataFrame 读写 Hudi 表提供了完整的 DataSource API 支持。
从仓库目录结构看,该模块采用**聚合模块(aggregator pom)**的形态组织:父级 pom.xml 声明了hudi-spark-common与hudi-spark两个子模块,而根 pom.xml 中的各spark3.x/spark4.xProfile 则按需引入对应的版本特定适配模块(hudi-spark3.3.x、hudi-spark4.1.x等)。这种"聚合 + 按 Profile 激活"的组合方式,是理解整个模块体系的关键起点。
分层模块架构:最大化复用,最小化版本差异
该模块的核心设计思想是分层架构:把与 Spark 版本无关的通用逻辑下沉到公共模块,把版本相关的差异收敛到各自的适配器模块,从而在不同 Spark 版本之间最大化代码复用,同时保留版本特定的优化。
| 模块 | 职责定位 |
|---|---|
hudi-spark-common | 跨所有 Spark 版本共享的核心 Spark 集成代码,包含 DataSource V1/V2 实现、文件索引、SQL 写入器和增量读支持 |
hudi-spark3-common | Spark 3.x 版本共享代码,包含 Spark 3 适配器接口、DML 命令与分区映射 |
hudi-spark4-common | Spark 4.x 版本共享代码,包含 Spark 4 适配器接口及 4.x 特有实现 |
hudi-spark3.3.x | Spark 3.3.x 特定适配器实现,含版本相关的 SQL 解析器与文件读取器 |
hudi-spark3.4.x | Spark 3.4.x 特定适配器实现 |
hudi-spark3.5.x | Spark 3.5.x 特定适配器实现(默认) |
hudi-spark4.0.x | Spark 4.0.x 特定适配器实现 |
hudi-spark4.1.x | Spark 4.1.x 特定适配器实现 |
hudi-spark4.2.x | Spark 4.2.x 特定适配器实现 |
hudi-spark | 主 Spark DataSource 模块,包含 Spark Session 扩展、存储过程、SQL 解析器与逻辑计划 |
公共层:hudi-spark-common
hudi-spark-common是整个 Spark 集成的心脏,其源码位于 hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi。从源码结构看,该模块中的关键构件包括:
- DataSource 入口与关系实现:DefaultSource.scala、BaseDefaultSource.scala、HoodieBaseRelation.scala、HoodieHadoopFsRelationFactory.scala 等共同支撑 DataSource V1/V2 的读写路径;
- 写入器:HoodieSparkSqlWriter.scala 与 HoodieWriterUtils.scala 负责将 Spark DataFrame 高效落盘为 Hudi 表;
- 文件索引:HoodieFileIndex.scala、HoodieIncrementalFileIndex.scala、SparkHoodieTableFileIndex.scala 承担表文件定位与查询裁剪;
- 增量读:IncrementalRelationV1.scala、IncrementalRelationV2.scala、MergeOnReadIncrementalRelationV1.scala 等实现增量查询语义;
- 流式写入:HoodieStreamingSink.scala 提供 Structured Streaming Sink。
版本特定层:Spark 3.x / 4.x 适配器
hudi-spark3-common(Spark 3 适配器接口、DML 命令、分区映射)与hudi-spark4-common(Spark 4 适配器接口)负责屏蔽引擎 API 差异;hudi-spark3.3.x/hudi-spark3.4.x/hudi-spark3.5.x与hudi-spark4.0.x/hudi-spark4.1.x/hudi-spark4.2.x则各自携带版本特定的 SQL 解析器(g4 语法文件)与文件读取器实现。README 明确指出 Spark 3.5.x 是默认适配版本,仓库根 pom.xml 中spark3.5Profile 的activeByDefault配置与此一致。
主模块:hudi-spark
hudi-spark是最终面向用户打包的主模块,聚合了上述公共层与版本适配层。其 pom.xml 通过${hudi.spark.module}_${scala.binary.version}动态依赖对应版本的适配器,并引入hudi-spark-client、hudi-client-common、hudi-hive-sync等基础能力,同时使用 ANTLR 编译 Hudi 自有 SQL 语法(src/main/antlr4/)。
Spark 版本支持矩阵与构建 Profile
README 给出了官方的版本支持矩阵,这也是选择 Hudi 构建参数时的权威依据:
| Spark 版本 | 模块 | Scala 版本 | Java 版本 | 构建 Profile |
|---|---|---|---|---|
| 3.3.x | hudi-spark3.3.x | 2.12 | 11+ | -Dspark3.3 |
| 3.4.x | hudi-spark3.4.x | 2.12 | 11+ | -Dspark3.4 |
| 3.5.x(默认) | hudi-spark3.5.x | 2.12、2.13 | 11+ | -Dspark3.5 |
| 4.0.x | hudi-spark4.0.x | 2.13 | 17+ | -Dspark4.0 |
| 4.1.x | hudi-spark4.1.x | 2.13 | 17+ | -Dspark4.1 |
| 4.2.x | hudi-spark4.2.x | 2.13 | 17+ | -Dspark4.2 |
对应地,根 pom.xml 中定义了spark3.3、spark3.4、spark3.5、spark4.0、spark4.1、spark4.2等构建 Profile(spark3.5为默认激活),并在各 Profile 内绑定该版本配套的生态组件版本。例如:
spark3.3Profile 使用parquet.version1.12.2、avro.version1.11.4、Scala 2.12.15,并跳过 Lance / Vortex 测试;spark3.5Profile 使用parquet.version1.13.1、Scala 2.12.18 与 2.13.8,启用 Lance / Vortex 测试;spark4.1Profile 使用parquet.version1.16.0、Scala 2.13.17、hadoop.version3.4.2。
这意味着你需要根据实际部署的 Spark 版本选择对应 Profile 进行构建,例如mvn clean install -Dspark3.5(默认即可省略该参数)或mvn clean install -Dspark4.1。若要在 Spark 3.x 上同时覆盖 2.12 与 2.13 两个 Scala 二进制版本,则需分别按各自 Profile 构建产物(hudi-spark_2.12与hudi-spark_2.13,见 hudi-spark/pom.xml 的 artifactId 定义)。
Parquet 压缩编码兼容性:ZSTD 与 GZIP 的版本差异
这是一个在部署 Hudi + Spark 时容易被忽视、但影响内存安全的重要细节。README 明确给出了以下结论:
- Hudi 将
hoodie.parquet.compression.codec的默认值设置为:Spark 3.5 及更新版本为 ZSTD,Spark 3.3 与 3.4 为 GZIP。 - Spark 3.3 / 3.4 的非向量化 Hudi file-group 读取器使用 parquet-java 1.12.x,该版本受PARQUET-2160问题影响,在读取 ZSTD 文件时可能发生堆外内存泄漏。
- 因此:在使用 ZSTD 之前,应先升级到 Spark 3.5 或更新版本。
- 旧版本 Spark 上的 GZIP 写入默认值并不能消除读取"其他引擎产生的 ZSTD 文件"的风险——换言之,文件是哪个引擎写的不重要,关键是读取时所用 Spark 版本的 parquet-java 是否受影响。
- 显式设置
hoodie.parquet.compression.codec会覆盖版本相关的默认值,这是规避该风险最直接的手段。
该行为在源码层有直接佐证:HoodieStorageConfig.java 中PARQUET_COMPRESSION_CODEC_NAME的配置声明写明默认值为 "ZSTD for Flink and Spark 3.5 or newer; GZIP for Java and older Spark versions",并重复了 PARQUET-2160 的警告与"显式配置值总是优先于引擎默认值"的规则。同文件还提供了配套的 ZSTD 级别控制项hoodie.logfile.parquet.compression.codec.zstd.level(默认 1),用于调节原生 Parquet 日志文件的 Zstandard 压缩级别。
关键能力全景与源码落点
README 将hudi-spark-datasource的关键能力归纳为八项,它们并非抽象承诺,而是在仓库源码中均有具体实现:
DataSource V1/V2 支持
完整的 Spark DataSource 集成由hudi-spark-common的 DefaultSource.scala 等入口类驱动,配合 HoodieBaseRelation.scala 与 HoodieHadoopFsRelationFactory.scala,同时覆盖 V1 与 V2 两条 API 路径。
Spark SQL 原生集成
HoodieSparkSessionExtension.scala 通过SparkSessionExtensions机制注入自定义能力:injectParser引入 HoodieCommonSqlParser 解析 Hudi 专属 SQL 语法,injectResolutionRule/injectPostHocResolutionRule/injectOptimizerRule挂载 HoodieAnalysis 中的解析与优化规则,最后通过sparkAdapter.injectTableFunctions / injectScalarFunctions / injectPlannerStrategies注入表函数、标量函数与执行策略。使用时只需在启动时配置spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension。
存储过程(Stored Procedures)
hudi-spark内置了丰富的表管理与运维存储过程,全部位于 hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures 目录下,涵盖:
- 查询类:ShowCommitsProcedure.scala、ShowTimelineProcedure.scala、ShowFsPathDetailProcedure.scala、ShowHoodieLogFileRecordsProcedure.scala;
- 执行类:RunCompactionProcedure.scala、RunClusteringProcedure.scala、RunCleanProcedure.scala、RunBootstrapProcedure.scala;
- 运维/修复类:RestoreToInstantProcedure.scala、RollbackToInstantTimeProcedure.scala、CreateSavepointProcedure.scala、RepairOrphanFilesProcedure.scala 等。
这些过程由 HoodieProcedures.scala 统一注册,通过CALL <procedure_name>(...)语法调用,是 Spark SQL 侧做表运维的标准入口。
时间旅行(Time Travel)
通过读取历史版本数据满足回放与审计需求,hudi-spark-common中的 TimelineRelation.scala 等实现支撑了该能力。
增量查询(Incremental Queries)
增量读(Incremental Query)用于高效消费自某个 commit 以来的变更数据,是实现变更数据捕获类管道的基础。IncrementalRelationV1.scala、IncrementalRelationV2.scala 与 MergeOnReadIncrementalRelationV1.scala 分别处理 COW / MOR 表的增量读,配合 HoodieIncrementalFileIndex.scala 定位增量文件。
索引支持(Index Support)
hudi-spark-common集中实现了多种索引能力的 Spark 侧支持,包括布隆过滤器 BloomFiltersIndexSupport.scala、列统计 ColumnStatsIndexSupport.scala、记录级索引 RecordLevelIndexSupport.scala(含分区版本 PartitionedRecordLevelIndexSupport.scala 与全局版本 GlobalRecordLevelIndexSupport.scala)、分区统计 PartitionStatsIndexSupport.scala、表达式索引 ExpressionIndexSupport.scala、桶索引 BucketIndexSupport.scala 以及辅助索引 SecondaryIndexSupport.scala 等,这些实现共同支撑查询阶段的文件与分区裁剪。
流式写入(Streaming Support)
HoodieStreamingSink.scala 将 Hudi 接入 Spark Structured Streaming 的 Sink 体系,持续接收微批数据写入 Hudi 表;从源码看,其内部会强制MarkerType.DIRECT、启用AUTO_ADJUST_LOCK_CONFIGS并注入批次号(SPARK_STREAMING_BATCH_ID),以保证流式写入下的并发安全与幂等。
CDC 支持
hudi-spark-common源码目录中专门的 cdc 子目录承载 CDC 读取相关实现,用于追踪行级变更,可配合增量查询构建实时数仓的变更流。
仓库内的实战示例与运行脚本
hudi-spark-datasource并非只有源码,仓库还提供了可运行的示例与脚本,便于快速验证:
- Spark SQL 批量查询:docker/demo/sparksql-batch1.commands 演示了 COW、MOR(
_ro与_rt)以及 Bootstrap 表上按股票代码查询max(ts)与_hoodie_commit_time等典型 SQL; - 增量查询(DataFrame API):docker/demo/sparksql-incremental.commands 展示了完整的增量读流程:通过
HoodieDataSourceHelpers.streamCompletionTimeSince计算起始 completion time,再以spark.read.format("hudi")配合QUERY_TYPE_INCREMENTAL_OPT_VAL与START_COMMIT选项加载增量数据; - 本地运行脚本:hudi-spark-datasource/hudi-spark/run_hoodie_app.sh 与 hudi-spark-datasource/hudi-spark/run_hoodie_streaming_app.sh 演示了如何基于
packaging/hudi-spark-bundle打出的 bundle jar 组装 classpath 并运行示例应用,其中流式脚本对应HoodieJavaStreamingApp。
小结
hudi-spark-datasource通过"公共层 + 版本适配层 + 主模块"的三层结构,把 Hudi 的 DataSource、SQL 扩展、存储过程、索引、增量与流式能力完整地带给了从 Spark 3.3 到 Spark 4.2 的多个版本。选用时请务必对照版本支持矩阵选择正确的构建 Profile,并在 Spark 3.3/3.4 上关注 ZSTD 压缩的 PARQUET-2160 内存风险——必要时显式设置hoodie.parquet.compression.codec。以上配置与限制均以当前仓库实际内容为准,部署前建议结合目标环境的 Spark 发行版再做一次版本核对。
- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
相关推荐
Apache Zeppelin Spark Interpreter 模块架构与多版本支持深度解析
Apache Zeppelin Spark Interpreter 模块架构与多版本支持深度解析 Spark interpreter 是 Apache Zepp
数据分析数据可视化大数据后端前端任务调度OpenDrop版本支持矩阵:Python版本与OS兼容性
OpenDrop版本支持矩阵:Python版本与OS兼容性 你是否在部署OpenDrop时遇到过Python版本不兼容或操作系统支持问题?本文将详细解析Open
网络通信探索KiCad 4.0核心资源:gh_mirrors/ki/kicad-library完全解析
探索KiCad 4.0核心资源:gh_mirrors/ki/kicad library完全解析 gh_mirrors/ki/kicad library是KiCa
硬件开发
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考