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

如何查看Spark DataFrame.write()写入文件的进度百分比或完成分区数?

嘿,这个场景我太熟悉了——大DataFrame写入的时候盯着屏幕等半天,总想知道到底跑了多少,下面几个实用方法帮你搞定:

方法1:用Spark Web UI直观查看

这是最简单的方式,Spark自带的Web UI(默认端口4040)就能实时追踪:

  • 打开 http://<你的Driver节点IP>:4040,切换到Jobs页面
  • 找到触发写入操作的那个Job(就是你调用save()/write()后生成的Job),点进去看对应的Stage
  • 每个Stage里的Task数量基本就对应你的DataFrame分区数,已完成的Task数除以总Task数,就是当前的写入进度百分比
  • 要是你用YARN/Mesos集群,也可以通过集群的管理UI(比如YARN的ResourceManager页面)查看任务整体进度
方法2:自定义监听器,代码里实时输出进度

如果想在代码里主动获取进度(比如打印到控制台或者写入监控系统),可以写个自定义的SparkListener来追踪Task完成情况:

import org.apache.spark.scheduler._

class WriteProgressTracker extends SparkListener {
  override def onTaskEnd(taskEnd: SparkListenerTaskEnd): Unit = {
    val stageId = taskEnd.stageId
    val sc = SparkContext.getOrCreate()
    sc.statusTracker.getStageInfo(stageId).foreach { stageInfo =>
      val completed = stageInfo.numCompletedTasks
      val total = stageInfo.numTasks
      val progress = (completed.toDouble / total) * 100
      println(s"当前写入进度: %.2f%% (已完成 ${completed}/${total} 个分区)".format(progress))
    }
  }
}

// 注册监听器
val spark = SparkSession.builder().getOrCreate()
spark.sparkContext.addSparkListener(new WriteProgressTracker())

// 执行写入
df.write.format("parquet").mode("overwrite").save("/target/path")

注:Python/Java版本的逻辑类似,只是语法不同,核心都是通过监听TaskEnd事件来统计已完成的分区数

方法3:提前拿分区数+查日志

如果不想改代码,也可以这么做:

  • 先通过 df.rdd.getNumPartitions() (Scala/Python)获取DataFrame的总分区数
  • 然后查看Spark Driver或者Executor的日志,里面会记录每个Task(对应一个分区)的完成日志
  • 数一下已完成的Task条目,就能算出已完成的分区数和进度

小提醒:如果写入操作涉及shuffle(比如之前有groupBy、join等操作),要注意找到最后一个Stage——那才是实际执行文件写入的阶段,对应的Task数才是你要关注的分区数。

内容的提问来源于stack exchange,提问作者Milind Jindal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:19:09