【问题标题】:Why does RDD.groupBy return an empty RDD if the initial RDD wasn't empty?如果初始 RDD 不为空,为什么 RDD.groupBy 返回一个空 RDD?
【发布时间】:2015-09-30 19:04:54
【问题描述】:

我有一个用来加载二进制文件的 RDD。每个文件被分成多个部分并进行处理。在处理步骤之后,每个条目是:

(filename, List[Results])

由于文件被分成几个部分,因此 RDD 中的多个条目的文件名相同。我正在尝试使用 reduceByKey 将每个部分的结果重新组合在一起。但是,当我尝试对此 RDD 进行计数时,它返回 0:

val reducedResults = my_rdd.reduceByKey((resultsA, resultsB) => resultsA ++ resultsB)
reducedResults.count() // 0

我尝试更改它使用的密钥,但没有成功。即使非常简单地尝试对结果进行分组,我也没有得到任何输出。

val singleGroup = my_rdd.groupBy((k, v) => 1) 
singleGroup.count() // 0

另一方面,如果我只是收集结果,那么我可以将它们分组到 Spark 之外,一切正常。但是,我仍然需要对收集的结果进行额外的处理,所以这不是一个好的选择。

如果初始 RDD 不为空,什么会导致 groupBy/reduceBy 命令返回空 RDD?

【问题讨论】:

  • @zero323 我使用的是 Spark 1.2.0,所以我不能使用那个库。
  • 我的错误网址:stackoverflow.com/help/mcve
  • 从您在此处发布的内容来看,my_rdd 似乎是空的 - 请通过提供您的数据样本来说服我,以便我们尝试重现您的问题。

标签: scala apache-spark rdd


【解决方案1】:

事实证明,我为该特定作业生成 Spark 配置的方式存在错误。而不是将spark.default.parallelism 字段设置为合理的值,而是将其设置为 0。

来自spark.default.parallelism 上的 Spark 文档:

当用户未设置时,由 join、reduceByKey 和并行化等转换返回的 RDD 中的默认分区数。

因此,虽然像collect() 这样的操作工作得非常好,但任何在不指定分区数量的情况下重新洗牌的尝试都会给我一个空 RDD。这将教会我信任旧的配置代码。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多