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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 20:36:03