【问题标题】:How to create a nested forloop on a RDD (Scala)如何在 RDD (Scala) 上创建嵌套的 for 循环
【发布时间】:2018-03-24 19:24:21
【问题描述】:

我有一个具有以下结构的 RDD:
((ByteArray, Idx), ((srcIdx,srcAdress), (destIdx,destAddress)))

这是比特币区块链边缘(交易)的表示。 (ByteArray, Idx) 可以看作是一个标识符,剩下的就是一个边。我的最终目标是在区块链的图形表示中聚合节点。为此我需要对结构进行的第一个修改是将位于同一比特币交易中的源放在一个边缘(最终在一个节点中)。通过这种方式,我将“聚集”属于同一用户的公钥。 此修改的结果将具有以下结构:
((ByteArray, Idx), (List((srcIdx, srcAddress)), (destIdx, destAddress)))
或者以任何其他形式具有相同的功能(例如,如果这在 Scala 中是不可能的或不符合逻辑的)。

我目前的思考过程如下。在 Java 中,我会对 RDD 中的项目进行嵌套 for 循环,每个循环都会为具有相同键 ((ByteArray, Idx)) 的那些项目创建一个列表。在此之后删除任何重复项。 但是,由于我正在处理 RDD 和 Scala,所以这是不可能的。接下来,我尝试在我的 RDD 上执行 .collect() 和单独的 .map() 函数,使用要在我的 map 函数中循环的集合。然而,Spark 不喜欢这样,因为显然集合不能被序列化。 接下来我尝试创建一个“嵌套”映射函数,如下所示:

val aggregatedTransactions = joinedTransactions.map( f => {
  var list = List[Any](f._2._1)

  val filtered = joinedTransactions.filter(t => f._1 == t._1)

  for(i <- filtered){
    list ::= i._2._1
  }

  (f._1, list, f._2._2)
})

这是不允许的,因为过滤器(或映射)功能在 .map() 中不可用。有哪些替代方案?

我对 Scala 还很陌生,因此非常感谢任何有用的背景信息。

【问题讨论】:

  • 我认为,鉴于您的问题的性质,提供输入+输出示例以避免误解会很有用

标签: scala apache-spark rdd


【解决方案1】:

嵌套 RDD 是不可能的,但是 RDD 中的集合是 可能。

您可以使用嵌套的 for 循环 cartesian

def笛卡尔[U](其他:RDD[U])(隐式arg0:ClassTag[U]):RDD[(T, U)] 永久链接 返回此 RDD 和另一个的笛卡尔积 一,即a所在的所有元素对(a, b)的RDD this 和 b 在 other 中。

val nestedForRDD = rdd1.cartesian(rdd2)

nestedForRDD.map((rdd1TypeVal, rdd2TypeVal) => {
  //Do your inner-nested evaluation code here
})

使用 Spark SQL 也可以实现。

http://bigdatums.net/2016/02/12/how-to-extract-nested-json-data-in-spark/

【讨论】:

    【解决方案2】:

    我的最终目标是在区块链的图形表示中聚合节点。为此我需要对结构进行的第一个修改是将位于同一比特币交易中的资源放在一个边缘(最终在一个节点中)。

    所以本质上你想groupByKey:

    joinedTransactions.groupByKey().map {
       // process data to get desired shape
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-07-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-11-02
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多