【问题标题】:How to transform RDD[(Key, Value)] into Map[Key, RDD[Value]]如何将 RDD[(Key, Value)] 转换为 Map[Key, RDD[Value]]
【发布时间】:2015-03-22 14:00:28
【问题描述】:

我找了很长时间的解决方案,但没有得到任何正确的算法。

在 scala 中使用 Spark RDD,我如何将 RDD[(Key, Value)] 转换为 Map[key, RDD[Value]],知道我不能使用 collect 或其他可能将数据加载到内存中的方法?

事实上,我的最终目标是按键循环 Map[Key, RDD[Value]] 并为每个 RDD[Value] 调用 saveAsNewAPIHadoopFile

例如,如果我得到:

RDD[("A", 1), ("A", 2), ("A", 3), ("B", 4), ("B", 5), ("C", 6)]

我想要:

Map[("A" -> RDD[1, 2, 3]), ("B" -> RDD[4, 5]), ("C" -> RDD[6])]

我想知道在RDD[(Key, Value)] 的每个键 A、B、C 上使用 filter 是否会花费太多,但我不知道是否多次调用过滤器会有不同的键高效的 ? (当然不是,但也许使用cache ?)

谢谢

【问题讨论】:

  • “知道我不能使用 collect 或其他可能将数据加载到内存中的方法吗?”。这没有意义。无论如何,生成的地图必须适合内存。
  • 只是黑暗中的狂野刺; groupBy(...) 不会给你一些你可以使用的东西吗?它应该给你 RDD[key, Iterable[values]]
  • @thoredge 我不确定迭代是否应该适合大量数据的内存,但实际上根据我的输入量,这可能是一个解决方案

标签: scala bigdata apache-spark rdd


【解决方案1】:

你应该使用这样的代码(Python):

rdd = sc.parallelize( [("A", 1), ("A", 2), ("A", 3), ("B", 4), ("B", 5), ("C", 6)] ).cache()
keys = rdd.keys().distinct().collect()
for key in keys:
    out = rdd.filter(lambda x: x[0] == key).map(lambda (x,y): y)
    out.saveAsNewAPIHadoopFile (...)

一个 RDD 不能是另一个 RDD 的一部分,您无法选择仅收集键并将其相关值转换为单独的 RDD。在我的示例中,您将遍历缓存的 RDD,这没问题并且可以快速运行

【讨论】:

  • 我不确定过滤器的效率,但我认为这是我将实施的解决方案。
  • 没有为你的逻辑准备好转换,恐怕如果你想要更高效的东西你必须自己实现它
  • 这基本上是一个次优的解决方案。您可以满足他的最终目标,即使用 MultipleTextOutput 一次通过每个键写入一个单独的文件。
  • 同意,你可以有另一个解决方案:stackoverflow.com/questions/23995040/…
  • 您在生产中运行此代码时应该注意,因为您正在执行在 master 上运行的收集操作。这可能会导致你的主人很快失去记忆。
【解决方案2】:

听起来您真正想要的是将您的 KV RDD 保存到每个密钥的单独文件中。与其创建Map[Key, RDD[Value]],不如考虑使用MultipleTextOutputFormat similar to the example here. 代码几乎都在示例中。

这种方法的好处是,您可以保证在 shuffle 之后只通过 RDD 一次,并获得您想要的相同结果。如果您按照另一个答案中的建议通过过滤和创建多个 ID 来做到这一点(除非您的源支持下推过滤器),那么您最终会为每个单独的键遍历数据集,这会慢得多。

【讨论】:

    【解决方案3】:

    这是我的简单测试代码。

    val test_RDD = sc.parallelize(List(("A",1),("A",2), ("A",3),("B",4),("B",5),("C",6)))
    val groupby_RDD = test_RDD.groupByKey()
    val result_RDD = groupby_RDD.map{v => 
        var result_list:List[Int] = Nil
        for (i <- v._2) {
            result_list ::= i
        }
        (v._1, result_list)
    }
    

    结果如下

    result_RDD.take(3)
    >> res86: Array[(String, List[Int])] = Array((A,List(1, 3, 2)), (B,List(5, 4)), (C,List(6)))
    

    或者你可以这样做

    val test_RDD = sc.parallelize(List(("A",1),("A",2), ("A",3),("B",4),("B",5),("C",6)))
    val nil_list:List[Int] = Nil
    val result2 = test_RDD.aggregateByKey(nil_list)(
        (acc, value) => value :: acc,
        (acc1, acc2) => acc1 ::: acc2 )
    

    结果是这样的

    result2.take(3)
    >> res209: Array[(String, List[Int])] = Array((A,List(3, 2, 1)), (B,List(5, 4)), (C,List(6)))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-08-08
      • 2021-03-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多