【发布时间】: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