【发布时间】:2015-09-28 05:08:50
【问题描述】:
我一直在尝试学习如何使用 Apache Spark,但在尝试对 Cassandra 列中的所有值求和时遇到问题(使用 datastax spark-cassandra-connector)。我尝试的所有操作都会导致 java.lang.OutOfMemoryError: Java heap space。
这是我提交给 spark master 的代码:
object Benchmark {
def main( args: Array[ String ] ) {
val conf = new SparkConf()
.setAppName( "app" )
.set( "spark.cassandra.connection.host", "ec2-blah.compute-1.amazonaws.com" )
.set( "spark.cassandra.auth.username", "myusername" )
.set( "spark.cassandra.auth.password", "mypassword" )
.set( "spark.executor.memory", "4g" )
val sc = new SparkContext( conf )
val tbl = sc.cassandraTable( "mykeyspace", "mytable" )
val res = tbl.map(_.getFloat("sclrdata")).sum()
println( "sum = " + res )
}
}
现在我的集群中只有一个 spark 工作节点,而且鉴于表的大小,绝对有可能不是所有的都可以同时放入内存中。但是我不认为这会是一个问题,因为 spark 应该懒惰地评估命令,并且对列中的所有值求和不需要让整个表一次驻留在内存中。
我是这个主题的新手,所以任何关于为什么这不起作用的澄清或关于如何正确执行它的帮助将不胜感激。
谢谢
【问题讨论】:
-
你是完全正确的,不应该将所有内容都加载到内存中。您可以启用调试日志记录以查看拆分大小吗?您使用的是哪个版本的连接器?创建了多少个拆分(火花分区/任务) - 您可以在火花 Web 控制台中看到?您从哪里获得 OOM - 它是在执行程序还是驱动程序应用程序上?
标签: scala cassandra apache-spark datastax spark-cassandra-connector