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

对应的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
相关产品推荐
相关产品推荐

