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

如何将EMR上PySpark作业的峰值内存利用率写入文件?

如何在EMR的PySpark作业中自动记录峰值内存利用率?

我们在EMR上运行大量PySpark作业,执行流程固定,但输入数据会大幅改变峰值内存利用率,且该利用率呈上升趋势。希望自动将每个作业步骤的峰值内存利用率写入文件(比如记录“作业峰值使用了10TB内存”)。作业不受CPU或其他指标限制,仅关注内存,若输出所有聚合指标能简化实现方案也可以接受。

说明:采用YARN作为集群管理器的集群模式,且以Docker容器形式提交作业。

我怀疑自己可能过度复杂化了问题,原本以为会有更简单的方式,比如:

spark = init_spark()

do_some_stuff()

spark.metrics.peakMemoryUtilization()

如果有这类遗漏的方法,还请指教。


尝试1:继承JVM Spark Listener的自定义类

根据GitHub Copilot的建议,尝试用SparkListeners实现但未成功,需要创建自定义类借助JVM访问任务指标。

代码实现

from pyspark.sql import SparkSession
from pyspark import SparkConf
from py4j.java_gateway import java_import

class MemoryUsageListener:
    def __init__(self):
        self.peak_memory_usage = 0

    def on_task_end(self, task_end):
        metrics = task_end.taskMetrics()
        memory_used = metrics.peakExecutionMemory()
        if memory_used > self.peak_memory_usage:
            self.peak_memory_usage = memory_used

    # 用于将监听器添加到SparkContext的方法
    def get_listener(self, sc):
        gw = sc._gateway
        java_import(gw.jvm, "org.apache.spark.scheduler.SparkListener")
        java_import(gw.jvm, "org.apache.spark.scheduler.SparkListenerTaskEnd")

        class JavaListener(gw.jvm.SparkListener):
            def __init__(self, parent):
                super(gw.jvm.SparkListener, self).__init__()
                self.parent = parent

            def onTaskEnd(self, taskEnd):
                self.parent.on_task_end(taskEnd)

        return JavaListener(self)

    def get_peak_memory_usage(self):
        return self.peak_memory_usage

初始化代码

conf = SparkConf().setAppName("MemoryUsageTracker")
spark = SparkSession.builder.config(conf=conf).getOrCreate()
sc = spark.sparkContext

# 创建并注册监听器
memory_listener = MemoryUsageListener()
java_listener = memory_listener.get_listener(sc)
sc._jsc.sc().addSparkListener(java_listener)

... 执行实际Spark作业逻辑 ...

peak_memory_usage = memory_listener.get_peak_memory_usage()
write_to_file(peak_memory_usage)

遇到的问题

初始化类时出现JavaClass.__init__() "takes 3 positional arguments but 4 were given"错误,根源在class JavaListener(gw.jvm.SparkListener)这一行。尝试过调整super声明位置、将Java类移至独立函数,但都未解决。


尝试2:调用YARN API

既然直接从Spark获取指标行不通,尝试用requests库调用API端点获取执行器详情。

代码实现

import requests
import socket
from pyspark.sql import SparkSession
from pyspark import SparkConf

def get_executor_memory_metrics(app_id, spark_history_server_url):
    response = requests.get(f"{spark_history_server_url}/api/v1/applications/{app_id}/executors")
    if response.status_code != 200:
        raise Exception(f"Failed to fetch executor metrics, status code: {response.status_code}")
    executors = response.json()
    return executors

def calculate_peak_memory_usage(executors):
    # 不确定这个字典访问是否正确,是Copilot生成的,目前不是主要问题
    peak_memory_usage = max(executor['peakMemoryMetrics']['JVMHeapMemory'] for executor in executors if 'peakMemoryMetrics' in executor)
    return peak_memory_usage

sc = spark.sparkContext
app_id = sc.applicationId

master_url = sc.master
# 不确定这个逻辑是否合理,但没有权限修改集群/YARN设置
yarn_resource_manager_host = socket.getfqdn()
yarn_port = 18080

# 在作业结束时
executors = get_executor_memory_metrics(app_id, f"{yarn_resource_manager_host}:{yarn_port}")
peak_memory_usage = calculate_peak_memory_usage(executors)

遇到的问题

请求部分持续出现“Connection Refused”错误。已确认集群IP和应用ID正确传入URL,但不确定端口号是否正确,也不清楚运行时是否有更好的访问方式。另外,这段代码本来需要封装到线程中循环监控峰值利用率,但目前连请求都无法成功。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 02:15:08