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