PySpark升级3.4.1后onExecutorMetricsUpdate获取executor_id异常问题
PySpark 3.4.1中Executor指标获取异常的原因及解决方法
原因分析
- Spark 3.x版本对
ExecutorMetricsUpdate事件的处理逻辑做了调整,Yarn模式下Driver端接收的该事件默认不再携带真实Executor ID,统一标记为'driver',这是由于3.x优化了Metrics传输机制,部分指标聚合逻辑迁移到Driver端,但存在Executor ID映射丢失的场景。 peakExecutionMemory为0、/executors端点仅显示driver的问题,主要是因为默认配置下未开启完整的Executor级别Metrics采集,导致Driver无法接收Executor上报的真实执行内存数据;同时Yarn集群的Metrics上报限制也可能导致UI端点无法展示所有Executor信息。
解决方法
1. 开启完整的Metrics采集配置
在spark-submit参数或spark-defaults.conf中添加以下配置,确保Executor能正确上报Metrics:
spark.metrics.conf.*.sink.servlet.class=org.apache.spark.metrics.sink.MetricsServlet spark.executor.metrics.enabled=true spark.sql.metrics.enabled=true spark.executor.memoryMetrics.enabled=true
这些配置会触发Executor端的Metrics采集,并确保Spark UI的/executors端点展示所有Executor的信息。
2. 改用正确的方式获取Executor ID与指标
- 优先监听
TaskEnd事件:从taskInfo.executorId获取真实的Executor ID,同时从taskMetrics.peakExecutionMemory获取该Task的执行内存峰值,再聚合到对应Executor的统计报告中,这是Spark 3.x中获取Task级内存指标最可靠的方式。 - 若需使用
ExecutorMetricsUpdate事件:可通过executorMetricsUpdate.sourceName()区分来源,部分场景下该字段会包含Executor标识;也可结合ExecutorAllocationManager获取已注册的Executor列表,通过Metrics的tag(如executorId标签)匹配对应Executor。
3. 优化自定义MemoryReporter
- 监听
ExecutorAdded事件,记录所有注册的Executor ID(从executorInfo.executorId获取),后续处理Metrics事件时,通过Metrics上下文的tag信息关联到对应Executor,不再依赖executorMetricsUpdate.execId()。 - 对于Executor级内存指标,开启
spark.executor.memoryMetrics.enabled=true后,可从executorMetricsUpdate.metrics()中获取如JVMHeapMemory、JVMOffHeapMemory等非0数据,用于统计报告。
内容的提问来源于stack exchange,提问作者fersarr
相关产品推荐
相关产品推荐

