如何查看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
相关产品推荐
相关产品推荐

