【问题标题】:How to get new RDD from PairRDD based on Key如何根据 Key 从 PairRDD 中获取新的 RDD
【发布时间】:2015-06-07 06:25:30
【问题描述】:

在我的 Spark 应用程序中,我使用了一个 JavaPairRDD<Integer, List<Tuple3<String, String, String>>>,它具有大量数据。

我的要求是我需要一些其他的 RDD JavaRDD<Tuple3<String, String, String>> 来自基于键的大型 PairRDD。

【问题讨论】:

  • 为什么不只过滤base rdd?
  • 使用 java.util.stream.Stream 过滤数据。请看link
  • 在 PairRDD 中,我使用的 List 包含数百万个 Tuple3。但基于元组的第三个参数,我只需要该列表中的 50 条排序记录。所以为此,我只是想创建一些新的 Rdd,而不是在对元组进行排序之后。如果有其他方法,请告诉我。
  • @PrakharAsthana 谢谢,但我使用的是 Java 7,而不是 Java 8。Spark 中还有其他方法吗?
  • 到目前为止你尝试过什么?还请添加当前状态和预期输出的示例。目前尚不清楚您是想要 RDD 的一个元素还是多个元素的合并。

标签: java apache-spark rdd


【解决方案1】:

我不知道 Java API,但这是你在 Scala 中的做法(spark-shell):

def rddByKey[K: ClassTag, V: ClassTag](rdd: RDD[(K, Seq[V])]) = {
  rdd.keys.distinct.collect.map {
    key => key -> rdd.filter(_._1 == key).values.flatMap(identity)
  }
}

您必须为每个键 filter 并将 Lists 与 flatMap 展平。

不得不提的是,这不是一个有用的操作。如果您能够构建原始 RDD,这意味着每个 List 都足够小以适合内存。所以我不明白你为什么要把它们变成 RDD。

【讨论】:

  • 也许我误解了这个问题?告诉我。
猜你喜欢
  • 1970-01-01
  • 2017-10-12
  • 1970-01-01
  • 1970-01-01
  • 2017-12-19
  • 1970-01-01
  • 2015-08-19
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多