使用Spark Shell创建Hudi表
操作场景
在数据湖开发场景中,用户需要通过编程方式灵活地操作Hudi表,包括写入、更新、删除和增量查询数据。Spark SQL方式适用于标准SQL操作,但对于需要更细粒度控制写入参数、实现复杂数据处理逻辑的场景,SQL方式灵活性不足。那么,如何通过编程API方式灵活操作Hudi表?
本节介绍如何在spark-shell中通过Spark DataSource API操作Hudi COW表。
使用Spark数据源,通过代码段展示如何插入和更新Hudi的默认存储类型数据集COW表,以及每次写操作之后如何读取快照和增量数据。
前提条件
如果集群已开启Kerberos认证,需在Manager界面创建1个人机用户并关联到hadoop和hive用户组,主组为hadoop,具体操作请参考创建MRS集群用户。
约束与限制
本章节适用于MRS 3.3.1-LTS及之前版本,不同版本集群由于组件内核差异,相关命令可能有差异。
操作步骤
- 下载并安装Hudi客户端,具体请参考安装MRS客户端章节。
目前Hudi集成在Spark/Spark2x组件中,用户从Manager页面下载Spark/Spark2x客户端即可,例如客户端安装目录为:“/opt/client”。
- 使用客户端安装用户登录客户端节点,执行如下命令: 进入客户端目录:
cd /opt/hadoopclient执行以下命令加载环境变量:
source bigdata_env
source Hudi/component_env
安全认证:
kinit 创建的业务用户
- 新创建的用户首次认证需要修改密码。
- 普通模式(未开启kerberos认证)集群无需执行kinit命令。
- 进入spark-shell,然后引入Hudi相关软件包并生成测试数据。
spark-shell --master yarn-client
- 引入需要的包。
执行Spark操作Hudi 表前,需预先导入以下核心依赖包,覆盖测试工具、集合兼容、Spark写入模式、Hudi读写参数、高级写入配置能力。
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._
- 定义表名,存储路径,生成测试数据。
// 定义Hudi表基础信息:表名、HDFS存储根路径 val tableName = "hudi_cow_table" val basePath = "hdfs://hacluster/tmp/hudi_cow_table" // 初始化Hudi内置数据生成器,构造测试写入数据集 val dataGen = new DataGenerator // 生成10条测试增量数据,并转换为Java字符串集合 val inserts = convertToStringList(dataGen.generateInserts(10)) // 将数据集转为RDD并并行度2,加载为JSON格式DataFrame val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
- 引入需要的包。
- 将测试数据以OVERWRITE模式写入Hudi表,完成表的初始化写入。执行以下命令写入Hudi表,模式为OVERWRITE。
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option(PRECOMBINE_FIELD_OPT_KEY, "ts"). option(RECORDKEY_FIELD_OPT_KEY, "uuid"). option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath"). option(TABLE_NAME, tableName). mode(Overwrite). save(basePath)关键参数说明:
- PRECOMBINE_FIELD_OPT_KEY:该参数用于指定预合并字段,同主键多条数据按该字段值大小取最新记录,此处指定为“ts”(时间戳)。
- RECORDKEY_FIELD_OPT_KEY:该参数用于指定Hudi表的主键字段,此处指定为“uuid”,用于记录级更新和删除。
- PARTITIONPATH_FIELD_OPT_KEY:该参数用于指定分区字段,数据将按该字段值进行分区存储,此处指定为“partitionpath”。
- TABLE_NAME:该参数用于指定Hudi表名称。
- 注册临时表并查询数据,验证写入操作是否成功。执行以下命令注册临时表并查询。
val roViewDF = spark.read.format("org.apache.hudi").load(basePath + "/*/*/*/*") roViewDF.createOrReplaceTempView("hudi_ro_table") spark.sql("select fare, begin_lon, begin_lat, ts from hudi_ro_table where fare > 20.0").show() - 生成更新数据并以APPEND模式写入,完成数据更新操作。执行以下命令生成更新数据并更新Hudi表,模式为APPEND。
val updates = convertToStringList(dataGen.generateUpdates(10)) val df = spark.read.json(spark.sparkContext.parallelize(updates, 1)) df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option(PRECOMBINE_FIELD_OPT_KEY, "ts"). option(RECORDKEY_FIELD_OPT_KEY, "uuid"). option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath"). option(TABLE_NAME, tableName). mode(Append). save(basePath) - 查询Hudi表增量数据。
- 重新加载:
spark.read.format("org.apache.hudi").load(basePath + "/*/*/*/*").createOrReplaceTempView("hudi_ro_table") - 进行增量查询:
val commits = spark.sql("select distinct(_hoodie_commit_time) as commitTime from hudi_ro_table order by commitTime").map(k => k.getString(0)).take(50) val beginTime = commits(commits.length - 2) val incViewDF = spark.read.format("org.apache.hudi"). option(VIEW_TYPE_OPT_KEY, VIEW_TYPE_INCREMENTAL_OPT_VAL). option(BEGIN_INSTANTTIME_OPT_KEY, beginTime). load(basePath); incViewDF.registerTempTable("hudi_incr_table") spark.sql("select `_hoodie_commit_time`, fare, begin_lon, begin_lat, ts from hudi_incr_table where fare > 20.0").show()
- 重新加载:
- 进行指定时间点提交的查询。
val beginTime = "000" val endTime = commits(commits.length - 2) val incViewDF = spark.read.format("org.apache.hudi"). option(VIEW_TYPE_OPT_KEY, VIEW_TYPE_INCREMENTAL_OPT_VAL). option(BEGIN_INSTANTTIME_OPT_KEY, beginTime). option(END_INSTANTTIME_OPT_KEY, endTime). load(basePath); incViewDF.registerTempTable("hudi_incr_table") spark.sql("select `_hoodie_commit_time`, fare, begin_lon, begin_lat, ts from hudi_incr_table where fare > 20.0").show() - 删除测试数据。
- 准备删除的数据。
val df = spark.sql("select uuid, partitionpath from hudi_ro_table limit 2") val deletes = dataGen.generateDeletes(df.collectAsList()) - 执行删除操作。
val df = spark.read.json(spark.sparkContext.parallelize(deletes, 2));
df.write.format("org.apache.hudi"). options(getQuickstartWriteConfigs). option(OPERATION_OPT_KEY,"delete"). option(PRECOMBINE_FIELD_OPT_KEY, "ts"). option(RECORDKEY_FIELD_OPT_KEY, "uuid"). option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath"). option(TABLE_NAME, tableName). mode(Append). save(basePath); - 重新查询数据。
val roViewDFAfterDelete = spark.read.format("org.apache.hudi"). load(basePath + "/*/*/*/*") roViewDFAfterDelete.createOrReplaceTempView("hudi_ro_table") spark.sql("select uuid, partitionPath from hudi_ro_table").show()
- 准备删除的数据。
验证删除操作是否生效:
查询数据,确认上述步骤中选取的2条记录已不在查询结果中,表示删除操作成功。
常见问题
- spark-shell启动时报错"类未找到"如何处理?
请检查是否已正确执行source Hudi/component_env加载Hudi组件环境变量。如果未加载该环境变量,spark-shell将无法找到Hudi相关依赖类。
- OVERWRITE和APPEND模式有什么区别?
OVERWRITE模式会覆盖目标表中的所有数据后重新写入;APPEND模式在已有数据基础上追加写入,不会清除已有数据。首次建表写入时建议使用OVERWRITE模式,后续更新数据时使用APPEND模式。
相关文档
- 使用Spark SQL方式操作Hudi表,请参考使用Spark SQL操作Hudi表。
- 了解更多Hudi SQL语法,请参考Hudi SQL语法参考。