
# 使用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集群用户](https://support.huaweicloud.com/usermanual-mrs/mrs_01_0345.html)。
#### 约束与限制
本章节适用于MRS 3.3.1-LTS及之前版本，不同版本集群由于组件内核差异，相关命令可能有差异。
#### 操作步骤
1. 下载并安装Hudi客户端，具体请参考[安装MRS客户端](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_2127.html)章节。
   
   目前Hudi集成在Spark/Spark2x组件中，用户从Manager页面下载Spark/Spark2x客户端即可，例如客户端安装目录为："/opt/client"。
   
   
2. 使用客户端安装用户登录客户端节点，执行如下命令： 
   进入客户端目录：
   ```
   cd /opt/hadoopclient
   ```
   执行以下命令加载环境变量：
   ```
   source bigdata_env
   ```
   ```
   source Hudi/component_env
   ```
   安全认证：
   ```
   kinit 创建的业务用户
   ```
   ![](https://support.huaweicloud.com/cmpntguide-lts-mrs/public_sys-resources/note_3.0-zh-cn.png)
   - 新创建的用户首次认证需要修改密码。
   
   - 普通模式（未开启kerberos认证）集群无需执行**kinit**命令。
    
   
   
3. 进入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))
     ```
     
   
   
   
   
4. 将测试数据以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表名称。
   
   
   
   
5. 注册临时表并查询数据，验证写入操作是否成功。执行以下命令注册临时表并查询。 
   ```
   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()
   ```
   
   
6. 生成更新数据并以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)
   ```
   
   
7. 查询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()
     ```
     
   
   
   
   
8. 进行指定时间点提交的查询。 
   ```
   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()
   ```
   
   
9. 删除测试数据。 
   - 准备删除的数据。
     ```
     val df = spark.sql("select uuid, partitionpath from hudi_ro_table limit 2")
     val deletes = dataGen.generateDeletes(df.collectAsList())
     ```
     
   
   - 执行删除操作。 将删除数据集转换为JSON格式的DataFrame：
     ```
     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表](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_300441.html)。
- 了解更多Hudi SQL语法，请参考[Hudi SQL语法参考](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_24261.html)。
 
