PySpark实现DataFrame中ID与Link的关联关系构建求助
问题:构建DataFrame中ID与Link的全关联关系
关联逻辑
- 若ID 1关联Link 2,ID 2关联Link 3,则所有关联关系为:1->2、1->3、2->1、2->3、3->1、3->2(即同一连通分量内的所有ID两两互相关联)
- 同理,若1关联4、4关联7、7关联5,则1、4、7、5四个ID两两互相关联,生成所有双向组合
输入DataFrame
+---+----+ | id|link| +---+----+ | 1| 2| | 3| 1| | 4| 2| | 6| 5| | 9| 7| | 9| 10| +---+----+
期望输出DataFrame
+---+----+ | Id|Link| +---+----+ | 1| 2| | 1| 3| | 1| 4| | 2| 1| | 2| 3| | 2| 4| | 3| 1| | 3| 2| | 3| 4| | 4| 1| | 4| 2| | 4| 3| | 5| 6| | 6| 5| | 7| 9| | 7| 10| | 9| 7| | 9| 10| | 10| 7| | 10| 9| +---+----+
已尝试的代码
尝试1:全组合拼接
df = spark.createDataFrame([(1, 2), (3, 1), (4, 2), (6, 5), (9, 7), (9, 10)], ["id", "link"]) ids = df.select("Id").distinct().rdd.flatMap(lambda x: x).collect() links = df.select("Link").distinct().rdd.flatMap(lambda x: x).collect() combinations = [(id, link) for id in ids for link in links] df_combinations = spark.createDataFrame(combinations, ["Id", "Link"]) result = df_combinations.join(df, ["Id", "Link"], "left_anti").union(df).dropDuplicates() result = result.sort(asc("Id"), asc("Link"))
尝试2:交叉连接分组
from pyspark.sql import functions as F from pyspark.sql.window import Window df = spark.createDataFrame([(1, 2), (3, 1), (4, 2), (6, 5), (9, 7), (9, 10)], ["id", "link"]) combinations = df.alias("a").crossJoin(df.alias("b")) \ .filter(F.col("a.id") != F.col("b.id"))\ .select(F.col("a.id").alias("a_id"), F.col("b.id").alias("b_id"), F.col("a.link").alias("a_link"), F.col("b.link").alias("b_link")) window = Window.partitionBy("a_id").orderBy("a_id", "b_link") paths = combinations.groupBy("a_id", "b_link") \ .agg(F.first("b_id").over(window).alias("id")) \ .groupBy("id").agg(F.collect_list("b_link").alias("links")) result = paths.select("id", F.explode("links").alias("link")) result = result.union(df.selectExpr("id as id_", "link as link_"))
解决方案:基于连通分量的全关联生成
你的需求本质是找出所有连通的ID组,然后组内所有ID两两生成双向关联。可以用Spark的GraphFrames库计算连通组件,再生成全组合:
代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from graphframes import GraphFrame # 初始化SparkSession spark = SparkSession.builder.appName("IDLinkAssociation").getOrCreate() # 输入数据 df = spark.createDataFrame([(1, 2), (3, 1), (4, 2), (6, 5), (9, 7), (9, 10)], ["id", "link"]) # 1. 提取所有节点(id和link的去重集合) nodes = df.select("id").union(df.select("link")).distinct().withColumnRenamed("id", "id") # 2. 构建边(原数据的id->link) edges = df.select(F.col("id").alias("src"), F.col("link").alias("dst")) # 3. 创建图并计算连通组件 g = GraphFrame(nodes, edges) connected_components = g.connectedComponents() # 4. 按连通组件分组,收集组内所有节点 grouped_nodes = connected_components.groupBy("component") \ .agg(F.collect_list("id").alias("node_list")) # 5. 生成组内所有两两双向关联(排除自身关联) result = grouped_nodes.withColumn( "pair", F.explode(F.expr("array_zip(node_list, array_repeat(node_list, size(node_list)))")) ).select( F.col("pair.0").alias("Id"), F.col("pair.1").alias("Link") ).filter(F.col("Id") != F.col("Link")) \ .orderBy("Id", "Link") # 查看结果 result.show()
关键说明
- 连通组件计算:GraphFrames的
connectedComponents会给同一连通子图的节点分配相同的component ID,自动将1、2、3、4归为一组,5、6一组,7、9、10一组。 - 生成双向关联:通过
array_zip和explode生成组内所有节点对,过滤自身关联后得到完全匹配期望的双向ID-Link组合。
注意事项
- 需要提前安装GraphFrames:
pip install graphframes,启动Spark时需匹配对应版本的依赖包,例如:--packages graphframes:graphframes:0.8.2-spark3.2-s_2.12(版本根据你的Spark版本调整)。 - 若无法使用GraphFrames,可通过Spark SQL递归CTE实现连通分量计算,原理一致但代码稍复杂。
内容的提问来源于stack exchange,提问作者Avijit
相关产品推荐
相关产品推荐

