PySpark如何基于多列值按dotId分组聚合拼接codePp字段?
实现代码
直接上可运行的PySpark代码示例:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, concat, lit, concat_ws, collect_list # 初始化SparkSession spark = SparkSession.builder.appName("codePp_concat").getOrCreate() # 构造示例数据,实际使用时替换为你自己的数据源读取逻辑即可 data = [ ("dot0001", "Pp3523", "start"), ("dot0001", "Pp3524", "stop"), ("dot0020", "Pp3522", "start"), ("dot0020", "Pp3556", "stop"), ("dot9999", "Pp3545", "stop"), ("dot9999", "Pp3523", "start"), ("dot9999", "Pp3587", "stop"), ("dot9999", "Pp3567", "start") ] df = spark.createDataFrame(data, schema=["dotId", "codePp", "status"]) # 核心处理逻辑 result_df = df.withColumn("processed_code", when(col("status") == "stop", concat(col("codePp"), lit("(stop)"))) .otherwise(col("codePp"))) \ .groupBy("dotId") \ .agg(concat_ws(", ", collect_list("processed_code")).alias("codePp")) # 输出结果验证 result_df.show(truncate=False)
逻辑说明
- 先用
when函数做单行值处理:status为stop的行给codePp拼接对应后缀,其余行保留原codePp值 - 按dotId分组后,用
collect_list把同组处理好的codePp收集为数组,再用concat_ws按,分隔符拼接成单个字符串,输出完全符合要求 - 如果有保持原行顺序拼接的需求,默认
collect_list会按数据原始顺序收集,不需要额外调整;如果有自定义排序要求,可以在分组前先对全局做排序即可
内容的提问来源于stack exchange,提问作者YoYoDev
相关产品推荐
相关产品推荐

