【发布时间】: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