【问题标题】:How do I split a Spark rdd Array[(String, Array[String])]?如何拆分 Spark rdd Array[(String, Array[String])]?
【发布时间】:2016-04-26 08:47:15
【问题描述】:

我正在练习在 Spark shell 中进行排序。我有一个大约 10 列/变量的 rdd。我想根据第 7 列的值对整个 rdd 进行排序。

rdd
org.apache.spark.rdd.RDD[Array[String]] = ...

据我所知,这样做的方法是使用 sortByKey,而它又只适用于对。所以我映射了它,所以我有一对由 column7(字符串值)和完整的原始 rdd(字符串数组)组成

rdd2 = rdd.map(c => (c(7),c))
rdd2: org.apache.spark.rdd.RDD[(String, Array[String])] = ...

然后我申请sortByKey,还是没问题...

rdd3 = rdd2.sortByKey()
rdd3: org.apache.spark.rdd.RDD[(String, Array[String])] = ...

但是现在我如何从 rdd3 (Array[String]) 中分离、收集和保存排序后的原始 rdd?每当我尝试对 rdd3 进行拆分时,都会出现错误:

val rdd4 = rdd3.map(_.split(',')(2))
<console>:33: error: value split is not a member of (String, Array[String])

我在这里做错了什么?还有其他更好的方法来对其中一列的 rdd 进行排序吗?

【问题讨论】:

  • 我不明白你到底想要什么。您的意思是要拆分 Array[String] 中的每个字符串?
  • 你试图拆分元组,这就是错误的原因
  • @John 不,我想拆分 rdd3(一对已排序的 column7 和原始 rdd),所以我会返回原始 rdd 但仍然在第 7 列上排序......实际上没有列7 前缀(如在 rdd3 中)。我稍微编辑了问题,现在更清楚了吗?

标签: scala apache-spark rdd


【解决方案1】:

您对rdd2 = rdd.map(c =&gt; (c(7),c)) 所做的是将其映射到一个元组。 rdd2: org.apache.spark.rdd.RDD[(String, Array[String])] 正如它所说:)。 现在,如果您想拆分记录,则需要从此元组中获取它。 您可以再次映射,只取元组的第二部分(即 Array[String]... 的数组),如下所示:rdd3.map(_._2)

但我强烈建议使用 try rdd.sortBy(_(7)) 或类似的东西。这样你就不需要用元组之类的东西来打扰自己了。

【讨论】:

  • 我尝试了您对 rdd.sortBy(.=>.7) 的建议,但结果显示“错误:需要标识符,但找到了 '=>'”。您可以编辑它以便我接受您的答案吗?正如您所建议的, rdd3.map(_._2) 也可以完成这项工作,但需要做更多的工作。
  • .sortBy(c =&gt; c._7) 不会像.sortBy(_._7) 那样工作,因为 rdd 中的元素具有数组结构。 @KoenDeCouck,我已经发布了我的答案。你可能想检查一下。 :)
  • 它没有,但是@John Titus Jungao 的答案有解决方案:rdd.sortBy(_(7))。我会接受这个答案,因为问题毕竟集中在拆分上,并且您提供了一些关于为什么这不起作用的信息。
  • 当然可以。现在我看到了我的错误。经@JohnTitusJungao 许可,我可以对其进行编辑以供将来参考。
  • @ZahiroMor,当然!前进。很高兴这有帮助。 :)
【解决方案2】:

如果你想使用数组中的第7个字符串对rdd进行排序,你可以直接通过

rdd.sortBy(_(6)) // array starts at 0 not 1

rdd.sortBy(arr => arr(6))

这将为您省去进行多次转换的所有麻烦。 rdd.sortBy(_._7)rdd.sortBy(x =&gt; x._7) 不起作用的原因是因为这不是您访问数组中元素的方式。要访问数组的第 7 个元素,比如arr,您应该使用arr(6)

为了测试这一点,我做了以下操作:

val rdd = sc.parallelize(Array(Array("ard", "bas", "wer"), Array("csg", "dip", "hwd"), Array("asg", "qtw", "hasd")))

// I want to sort it using the 3rd String
val sorted_rdd = rdd.sortBy(_(2))

结果如下:

Array(Array("ard", "bas", "wer"), Array("csg", "dip", "hwd"), Array("asg", "qtw", "hasd"))

【讨论】:

  • 谢谢约翰!这个解决方案看起来是更好的排序方式。我会接受 Zahiro 的回答,但是由于问题的措辞方式,并附有您的解决方案。 (点赞)
【解决方案3】:

这样做:

val rdd4 = rdd3.map(_._2)

【讨论】:

    【解决方案4】:

    我以为你不熟悉 Scala, 所以,下面应该可以帮助您了解更多,

    rdd3.map(kv => {
      println(kv._1) // This represent String 
      println(kv._2) // This represent Array[String]
    })
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-12-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-02-06
      • 2017-01-29
      • 1970-01-01
      相关资源
      最近更新 更多