You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何确认insertInto是否启用了动态分区覆写功能?

如何在Spark UI中验证动态分区更新功能是否生效

我有一个可复用函数,通过saveAsTable处理初始表加载,通过insertInto处理增量更新,并且在insertInto中设置了partitionOverwriteMode=dynamic来启用动态分区更新(仅覆盖DataFrame中包含数据的分区)。现在想确认:能不能在Spark UI里检查这个功能是否正在生效?我怀疑DAG截图里的number of dynamic part可以作为依据,但不确定。

DAG中的动态分区指标

对应的Scala函数代码如下:

def saveAsManagedTable(df: DataFrame, dbOrSchemaName: String, tableName: String, partitionColNames: Array[String],
                         comment: String, numPartitions: Int = -1, saveMode: SaveMode = SaveMode.Overwrite,
                         fileFormat: StorageFormat = StorageFormat.Delta, otherWriterOpts: Map[String, String] = Map.empty): Unit = {
    //Check pre-conditions
    require((fileFormat.equals(StorageFormat.Delta) || fileFormat.equals(StorageFormat.Parquet)), "StorageFormat can be one of Delta or Parquet.")

    val spark = SynapseSpark.getActiveSession

    val fullTableName = s"$dbOrSchemaName.$tableName"
    println(
      s"""Saving as : $fullTableName, at location: ${spark.conf.get("spark.sql.warehouse.dir")},
         | with partition cols: ${partitionColNames.mkString(",")} and save mode: $saveMode""".stripMargin)
    
    if (!spark.catalog.databaseExists(dbOrSchemaName)) {
      spark.sql(
        s"""CREATE DATABASE IF NOT EXISTS $dbOrSchemaName
           |COMMENT '$comment'
           |)
           |""".stripMargin)
    }
    val tempDf = Partitions.selectCoalesceOrRepartition(df, numPartitions)

    if (!spark.catalog.tableExists(fullTableName)) {
      println(s"$fullTableName doesn't exist in catalog. Creating it.")
      tempDf
        .write
        .partitionBy(partitionColNames: _*)
        .format(fileFormat.format)
        .options(otherWriterOpts)
        .option("encoding", StandardCharsets.UTF_8.name())
        .mode(saveMode)
        .saveAsTable(fullTableName)
    }
    else {
      //Raises "AnalysisException: Table not found." if the table already doesn't exist
     println(s"$fullTableName already exists in catalog. Appending to it.")
      tempDf
        .write
        //.partitionBy(partitionCol) Not required as it retrieves the partition info from catalog
        .format(fileFormat.format)
        .options(otherWriterOpts)
        .option("partitionOverwriteMode", "dynamic") //The dynamic value makes sure that Spark will overwrite only partitions that we have data for in our DataFrame.
        .option("encoding", StandardCharsets.UTF_8.name())
        .mode(saveMode) //mode is not needed as it is "Append" by default
        .insertInto(fullTableName)
    }
  }

验证方法:

  • DAG中的number of dynamic part是直接证据:你观察的点是对的,当partitionOverwriteMode=dynamic生效时,Spark会在写入阶段的DAG节点上显示number of dynamic part指标,数值代表本次作业实际处理的动态分区数量。如果是静态模式(static),这个指标不会出现,且会处理全部分区。
  • 查看Stage详情的任务明细:进入写入对应的Stage(通常是最后一个Stage),查看每个任务的输入输出信息。动态分区模式下,任务只会针对DataFrame中存在数据的分区生成/覆盖文件,不会涉及全表所有分区。
  • 检查环境配置参数:在Spark UI的「Environment」标签页,搜索spark.sql.sources.partitionOverwriteMode,确认其值为dynamic,这能证明该配置已在作业中生效。
  • 对比写入前后的分区文件:如果是Parquet表,可以在存储路径下看对应分区的文件更新时间;如果是Delta表,Spark UI的「SQL」标签页(若开启)会显示Delta的操作日志,其中会明确记录是动态覆盖分区还是全表覆盖。

内容的提问来源于stack exchange,提问作者Aravind Yarram

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 03:35:41