如何在Spark中通过筛选操作生成NodeId节点对?
解决Spark Scala中生成同类型NodeId节点对的问题
嘿,我来帮你搞定这个需求!首先咱们得明确核心目标:只筛选出typeFrom和typeTo取值相同的记录,然后把这些记录里同类型的NodeId两两组成节点对。下面我用具体的代码示例来一步步实现:
第一步:先过滤符合条件的记录
首先,我们需要把typeFrom等于typeTo的行筛选出来,因为只有这些行才需要生成节点对。假设你的DataFrame结构包含NodeId、typeFrom、typeTo这三列,代码如下:
import org.apache.spark.sql.functions._ // 先模拟你的原始DataFrame(你可以替换成自己的实际DF) val originalDF = spark.createDataFrame(Seq( (1, "A", "A"), (2, "A", "A"), (3, "B", "C"), (4, "B", "B"), (5, "B", "B") )).toDF("NodeId", "typeFrom", "typeTo") // 过滤出typeFrom和typeTo相等的记录 val validDF = originalDF.filter($"typeFrom" === $"typeTo")
第二步:生成同类型的NodeId节点对
接下来,我们通过自连接的方式,把同一类型下的所有NodeId两两配对。自连接的条件是两个表的typeFrom(因为已经过滤过,typeFrom=typeTo)相等:
// 自连接生成节点对 val nodePairsDF = validDF.as("df1") .join(validDF.as("df2"), $"df1.typeFrom" === $"df2.typeFrom") .select( $"df1.NodeId".alias("NodeId_Source"), $"df2.NodeId".alias("NodeId_Target"), $"df1.typeFrom".alias("Common_Type") )
可选:排除自身配对
如果你不需要NodeId和自己组成的配对(比如(1,1)这种),可以再加一层过滤:
val nodePairsWithoutSelf = nodePairsDF.filter($"NodeId_Source" < $"NodeId_Target")
示例输出
拿我们模拟的原始数据来说,validDF里的记录是:
| NodeId | typeFrom | typeTo |
|---|---|---|
| 1 | A | A |
| 2 | A | A |
| 4 | B | B |
| 5 | B | B |
生成的nodePairsDF会是:
| NodeId_Source | NodeId_Target | Common_Type |
|---|---|---|
| 1 | 1 | A |
| 1 | 2 | A |
| 2 | 1 | A |
| 2 | 2 | A |
| 4 | 4 | B |
| 4 | 5 | B |
| 5 | 4 | B |
| 5 | 5 | B |
如果加上排除自身的过滤,nodePairsWithoutSelf的结果就是:
| NodeId_Source | NodeId_Target | Common_Type |
|---|---|---|
| 1 | 2 | A |
| 4 | 5 | B |
另一种实现方式:用窗口函数收集后展开
如果你觉得自连接有点重,也可以用窗口函数先收集同类型的所有NodeId,再通过explode和数组操作生成配对:
// 用窗口函数收集同类型的所有NodeId val groupedDF = validDF .withColumn("all_node_ids", collect_list($"NodeId").over(Window.partitionBy($"typeFrom"))) // 生成节点对 val pairDF = groupedDF .withColumn("NodeId_Target", explode($"all_node_ids")) .select($"NodeId".alias("NodeId_Source"), $"NodeId_Target", $"typeFrom".alias("Common_Type")) // 同样可以加过滤排除自身配对
这种方式和自连接的效果是一样的,你可以根据自己的数据集大小选择更高效的方式~
内容的提问来源于stack exchange,提问作者Markus
相关产品推荐
相关产品推荐

