【问题标题】:Spark: Get elements of an RDD based on the elements of an array in another RDDSpark:根据另一个RDD中数组的元素获取RDD的元素
【发布时间】:2016-09-16 07:59:12
【问题描述】:

在 Spark Scala 框架中,我有一个 RDD,rdd1,其中每个元素代表矩阵 A 的单个元素:

val rdd1 = dist.map{case (((x,y),z,v)) => ((x,y),v)}

x 代表行, y 代表列和 v 表示矩阵A 中的值。

我还有另一个RDD,rdd2,形式为RDD[index, Array[(x, y)]],其中每个元素中的数组表示矩阵A的元素集合,存储在rdd1中,具体需要index 在该元素中表示。

现在我需要做的是获取每个index 的矩阵A 元素的值,保留包括index(x,y)v 在内的所有数据。这样做的好方法是什么?

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    如果我理解正确,您的问题归结为:

    val valuesRdd = sc.parallelize(Seq(
    //((x, y), v)
      ((0, 0), 5.5),            
      ((1, 0), 7.7)
    ))
    
    val indicesRdd = sc.parallelize(Seq(
    //(index, Array[(x, y)])
      (123, Array((0, 0), (1, 0))) 
    ))
    

    并且您想合并这些 RDD 以获取所有值 (index, (x, y), v),在本例中为 (123, (0,0), 5.5)(123, (1,0), 7.7)

    您绝对可以使用join 执行此操作,因为两个RDD 都有一个公共列(x, y),但由于其中一个实际上有一个Array[(x, y)],您必须先将其分解为一组行:

    val explodedIndices = indicesRdd.flatMap{case (index, coords: Array[(Int, Int)]) => coords.map{case (x, y) => (index, (x, y))}}
    // Each row exploded into multiple rows (index, (x, y))
    
    val keyedIndices = explodedIndices.keyBy{case (index, (x, y)) => (x, y)}
    // Each row keyed by the coordinates (x, y)
    
    val keyedValues = valuesRdd.keyBy{case ((x, y), v) => (x, y)}
    // Each row keyed by the coordinates (x, y)
    
    // Because we have common keys, we can join!
    val joined = keyedIndices.join(keyedValues)
    

    【讨论】:

    • 谢谢。它发出有关 flatMap 语句中使用的 _ 的错误:missing parameter of type for expanded function...
    • 好的。它进行了一些修改:val explodedIndices = qual.flatMap{case (index, coords: Array[(Long, Long)]) => coords.map{case(x,y) => (index, (x, y))}}。谢谢。
    • 太棒了!修复了答案,我实际上并没有尝试运行它。
    猜你喜欢
    • 1970-01-01
    • 2017-07-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-06-20
    • 2019-07-26
    • 1970-01-01
    相关资源
    最近更新 更多