PySpark关联ID分组实现:解决一列多键关联分组难题
PySpark 高效合并关联ID分组
原始数据与DataFrame
data = [ ['7933546917973caa8c2898c834446415', '3ef2e38d48a9af3e096ddd3bc3816afb', 1], ['7d693086c5b8f74cbe881166cf3c2a29', 'fcb907411aff4f44c599cf03d23327c0', 2], ['7e18b452bb1e2845800a71d9431033b6', '9bc9d06e0efb16abde20c35ba36a2f1b', 3], ['7e18b452bb1e2845800a71d9431033b6', 'ff351ada316cbb0f270f935adfd16ad4', 4], ['8240cf1e442a97aa91d1029270728bbb', '484f25e9ab91af2c116cd788c91bdc82', 5], ['8919d5fd5b6fd118c1c6b691c65c9df9', '8dc7dfb4466590375f1aaac7fc8cb987', 6], ['8919d5fd5b6fd118c1c6b691c65c9df9', '9b93e3cfc5605e74ce2ce4c9450fd622', 7], ['8dc7dfb4466590375f1aaac7fc8cb987', '9b93e3cfc5605e74ce2ce4c9450fd622', 8], ['8f459a7cff281bad73f604166841849e', '41f007c0cc45c228e246f1cc91145878', 9], ['99f70106443a6f3f5c69d99a49d22d01', 'be73ca52536d13dfea295d4fcd273fde', 10], ['a9781767ca4fe8fb1282ee003d2c06ac', 'cb6feb2f38731fc7832545cbe2ac881b', 11], ['f4901968c29e928fc7364411b03336d4', '6fa82a51f17f0bf258fe06befc661216', 12], ['f6da014449e6fa82c24d002b4a27b105', '41f007c0cc45c228e246f1cc91145878', 13], ['f6da014449e6fa82c24d002b4a27b105', '8f459a7cff281bad73f604166841849e', 14], ['f93c0028bb26bc9b99fca1db300c2ac1', 'ccce888c5813025e95434d7ceedf1db3', 15], ['ff351ada316cbb0f270f935adfd16ad4', '9bc9d06e0efb16abde20c35ba36a2f1b', 16], ['ffe20a2c61638bb10bf943c42b4d794f', '985e237162ccfc04874664648893c241', 17], ] df = spark.createDataFrame(data, schema=['id1', 'id2', 'grp']) df.show(truncate=False)
期望结果
+------------------------------------------------------------------------------------------------------+---+ |ID |grp| +------------------------------------------------------------------------------------------------------+---+ |[7d693086c5b8f74cbe881166cf3c2a29, fcb907411aff4f44c599cf03d23327c0] |2 | |[7933546917973caa8c2898c834446415, 3ef2e38d48a9af3e096ddd3bc3816afb] |1 | |[8240cf1e442a97aa91d1029270728bbb, 484f25e9ab91af2c116cd788c91bdc82] |5 | |[8dc7dfb4466590375f1aaac7fc8cb987, 9b93e3cfc5605e74ce2ce4c9450fd622, 8919d5fd5b6fd118c1c6b691c65c9df9]|8 | |[8f459a7cff281bad73f604166841849e, 41f007c0cc45c228e246f1cc91145878, f6da014449e6fa82c24d002b4a27b105]|9 | |[99f70106443a6f3f5c69d99a49d22d01, be73ca52536d13dfea295d4fcd273fde] |10 | |[a9781767ca4fe8fb1282ee003d2c06ac, cb6feb2f38731fc7832545cbe2ac881b] |11 | |[f4901968c29e928fc7364411b03336d4, 6fa82a51f17f0bf258fe06befc661216] |12 | |[ffe20a2c61638bb10bf943c42b4d794f, 985e237162ccfc04874664648893c241] |17 | |[ff351ada316cbb0f270f935adfd16ad4, 9bc9d06e0efb16abde20c35ba36a2f1b, 7e18b452bb1e2845800a71d9431033b6]|16 | |[f93c0028bb26bc9b99fca1db300c2ac1, ccce888c5813025e95434d7ceedf1db3] |15 | +------------------------------------------------------------------------------------------------------+---+
核心需求
将所有相互关联的id(无论出现在id1还是id2字段)合并到同一数组中,grp字段保留该关联组对应的任意一个值即可。
当前低效实现
df.alias('df1')\ .join(df.alias('df2'), (F.col('df1.ID1') == F.col('df2.ID2')), 'left')\ .select(F.array_distinct(F.array(F.col('df1.ID1'), F.col('df1.ID2'), F.col('df2.ID1'), F.col('df2.ID2'))).alias('ID'), F.col('df1.grp') )\ .show(truncate=False)
该方法仅能处理单层关联,无法覆盖链式关联(如A-B-C),且多次join会导致数据膨胀,效率极低。
高效解决方案
利用图论中的连通分量算法,通过GraphFrames实现分布式关联分组,适合大规模数据处理:
步骤1:安装并导入GraphFrames
# 若未安装GraphFrames,先执行安装 pip install graphframes
from graphframes import GraphFrame import pyspark.sql.functions as F
步骤2:构建图结构
- 顶点:所有唯一的id(从
id1和id2中收集) - 边:以
id1为源节点,id2为目标节点,保留grp字段
# 生成顶点表 vertices = df.selectExpr("id1 as id").union(df.selectExpr("id2 as id")).distinct() # 生成边表 edges = df.selectExpr("id1 as src", "id2 as dst", "grp")
步骤3:计算连通分量
# 创建图对象 g = GraphFrame(vertices, edges) # 计算每个节点所属的连通组件 connected_components = g.connectedComponents()
步骤4:分组合并ID并关联grp
按组件ID分组,收集所有关联id,再关联该组件对应的任意grp值:
result = connected_components.groupBy("component")\ .agg(F.collect_set("id").alias("ID"))\ # 每个组件保留一个grp值,这里取该组件中出现的第一个grp,可根据需求调整为max/min等 .join(connected_components.select("component", "grp").dropDuplicates(["component"]), on="component")\ .drop("component")\ .select("ID", "grp") result.show(truncate=False)
方案优势
- 基于Spark分布式计算,适合处理大规模数据集
- 自动处理多层链式关联,无需手动迭代join
- 时间复杂度远低于多次join的实现方式
内容的提问来源于stack exchange,提问作者bismi
相关产品推荐
相关产品推荐

