更新时间:2026-08-29 GMT+08:00
分享

使用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及之前版本,不同版本集群由于组件内核差异,相关命令可能有差异。

操作步骤

  1. 下载并安装Hudi客户端,具体请参考安装MRS客户端章节。

    目前Hudi集成在Spark/Spark2x组件中,用户从Manager页面下载Spark/Spark2x客户端即可,例如客户端安装目录为:“/opt/client”。

  2. 使用客户端安装用户登录客户端节点,执行如下命令:

    进入客户端目录:
    cd /opt/hadoopclient

    执行以下命令加载环境变量:

    source bigdata_env
    source Hudi/component_env

    安全认证:

    kinit 创建的业务用户
    • 新创建的用户首次认证需要修改密码。
    • 普通模式(未开启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模式。

相关文档

相关文档