【发布时间】:2018-03-30 22:51:24
【问题描述】:
我目前正在开发一个 Spark (v 2.2.0) Streaming 应用程序,并且在 Spark 似乎在整个集群中分配工作的方式上遇到了问题。此应用程序使用客户端模式提交到 AWS EMR,因此有一个驱动程序节点和几个工作程序节点。这是 Ganglia 的屏幕截图,显示了过去一小时的内存使用情况:
最左边的节点是“master”或“driver”节点,另外两个是worker节点。与通过流进入的工作负载相对应的所有三个节点的内存使用量都有峰值,但峰值不相等(即使缩放到内存使用百分比)。当一个大的工作负载进来时,驱动节点似乎过度工作,并且作业将崩溃并出现有关内存的错误:
OpenJDK 64-Bit Server VM warning: INFO: os::commit_memory(0x000000053e980000, 674234368, 0) failed; error='Cannot allocate memory' (errno=12)
我也遇到过:
Exception in thread "streaming-job-executor-10" java.lang.OutOfMemoryError: Java heap space 当 master 内存不足时,这同样令人困惑,因为我的理解是“客户端”模式不会使用驱动程序/主节点作为执行程序。
相关细节:
- 如前所述,本申请以客户端模式提交:
spark-submit --deploy-mode client --master yarn ...。 - 我在程序中没有运行
collect或coalesce - 我怀疑在单个节点上运行的任何工作(
jdbc主要读取)在完成后是repartition'd。 - 内存中有几个非常非常小的数据集
persist。 - 1 x 驱动程序规格:4 核,16GB RAM(m4.xlarge 实例)
- 2 x Worker 规格:4 核,30.5GB RAM(r3.xlarge 实例)
- 我尝试过允许 Spark 选择执行器大小/内核并手动指定它们。两种情况的行为相同。 (我手动指定了 6 个执行器,1 个核心,9GB RAM)
我肯定在这里不知所措。我不确定代码中发生了什么会触发驱动程序占用这样的工作量。
我能想到的唯一嫌疑人是类似于以下的代码sn-p:
val scoringAlgorithm = HelperFunctions.scoring(_: Row, batchTime)
val rawScored = dataToScore.map(scoringAlgorithm)
这里,正在从静态对象加载一个函数,并用于映射Dataset。据我了解,Spark 将在整个集群中序列化此函数(回复:http://spark.apache.org/docs/2.2.0/rdd-programming-guide.html#passing-functions-to-spark)。但是,也许我弄错了,它只是在驱动程序上运行此转换。
如果有人对此问题有任何见解,我很想听听!
【问题讨论】:
标签: apache-spark spark-streaming hadoop-yarn