
数据湖湖仓一体大数据数据存储【免费下载链接】hudiUpserts, Deletes And Incremental Processing on Big Data.项目地址https://gitcode.com/gh_mirrors/hud/hudi点击查看免费下载Apache Hudi 作为大数据场景下的数据湖框架提供了完整的 Upsert、Delete 与增量处理能力。本文以仓库内 python 目录下的 README 文档 为骨架结合 HoodiePySparkQuickstart.py 脚本源码系统讲解如何在 PySpark 环境下搭建 Hudi 运行环境、编译获取 Hudi Spark Bundle并逐段拆解快速入门脚本中覆盖的插入、更新、快照查询、时间旅行、增量查询、时间点查询、软删除、硬删除与插入覆盖等全部数据操作使读者能够独立复现完整的 Hudi PySpark 实战演练。一、环境前提Python 版本与 Hudi 构建仓库中的 PySpark 快速入门脚本以 Python 编写因此运行它的第一前提是系统安装 Python。这里需要特别留意一个版本兼容性约束文档明确指出PySpark 2.4.7 无法与较新的 Python3.8版本协同工作。如果你本机安装的是较新版本的 Python例如文档示例中的 3.5 之后的版本实际对应现代的 3.8就需要使用支持新 Python 版本的 PySpark 组合并据此重新构建 Hudi。在仓库根目录下执行以下命令即可构建出配套的 Hudi 产物cd $HUDI_DIR mvn clean install -DskipTests -Dspark3.5 -Dscala2.12参数含义说明参数作用-DskipTests跳过测试加快构建速度-Dspark3.5激活 Spark 3.5 构建 profile对应仓库根 pom.xml 中idspark3.5/id的 profile 配置产物将基于hudi-spark-datasource/hudi-spark3.5.x模块-Dscala2.12使用 Scala 2.12 构建最终产物形如hudi-spark3.5-bundle_2.12此外脚本运行时可能还需要若干 Python 依赖包。建议先安装 pip再通过pip install package name按需补齐。二、安装 PySpark 并配置环境变量构建好 Hudi 之后还需要本地拥有一个可用的 PySpark 运行环境。按照以下四步完成安装与配置1. 下载 PySpark从 Spark 官网下载与构建产物匹配的 PySpark 发行包仓库脚本默认面向 3.5.x 系列见 vector_blob_demo/requirements.txt 中pyspark3.5.*的约束。2. 解压并记录安装位置将下载的压缩包解压到你希望安装的目录并记住该路径后续环境变量会用到。3. 配置环境变量执行以下命令也可追加到~/.bashrc使其永久生效务必把SPARK_HOME替换为真实的 Spark 安装路径export SPARK_HOME/path/to/spark/home export PATH$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin export PYTHONPATH$SPARK_HOME/python/:$PYTHONPATH export PYTHONPATH$SPARK_HOME/python/lib/*.zip:$PYTHONPATH这里前三行分别设置了 Spark 主目录、将bin与sbin目录加入PATH、并把 Spark 的 Python 源码目录加入PYTHONPATH第四行则把 Spark 自带的各种 Python 依赖 zip 包如 py4j也加入PYTHONPATH这是import pyspark能正常工作的关键。4. 准备 Hudi Spark Bundle二选一快速入门脚本通过spark-submit参数将 Hudi 引入运行环境你可以选择以下两种形式之一Maven Package推荐格式形如org.apache.hudi:hudi-spark3.5-bundle_2.12:0.12.0本地 Jar 文件格式形如[HUDI_BASE_PATH]/packaging/hudi-spark-bundle/target/hudi-spark-bundle[VERSION].jar注意仓库中packaging/hudi-spark-bundle是统一的 bundle 构建模块构建产物会以具体版本号命名如果你采用第一节的源码构建方式实际产物名会形如hudi-spark3.5-bundle_2.12-版本-SNAPSHOT.jar。三、运行快速入门脚本进入 Hudi 仓库目录执行以下命令启动演练脚本为 HoodiePySparkQuickstart.pycd $HUDI_DIR python3 hudi-examples/hudi-examples-spark/src/test/python/HoodiePySparkQuickstart.py [-h] -t TABLE (-p PACKAGE | -j JAR)命令行参数说明脚本使用argparse定义见 HoodiePySparkQuickstart.py参数必填说明-t, --table是要创建的 Hudi 表名称-p, --package二选一Hudi Spark Bundle 的 Maven 坐标例如org.apache.hudi:hudi-spark3.5-bundle_2.12:0.12.0-j, --jar二选一Hudi Spark Bundle 本地 jar 的完整路径-h, --help否显示帮助信息-p与-j属于互斥参数add_mutually_exclusive_group二者只能提供一个。脚本会根据你传入的形式自动设置PYSPARK_SUBMIT_ARGS环境变量传-p时os.environ[PYSPARK_SUBMIT_ARGS] f--packages {package} pyspark-shell传-j时os.environ[PYSPARK_SUBMIT_ARGS] f--jars {jar} pyspark-shell值得一提的是脚本运行时会创建一个TemporaryDirectory作为表数据的基础路径即表实际写入路径为临时目录/表名因此无需手动清理每次运行都是干净环境。四、脚本内部的 SparkSession 配置要点进入 HoodiePySparkQuickstart.py可以看到构建SparkSession时设置了四项 Hudi 必需的配置理解它们对自建应用同样重要spark sql.SparkSession \ .builder \ .appName(Hudi Spark basic example) \ .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) \ .config(spark.kryo.registrator, org.apache.spark.HoodieSparkKryoRegistrar) \ .config(spark.kryoserializer.buffer.max, 512m) \ .config(spark.sql.extensions, org.apache.spark.sql.hudi.HoodieSparkSessionExtension) \ .getOrCreate()配置项作用spark.serializer KryoSerializer使用 Kryo 序列化器Hudi 在 Spark 下的标配spark.kryo.registrator HoodieSparkKryoRegistrar注册 Hudi 相关的 Kryo 类与 TestHoodieSparkQuickstart 中HoodieSparkKryoRegistrar$.MODULE$.register(sparkConf)的做法一致spark.kryoserializer.buffer.max 512m放宽 Kryo 缓冲区上限避免大记录序列化失败spark.sql.extensions HoodieSparkSessionExtension注册 Hudi 的 SQL 扩展Catalog、规则等是 Hudi 表能通过format(hudi)读写的基础五、七大核心数据操作逐段拆解快速入门脚本通过ExamplePySpark类的runQuickstart()方法串起全部演练其内部还嵌入了若干断言assert确保每一步操作结果符合预期——这与 Java 版 HoodieSparkQuickstart.java 的runQuickstart流程一一对应。5.1 基础写入配置ExamplePySpark.__init__中集中定义了写表时的 Hudi 选项这是理解整个演练的数据模型入口self.hudi_options { hoodie.table.name: tableName, hoodie.datasource.write.recordkey.field: uuid, hoodie.datasource.write.partitionpath.field: partitionpath, hoodie.datasource.write.operation: upsert, hoodie.table.ordering.fields: ts, hoodie.upsert.shuffle.parallelism: 2, hoodie.insert.shuffle.parallelism: 2 }选项说明hoodie.table.name表名hoodie.datasource.write.recordkey.field记录主键字段这里为uuidhoodie.datasource.write.partitionpath.field分区字段这里为partitionpathhoodie.datasource.write.operation写入操作类型默认upserthoodie.table.ordering.fields排序/预合并字段这里为ts时间戳用于处理同一主键的多次更新hoodie.upsert.shuffle.parallelism/hoodie.insert.shuffle.parallelism两类写操作的 shuffle 并行度脚本设为 2 以适配本地小规模演练测试数据由 Hudi 内置的QuickstartUtils.DataGenerator生成spark._jvm.org.apache.hudi.QuickstartUtils.DataGenerator()这是 Hudi 提供的确定性测试数据生成器生成的每条记录都包含uuid、ts、rider、driver、fare、partitionpath、经纬度等打车订单字段。5.2 插入数据Insertinserts self.spark._jvm.org.apache.hudi.QuickstartUtils.convertToStringList(self.dataGen.generateInserts(10)) df self.spark.read.json(self.spark.sparkContext.parallelize(inserts, 2)) df.write.format(hudi).options(**self.hudi_options).mode(overwrite).save(self.basePath)首轮写入生成 10 条 JSON 记录使用mode(overwrite)创建表。之后脚本会执行快照查询并把结果与插入的 DataFrame 做exceptAll比对断言两集合完全一致验证写入成功。5.3 更新数据Updateupdates self.spark._jvm.org.apache.hudi.QuickstartUtils.convertToStringList(self.dataGen.generateUniqueUpdates(5)) df self.spark.read.json(spark.sparkContext.parallelize(updates, 2)) df.write.format(hudi).options(**self.hudi_options).mode(append).save(self.basePath)生成 5 条唯一更新记录generateUniqueUpdates即基于已有主键产生新版本以append模式 upsert 进表。随后的断言验证更新后快照必须包含全部 5 条更新记录且除这 5 条外没有其他记录发生变化。5.4 快照查询Snapshot QuerytripsSnapshotDF self.spark.read.format(hudi).load(self.basePath) tripsSnapshotDF.createOrReplaceTempView(hudi_trips_snapshot) self.spark.sql(SELECT fare, begin_lon, begin_lat, ts FROM hudi_trips_snapshot WHERE fare 20.0).show() self.spark.sql(SELECT _hoodie_commit_time, _hoodie_record_key, _hoodie_partition_path, rider, driver, fare FROM hudi_trips_snapshot).show()快照查询是 Hudi 最基础的读取方式spark.read.format(hudi).load(basePath)读取当前最新提交的快照注册为临时视图后可执行任意 SQL。第二句 SQL 还展示了 Hudi 的三个系统列——_hoodie_commit_time提交时间、_hoodie_record_key记录主键、_hoodie_partition_path分区路径这些元数据列是理解 Hudi 时间线的钥匙。5.5 时间旅行查询Time Travel Queryself.spark.read.format(hudi).option(as.of.instant, 20210728141108).load(self.basePath).createOrReplaceTempView(time_travel_query) self.spark.read.format(hudi).option(as.of.instant, 2021-07-28 14:11:08.000).load(self.basePath).createOrReplaceTempView(time_travel_query) self.spark.read.format(hudi).option(as.of.instant, 2021-07-28).load(self.basePath).createOrReplaceTempView(time_travel_query)通过as.of.instant选项可以把表回放到任意历史时刻。脚本演示了三种时间格式Hudi 均可解析格式示例时间戳14 位20210728141108标准时间2021-07-28 14:11:08.000日期2021-07-28需要说明的是脚本中的示例时间戳是固定的演示值实际使用时应替换为你自己的提交时间可从_hoodie_commit_time列查询。5.6 增量查询Incremental Query增量查询是 Hudi 区别于普通存储的标志性能力它只返回指定提交时间之后的变更数据self.commits list(map(lambda row: row[0], self.spark.sql( SELECT DISTINCT(_hoodie_commit_time) AS commitTime FROM hudi_trips_snapshot ORDER BY commitTime).limit(50).collect())) beginTime self.commits[len(self.commits) - 2] incremental_read_options { hoodie.datasource.query.type: incremental, hoodie.datasource.read.begin.instanttime: beginTime, } tripsIncrementalDF self.spark.read.format(hudi).options(**incremental_read_options).load(self.basePath) tripsIncrementalDF.createOrReplaceTempView(hudi_trips_incremental) self.spark.sql(SELECT _hoodie_commit_time, fare, begin_lon, begin_lat, ts FROM hudi_trips_incremental WHERE fare 20.0).show()实现要点先从快照中收集去重后的提交时间列表取倒数第二个提交作为beginTime即上一次提交之后发生的变化设置hoodie.datasource.query.type incremental与hoodie.datasource.read.begin.instanttime两个读取选项结果注册为hudi_trips_incremental视图后可配合_hoodie_commit_time列做过滤分析。5.7 时间点查询Point-in-Time Query时间点查询是增量查询的变体通过同时指定起止提交时间把读取范围收敛到某一段时间线区间beginTime 000 endTime self.commits[len(self.commits) - 2] point_in_time_read_options { hoodie.datasource.query.type: incremental, hoodie.datasource.read.end.instanttime: endTime, hoodie.datasource.read.begin.instanttime: beginTime }begin.instanttime使用000表示从最早提交开始end.instanttime指定截止提交从而得到一个闭区间内的历史快照视图。5.8 软删除Soft Deletes软删除并不真正移除记录而是把记录中除主键外的字段置为null通过 upsert 覆盖旧版本meta_columns [_hoodie_commit_time, _hoodie_commit_seqno, _hoodie_record_key, _hoodie_partition_path, _hoodie_file_name] excluded_columns meta_columns [ts, uuid, partitionpath] nullify_columns list(filter(lambda field: field[0] not in excluded_columns, ...)) soft_delete_df reduce(lambda df, col: df.withColumn(col[0], lit(None).cast(col[1])), nullify_columns, ...) soft_delete_df.write.format(hudi).options(**hudi_soft_delete_options).mode(append).save(self.basePath)流程拆解从快照中取 2 条记录作为软删除对象把系统列_hoodie_commit_time等 5 列与主键、分区、排序字段ts/uuid/partitionpath排除在外其余业务字段全部置空以 upsert 模式写回。脚本用两次count()验证效果总记录数不变软删除不减少行但rider IS NOT null的记录数减少 2证明字段已被置空。5.9 硬删除Hard Deletes硬删除真正从表中移除记录通过hoodie.datasource.write.operation delete实现hudi_hard_delete_options { ... hoodie.datasource.write.operation: delete, ... } deletes list(map(lambda row: (row[0], row[1]), ds.collect())) hard_delete_df self.spark.sparkContext.parallelize(deletes).toDF([uuid, partitionpath]).withColumn(ts, lit(0.0)) hard_delete_df.write.format(hudi).options(**hudi_hard_delete_options).mode(append).save(self.basePath)要点是删除 DataFrame 只需携带主键uuid与分区partitionpath字段ts置为 0.0 占位即可。删除后重新查询记录数应减少 2。5.10 插入覆盖Insert Overwrite最后演示的是分区级覆盖写常用于整分区重建的场景hudi_insert_overwrite_options { ... hoodie.datasource.write.operation: insert_overwrite, ... } df self.spark.read.json(...).filter(partitionpath americas/united_states/san_francisco) df.write.format(hudi).options(**hudi_insert_overwrite_options).mode(append).save(self.basePath)这里仅对san_francisco分区生成 10 条新记录并执行insert_overwrite因此只会覆盖该分区脚本断言验证快照结果 除该分区外的原有数据 ∪ 新写入的该分区数据。注意insert_overwrite 覆盖的是整个分区而非单条记录这是与 upsert 的本质区别。六、与 Java 版 Quickstart 及测试用例的对应关系PySpark 快速入门并非孤立存在它与仓库内的 Java 版本实现同源同构可以作为交叉验证Java 版入口HoodieSparkQuickstart.java 的runQuickstart方法执行同样的 insert → update → 快照查询 → 增量查询 → 时间点查询 → 删除 → 插入覆盖流程并额外演示了deleteByPartition按分区删除能力自动化测试TestHoodieSparkQuickstart.java 通过Test testHoodieSparkQuickstart直接调用上述 Java 流程并在BeforeEach中完成HoodieSparkKryoRegistrar注册与SparkRDDReadClient.addHoodieSupport等 Spark 环境初始化——这正是 Python 脚本中 SparkSession 四项配置的底层依据。也就是说你可以把 PySpark 脚本当作用 Python 表述的同一套 Hudi 操作契约Java 源码与测试则提供了更底层的实现参考。七、进阶同目录下的 VECTOR BLOB 向量检索演练在 python 目录 下还提供了一个进阶演练 vector_blob_demo它基于 Apache Hudi 1.2.0 的VECTOR 类型、BLOB 类型与hudi_vector_search向量检索 TVF三个特性在 Oxford-IIIT Pet 数据集上做端到端演示包含hudi_blob_reader_demo.pyOUT_OF_LINE BLOB read_blob()、hudi_sql_vector_blob_demo.pySQL DDL INLINE BLOB 向量搜索与hudi_dataframe_vector_blob_demo.pyDataFrame API三个变体并可通过HUDI_BASE_FILE_FORMAT环境变量在 Lance 与 Parquet 基础文件格式之间切换。该演练的环境变量一览HUDI_BUNDLE_JAR、LANCE_BUNDLE_JAR、HUDI_LANCE_DEMO_N等与快速入门脚本的-p/-j参数思路一脉相承都是让 Python 脚本通过提交参数把 Hudi/Jar 注入 Spark 运行时。对于已经掌握本文基础流程的读者这是一个理想的进阶方向。八、常见问题与调优建议结合脚本源码与仓库配置汇总几个实战中高频问题问题原因与对策import pyspark失败SPARK_HOME/python/与python/lib/*.zip未正确加入PYTHONPATH检查环境变量四行配置找不到 Hudi 类如QuickstartUtils未在PYSPARK_SUBMIT_ARGS中注入 bundle确认-p/-j参数指向正确的包坐标或 jar 路径Python 版本过高导致 PySpark 异常PySpark 2.4.7 不支持 Python 3.8按第一节命令用-Dspark3.5 -Dscala2.12构建新版本 Hudi并匹配 Spark 3.5.x 的 PySpark数据量大时写性能差本地演练的 shuffle 并行度仅为 2生产环境应依据集群规模调大hoodie.upsert.shuffle.parallelism与hoodie.insert.shuffle.parallelism固定时间戳不生效时间旅行查询的as.of.instant必须使用表真实存在的提交时间可从_hoodie_commit_time列获取结语本文从环境构建、Bundle 准备、脚本运行到操作拆解完整还原了 Hudi PySpark 快速入门的全过程。核心要点可归纳为三条版本匹配Python 与 PySpark/Hudi 构建参数必须对应、运行时注入通过-p/-j把 Hudi Bundle 交给PYSPARK_SUBMIT_ARGS、操作即练习脚本内置断言确保 insert/update/query/time-travel/incremental/delete/insert-overwrite 每一步都被验证。以此为起点你可以进一步阅读 HoodiePySparkQuickstart.py 的完整源码或向 vector_blob_demo 的向量检索方向深入。赞分享数据湖湖仓一体大数据数据存储【免费下载链接】hudiUpserts, Deletes And Incremental Processing on Big Data.项目地址https://gitcode.com/gh_mirrors/hud/hudi点击查看免费下载相关推荐如何用Pipenv打造高效大数据项目环境PySpark与Dask的终极配置指南如何用Pipenv打造高效大数据项目环境PySpark与Dask的终极配置指南 Pipenv是Python开发工作流的最佳伴侣它将虚拟环境管理和依赖管理无缝开发工具CLI包管理器mmdetection3d快速入门从环境搭建到首个3D检测模型训练全流程mmdetection3d快速入门从环境搭建到首个3D检测模型训练全流程 1. 引言3D目标检测的痛点与解决方案 你是否在3D目标检测项目中遇到过以下问题人工智能计算机视觉深度学习自动驾驶Apache Hudi快速入门指南10分钟掌握Spark Shell数据操作技巧Apache Hudi快速入门指南10分钟掌握Spark Shell数据操作技巧 Apache Hudi是一个开源的分布式列存储系统专门用于处理大规模时间序数据湖湖仓一体大数据数据存储上一篇Thorium 浏览器实战指南更快、媒体更全的 Chromium 优化浏览器下一篇5分钟快速上手LinkSwift网盘直链下载助手终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考