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
相关产品推荐
相关产品推荐

