SparkListener.onStageCompleted获取的TaskMetrics是否为Stage所有Task的聚合指标?
Spark StageCompleted中TaskMetrics的含义
我正在编写自定义SparkListener以将指标发送至第三方系统,希望在Stage结束时发送该Stage内所有Task的聚合指标。SparkListener提供onStageCompleted()方法,可通过如下代码重写以获取TaskMetrics:
@Override public void onStageCompleted(SparkListenerStageCompleted stageCompleted) { TaskMetrics taskMetrics = stageCompleted.stageInfo().taskMetrics(); }
问题
上述taskMetrics对象的指标代表什么?是该Stage中最新Task的指标,还是所有Task的聚合指标?我使用的是Spark 3.2.3版本,查阅官方文档未找到相关说明。
解答
在Spark 3.2.3版本中,stageCompleted.stageInfo().taskMetrics()返回的是该Stage内所有Task的聚合指标,并非单个最新Task的指标。
这些聚合指标是Spark在Stage完成后,对该Stage下所有Task的各项指标(比如输入输出字节数、Shuffle数据量、执行时间等)按照对应逻辑汇总得到的结果:数据类指标(如字节数)一般是求和统计,时间类指标(如执行耗时)会包含最大值、最小值、平均值等维度,完全匹配你想要在Stage结束时发送全Stage Task聚合指标的需求。
内容的提问来源于stack exchange,提问作者matszewczyk
相关产品推荐
相关产品推荐

