lau*_*108 41 bigdata amazon-emr emr apache-spark
我正在AWS EMR上运行5节点Spark群集,每个群集大小为m3.xlarge(1个主4个从属).我成功地运行了一个146Mb的bzip2压缩CSV文件,最终获得了完美的聚合结果.
现在我正在尝试在此群集上处理~5GB bzip2 CSV文件但我收到此错误:
16/11/23 17:29:53 WARN TaskSetManager:阶段6.0中丢失的任务49.2(TID xxx,xxx.xxx.xxx.compute.internal):ExecutorLostFailure(执行者16退出由其中一个正在运行的任务引起)原因:容器由于超过内存限制而被YARN杀死.使用10.4 GB的10.4 GB物理内存.考虑提升spark.yarn.executor.memoryOverhead.
我很困惑为什么我在~75GB群集上获得~10.5GB内存限制(每3m.xlarge实例15GB)...
这是我的EMR配置:
[
{
"classification":"spark-env",
"properties":{
},
"configurations":[
{
"classification":"export",
"properties":{
"PYSPARK_PYTHON":"python34"
},
"configurations":[
]
}
]
},
{
"classification":"spark",
"properties":{
"maximizeResourceAllocation":"true"
},
"configurations":[
]
}
]
Run Code Online (Sandbox Code Playgroud)
根据我的阅读,设置maximizeResourceAllocation属性应告诉EMR配置Spark以充分利用群集上的所有可用资源.即,我应该有~75GB的内存......那么为什么我会得到~10.5GB的内存限制错误?这是我正在运行的代码:
def sessionize(raw_data, timeout):
# https://www.dataiku.com/learn/guide/code/reshaping_data/sessionization.html
window = (pyspark.sql.Window.partitionBy("user_id", "site_id")
.orderBy("timestamp"))
diff = (pyspark.sql.functions.lag(raw_data.timestamp, 1)
.over(window))
time_diff = (raw_data.withColumn("time_diff", raw_data.timestamp - diff)
.withColumn("new_session", pyspark.sql.functions.when(pyspark.sql.functions.col("time_diff") >= timeout.seconds, 1).otherwise(0)))
window = (pyspark.sql.Window.partitionBy("user_id", "site_id")
.orderBy("timestamp")
.rowsBetween(-1, 0))
sessions = (time_diff.withColumn("session_id", pyspark.sql.functions.concat_ws("_", "user_id", "site_id", pyspark.sql.functions.sum("new_session").over(window))))
return sessions
def aggregate_sessions(sessions):
median = pyspark.sql.functions.udf(lambda x: statistics.median(x))
aggregated = sessions.groupBy(pyspark.sql.functions.col("session_id")).agg(
pyspark.sql.functions.first("site_id").alias("site_id"),
pyspark.sql.functions.first("user_id").alias("user_id"),
pyspark.sql.functions.count("id").alias("hits"),
pyspark.sql.functions.min("timestamp").alias("start"),
pyspark.sql.functions.max("timestamp").alias("finish"),
median(pyspark.sql.functions.collect_list("foo")).alias("foo"),
)
return aggregated
spark_context = pyspark.SparkContext(appName="process-raw-data")
spark_session = pyspark.sql.SparkSession(spark_context)
raw_data = spark_session.read.csv(sys.argv[1],
header=True,
inferSchema=True)
# Windowing doesn't seem to play nicely with TimestampTypes.
#
# Should be able to do this within the ``spark.read.csv`` call, I'd
# think. Need to look into it.
convert_to_unix = pyspark.sql.functions.udf(lambda s: arrow.get(s).timestamp)
raw_data = raw_data.withColumn("timestamp",
convert_to_unix(pyspark.sql.functions.col("timestamp")))
sessions = sessionize(raw_data, SESSION_TIMEOUT)
aggregated = aggregate_sessions(sessions)
aggregated.foreach(save_session)
Run Code Online (Sandbox Code Playgroud)
基本上,只不过是窗口和groupBy来聚合数据.
它从一些错误开始,并停止增加相同错误的数量.
我已经尝试使用--conf spark.yarn.executor.memoryOverhead运行spark-submit,但这似乎也没有解决问题.
Duf*_*uff 53
我感觉到你的痛苦......
我们在YARN上使用Spark时遇到了类似的内存耗尽问题.我们有五个64GB,16个核心虚拟机,无论我们设置什么spark.yarn.executor.memoryOverhead,我们都无法为这些任务获得足够的内存 - 无论我们给他们多少内存,他们最终都会死掉.这是一个相对直接的Spark应用程序,导致这种情况发生.
我们发现虚拟机上的物理内存使用率非常低,但虚拟内存使用率非常高.我们建立yarn.nodemanager.vmem-check-enabled在yarn-site.xml对false我们的集装箱不再被杀害了,并且应用程序表现为预期工作.
做了更多的研究,我找到了为什么会发生这种情况的答案:https: //www.mapr.com/blog/best-practices-yarn-resource-management
由于操作系统行为,在Centos/RHEL 6上存在积极的虚拟内存分配,因此应禁用虚拟内存检查程序或将yarn.nodemanager.vmem-pmem-ratio增加到相对较大的值.
该页面链接到IBM非常有用的页面:https: //www.ibm.com/developerworks/community/blogs/kevgrig/entry/linux_glibc_2_10_rhel_6_malloc_may_show_excessive_virtual_memory_usage?lang = en
总之,glibc> 2.10改变了它的内存分配.虽然分配的大量虚拟内存不是世界末日,但它不能与YARN的默认设置一起使用.
yarn.nodemanager.vmem-check-enabled您也可以将MALLOC_ARENA_MAX环境变量设置为较小的数字,而不是设置为false hadoop-env.sh.
我建议阅读两个页面 - 信息非常方便.
lou*_*ton 15
如果你没有使用spark-submit,并且你正在寻找另一种方法来指定Duffyarn.nodemanager.vmem-check-enabled提到的参数,这里有另外两种方法:
如果您使用的是JSON配置文件(传递给AWS CLI或boto3脚本),则必须添加以下配置:
[{
"Classification": "yarn-site",
"Properties": {
"yarn.nodemanager.vmem-check-enabled": "false"
}
}]
Run Code Online (Sandbox Code Playgroud)
如果使用EMR控制台,请添加以下配置:
classification=yarn-site,properties=[yarn.nodemanager.vmem-check-enabled=false]
Run Code Online (Sandbox Code Playgroud)
看到,
我在一个现在正在工作的庞大集群中遇到了同样的问题.向工作人员添加内存不会解决问题.有时在进程聚合中,spark将使用比它更多的内存,并且spark作业将开始使用堆外内存.
一个简单的例子是:
如果你有一个你需要的数据集reduceByKey,有时会在一个worker中聚集更多的数据而不是其他的,如果这个数据占用了一个worker的内存,你会收到该错误消息.
spark.yarn.executor.memoryOverhead如果您为工作人员设置了50%的内存,那么添加该选项将有所帮助(仅用于测试,看看它是否有效,您可以通过更多测试添加更少的内容).
但您需要了解Spark如何与群集中的内存分配配合使用:
关于内存分配的一件好事,如果你没有在执行中使用缓存,你可以设置火花来使用该sotorage空间来处理执行以避免部分出现OOM错误.正如你在spark的文档中看到的那样:
该设计确保了几种理想的特性.首先,不使用缓存的应用程序可以使用整个空间执行,从而避免不必要的磁盘溢出.其次,使用缓存的应用程序可以保留最小存储空间(R),其中数据块不受驱逐.最后,这种方法为各种工作负载提供了合理的开箱即用性能,而无需用户内部划分内存的专业知识.
但是我们怎么能用呢?
您可以更改某些配置,将MemoryOverhead配置添加到作业调用中,但也可以考虑添加此配置:spark.memory.fraction更改为0.8或0.85并将其降低spark.memory.storageFraction到0.35或0.2.
其他配置可以提供帮助,但需要检查您的情况.在这里完成所有这些配置.
现在,在我的案例中有什么帮助.
我有一个拥有2.5K工作者和2.5TB RAM的集群.而且我们面临着像你这样的OOM错误.我们只增加到spark.yarn.executor.memoryOverhead2048.我们启用动态分配.当我们调用这个工作时,我们没有为工人设置内存,我们将其留给Spark来决定.我们只是设置了开销.
但对于我的小集群的一些测试,改变执行和存储内存的大小.这解决了这个问题.
尝试重新分区。它适用于我的情况。
加载write.csv(). 数据文件总共有 10 MB 左右,对于执行程序中的每个处理任务,可能需要说总共有几个 100 MB 的内存。我当时检查了分区数为 2。然后它在与其他表连接的以下操作中像滚雪球一样增长,添加新列。然后我在某个步骤遇到了内存超出限制的问题。我检查了分区数,它仍然是 2,我猜是从原始数据帧派生的。所以我一开始就尝试重新分区,现在没有问题了。
我还没有读过很多关于 Spark 和 YARN 的材料。我所知道的是节点中有执行者。执行器可以根据资源处理许多任务。我的猜测是一个分区将被原子映射到一项任务。它的体积决定了资源的使用。如果一个分区变得太大,Spark 将无法对其进行切片。
一个合理的策略是先确定节点和容器内存,10GB 或 5GB。理想情况下,两者都可以服务于任何数据处理工作,只是时间问题。给定 5GB 内存设置,你找到的一个分区的合理行,比如测试后是 1000(在处理过程中它不会失败任何步骤),我们可以用下面的伪代码来做:
RWS_PER_PARTITION = 1000
input_df = spark.write.csv("file_uri", *other_args)
total_rows = input_df.count()
original_num_partitions = input_df.getNumPartitions()
numPartitions = max(total_rows/RWS_PER_PARTITION, original_num_partitions)
input_df = input_df.repartition(numPartitions)
Run Code Online (Sandbox Code Playgroud)
希望能帮助到你!