【问题标题】:Spark throwing Out of Memory errorSpark抛出内存不足错误
【发布时间】:2014-10-13 18:49:09
【问题描述】:

我有一个带有 8 GB ram 的测试节点,我将仅 10 MB 的数据(来自 csv 文件)加载到 Cassandra(在同一节点本身上)。我正在尝试使用 spark 处理这些数据(在同一节点上运行)。

请注意,对于 SPARK_MEM,我分配了 1 GB 的 RAM,而我分配的 SPARK_WORKER_MEMORY 相同。分配任何额外的内存都会导致 spark 抛出“检查是否所有工作人员都已注册并有足够的内存错误”,这通常表明 Spark 试图寻找额外的内存(根据 SPARK_MEM 和 SPARK_WORKER_MEMORY 属性)并且出现短缺。

当我尝试使用 spark 上下文对象加载和处理 Cassandra 表中的所有数据时,我在处理过程中遇到错误。因此,我尝试使用循环机制从一个表中一次读取大块数据,处理它们并将它们放入另一个表中。

我的源代码结构如下

var data=sc.cassandraTable("keyspacename","tablename").where("value=?",1)
data.map(x=>tranformFunction(x)).saveToCassandra("keyspacename","tablename")

for(i<-2 to 50000){
    data=sc.cassandraTable("keyspacename","tablename").where("value=?",i)
    data.map(x=>tranformFunction(x)).saveToCassandra("keyspacename","tablename")    
}

现在,这可以运行一段时间,大约 200 个循环,然后抛出错误:java.lang.OutOfMemoryError: unable to create a new native thread。

我有两个问题:

Is this the right way to deal with data?
How can processing just 10 MB of data do this to a cluster?

【问题讨论】:

  • 这超出了主题,但即使是42 killobytes of data 也会让您的集群崩溃。 (换句话说,不要假设磁盘上的 10 MB == 内存中的 10 MB,这取决于您的使用模式,数据可能会疯狂增长)
  • 1 UP 用于清除它.. :-),我已经提到了我的基本使用模式(目前)作为代码。只需获取、转换和存储回来(中间只有 1 个中间变量)。
  • 在尝试上面显示的循环机制之前,您能否显示您第一次尝试的代码?您还可以总结一下 transformFunction() 对数据的作用吗?

标签: scala memory-management cassandra apache-spark


【解决方案1】:

您正在 for 循环中运行查询。如果“值”列不是键/索引列,Spark 会将表加载到内存中,然后过滤值。这肯定会导致 OOM。

【讨论】:

  • 但对我来说,“值”列已编入索引。但无论它是否被索引,在内部,查询不会像:select from where token() > 0 and value=?允许过滤。所以,我怀疑spark是否会首先加载整个表然后执行过滤。但是,我可能错了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-06-14
  • 2015-07-08
  • 1970-01-01
  • 2016-03-16
  • 1970-01-01
  • 2011-08-07
相关资源
最近更新 更多