【发布时间】:2015-09-26 00:25:12
【问题描述】:
我有一个 cassandra 数据库,其中包含大约 400 万条记录。我有 3 个从机和一个驱动程序。我想将此数据加载到火花内存中并对其进行处理。当我执行以下操作时,它会读取一台从机(6 Gb 中的 300 mb)中的所有数据,并且所有其他从机内存都未使用。我将数据帧重新分配为 3,但数据仍然存在于一台机器上。因此,由于每个作业都在一台机器上执行,因此处理数据需要大量时间。这就是我正在做的事情
val tabledf = _sqlContext.read.format("org.apache.spark.sql.cassandra").options(Map( "table" -> "events", "keyspace" -> "sams")).load
tabledf.registerTempTable("tempdf");
_sqlContext.cacheTable("tempdf");
val rdd = _sqlContext.sql(query);
val partitionedRdd = rdd.repartition(3)
val count = partitionedRdd.count.toInt
当我对 partitionedRdd 执行一些操作时,它只在一台机器上执行,因为所有数据都只存在于一台机器上
更新 我在配置中使用它--conf spark.cassandra.input.split.size_in_mb=32,我的所有数据仍然加载到一个执行器中
更新 我正在使用 spark 1.4 版和 spark cassandra 连接器 1.4 版发布
【问题讨论】:
-
你确定你的配置是正确的并且你在某处没有
val conf = new SparkConf().setMaster("local[*]")吗? -
不,我在集群模式下运行,Web UI 显示 3 个从机。我也在使用这个配置 spark.cassandra.input.split.size_in_mb=67108864
-
stackoverflow.com/questions/31583249/…,这就是我使用67108864的原因
-
哦,抱歉 - 现在是早上,我没有看到
rdd.repartition。我想你想增加分区的数量。我不知道你的奴隶是哪种类型的实例,但我猜他们有多个计算单元。分区数(您当前设置为 3)至少应为numberOfSlaves*numberOfComputeUnitsOnEachSlave,以便您以最佳方式利用集群。 -
我在从机上有双核 8 GB 机器。计算机单元数是否等于内核数?
标签: apache-spark