【问题标题】:Repartitioning of dataframe in spark does not workspark中数据帧的重新分区不起作用
【发布时间】: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


【解决方案1】:

如果“查询”仅访问单个 C* 分区键,您将只能获得单个任务,因为我们还没有办法自动并行获取单个 cassandra 分区。如果您正在访问多个 C* 分区,请尝试进一步缩小输入 split_size 以 mb 为单位。

【讨论】:

  • 是的,我正在尝试使用单个分区键。我尝试使用缓存将数据帧加载到内存中后对其进行重新分区,但这没有帮助。
  • 有没有一种方法可以让我将数据分散到其他机器上,或者我可以对特定列进行索引,以便我可以对该列进行范围查询。
  • 要并行化单个查询,您需要知道分区中的数据并执行并行范围查询
  • 如果单个分区键创建单个分区,那么如果我说 10 个分区键将创建 10 个分区,我没有错。那为什么我们需要 spark.xassandra.input.split.size_in_mb 配置变量
  • 还有在某处记录单个分区键只会在 spark 中创建单个分区
猜你喜欢
  • 1970-01-01
  • 2016-11-07
  • 1970-01-01
  • 2020-03-19
  • 1970-01-01
  • 1970-01-01
  • 2019-07-01
  • 2021-01-03
  • 2021-07-20
相关资源
最近更新 更多