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

