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

PySpark中聚合重复键(Port和Windows)记录并拼接字段的方法

在PySpark中聚合拼接重复键对应的字段

要实现按Port和Windows这两个键分组,将重复记录的Desc和Date字段聚合拼接成单个字段,你可以通过分组+聚合函数完成,核心是用collect_list(保留所有重复值)或collect_set(去重)收集字段值,再用concat_ws拼接成字符串。

步骤1:初始化环境并创建示例数据

先构造一个匹配你场景的数据集用于测试:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化SparkSession
spark = SparkSession.builder.appName("AggregateDuplicates").getOrCreate()

# 模拟包含重复键的数据集
sample_data = [
    (80, "Windows Server 2019", "Web Server", "2024-01-01"),
    (80, "Windows Server 2019", "HTTP Service", "2024-01-02"),
    (443, "Windows 10", "SSL Service", "2024-02-01"),
    (443, "Windows 10", "HTTPS Server", "2024-02-02"),
    (22, "Windows Server 2022", "SSH Service", "2024-03-01")
]

df = spark.createDataFrame(sample_data, ["Port", "Windows", "Desc", "Date"])

步骤2:执行聚合拼接

按Port和Windows分组,对Desc和Date分别进行拼接:

# 聚合拼接(保留所有重复值)
aggregated_df = df.groupBy("Port", "Windows") \
    .agg(
        # 用分号+空格拼接所有Desc内容
        F.concat_ws("; ", F.collect_list("Desc")).alias("Aggregated_Desc"),
        # 同理拼接Date内容
        F.concat_ws("; ", F.collect_list("Date")).alias("Aggregated_Date")
    )

# 查看结果
aggregated_df.show(truncate=False)

执行后输出结果:

+----+--------------------+--------------------------+--------------------------+
|Port|Windows             |Aggregated_Desc           |Aggregated_Date           |
+----+--------------------+--------------------------+--------------------------+
|80  |Windows Server 2019|Web Server; HTTP Service  |2024-01-01; 2024-01-02    |
|443 |Windows 10          |SSL Service; HTTPS Server |2024-02-01; 2024-02-02    |
|22  |Windows Server 2022|SSH Service               |2024-03-01                |
+----+--------------------+--------------------------+--------------------------+

可选:去重拼接

如果不需要保留重复的字段值,把collect_list换成collect_set即可:

# 去重后拼接
aggregated_df_distinct = df.groupBy("Port", "Windows") \
    .agg(
        F.concat_ws("; ", F.collect_set("Desc")).alias("Aggregated_Desc_Distinct"),
        F.concat_ws("; ", F.collect_set("Date")).alias("Aggregated_Date_Distinct")
    )

处理空值

如果数据中存在Null值,可以用coalesce替换为空字符串或默认值,避免拼接结果出现null:

# 处理空值的聚合逻辑
aggregated_df_with_null = df.groupBy("Port", "Windows") \
    .agg(
        F.concat_ws("; ", F.collect_list(F.coalesce("Desc", F.lit("N/A")))).alias("Aggregated_Desc"),
        F.concat_ws("; ", F.collect_list(F.coalesce("Date", F.lit("N/A")))).alias("Aggregated_Date")
    )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 21:42:24