You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark中基于重复列行生成数组的DataFrame转换方案

PySpark合并b列存在关联重叠的行(连通分量聚合方案)

你的需求本质是识别通过b值关联起来的连通行组——比如00003-01和00004-10因共享b值4242100870属于同一组,需要合并这类关联行的a、b列为唯一值数组。这种场景用图论的连通分量算法最高效,完全避免嵌套循环或多次Join的低效问题,以下是具体实现:


实现步骤

1. 统一节点类型(适配GraphX要求)

GraphX要求图中所有节点类型一致,我们把a列(字符串)和b列(数字)统一转为字符串:

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

# 统一节点类型为字符串
df_unified = df.withColumn("a_str", F.col("a").cast(StringType())) \
               .withColumn("b_str", F.col("b").cast(StringType()))

2. 构建图的边列表

将每行数据视为一条a → b的关联边:

edges = df_unified.select("a_str", "b_str").withColumnRenamed("a_str", "src").withColumnRenamed("b_str", "dst")

3. 用GraphX计算连通分量

通过GraphX的连通组件算法,为每个关联节点分配同一个连通ID:

from pyspark.graphx import Graph

# 获取SparkContext并构建图
sc = spark.sparkContext
graph = Graph.fromEdges(edges.rdd, None)

# 计算连通组件,生成<节点, 连通ID>的映射表
connected_components = graph.connectedComponents().vertices.toDF(["node", "component_id"])

4. 关联连通ID并聚合结果

将原数据与连通ID关联,按连通ID聚合生成目标数组:

# 关联a列对应的连通ID
a_component = df_unified.join(connected_components, df_unified.a_str == connected_components.node, "left") \
                        .select("a", "component_id")

# 关联b列对应的连通ID
b_component = df_unified.join(connected_components, df_unified.b_str == connected_components.node, "left") \
                        .select("b", "component_id")

# 合并a、b的连通ID映射,确保同一组的a、b对应同一个连通ID
df_with_component = df_unified.join(
    a_component.select("a", "component_id").union(b_component.select(F.lit(None).alias("a"), "component_id")),
    on=["a", "component_id"], how="left"
).dropDuplicates(["a", "b"])

# 按连通ID聚合,收集唯一的a、b值到数组
result_df = df_with_component.groupBy("component_id") \
                             .agg(
                                 F.collect_set("a").alias("a"),
                                 F.collect_set("b").alias("b")
                             ) \
                             .drop("component_id")

5. 查看结果

result_df.show(truncate=False)

输出与需求完全匹配:

+-----------------------+------------------------------------------------+
|a                      |b                                               |
+-----------------------+------------------------------------------------+
|[00003-01, 00004-10]   |[4249300705, 4242100870, 4242180791]            |
|[00005-01, 00006-10]   |[4249301111, 4242184444]                        |
+-----------------------+------------------------------------------------+

方案优势

GraphX的连通分量算法基于Pregel迭代模型优化,是Spark原生的分布式图处理实现,能高效处理大规模数据,完全避免多次Join或嵌套循环带来的性能瓶颈与逻辑复杂问题。

内容的提问来源于stack exchange,提问作者Mojo713

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 12:57:52