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

Spark中使用when()函数构建新列并按id统计增减次数

解决Spark DataFrame按ID统计增长/减少次数的问题

方法1:直接在聚合中使用条件统计

这种方法无需提前新增列,直接在groupBy后的agg里结合sum和when函数完成统计,逻辑更简洁:

from pyspark.sql import functions as F

# 假设原始DataFrame名为df
result_df = df.groupBy("id") \
    .agg(
        # 统计changes_count=1的次数(增长次数)
        F.sum(F.when(F.col("changes_count") == 1, 1).otherwise(0)).alias("increase"),
        # 统计changes_count=2的次数(减少次数)
        F.sum(F.when(F.col("changes_count") == 2, 1).otherwise(0)).alias("decrease")
    )

逻辑说明:when函数在条件满足时返回1,否则返回0,对每个ID分组后用sum累加,得到的就是对应状态的总次数,0值的行会自动被忽略(因为返回0,不影响总和)。

方法2:使用pivot实现列转行统计

如果需要更灵活的状态扩展,可以先过滤掉无变化的行,再用pivot将状态值转为列,最后统计数量:

from pyspark.sql import functions as F

# 先过滤掉changes_count=0的无变化行
filtered_df = df.filter(F.col("changes_count").isin(1, 2))

# 分组后pivot转列,统计次数并重命名列
result_df = filtered_df.groupBy("id") \
    .pivot("changes_count", [1, 2])  # 指定要转成列的状态值,避免自动排序
    .count() \
    .withColumnRenamed("1", "increase") \
    .withColumnRenamed("2", "decrease") \
    .fillna(0)  # 给没有增长/减少记录的ID补0,避免出现null

你之前可能踩的坑

如果用withColumn提前新增了increase和decrease列,但后续groupBy时直接选择这两列而没有用聚合函数,就会导致统计错误。正确的做法是新增列后,在groupBy里对这两列做sum聚合:

# 基于你之前思路的修正写法
temp_df = df.withColumn("increase", F.when(F.col("changes_count") == 1, 1).otherwise(0)) \
            .withColumn("decrease", F.when(F.col("changes_count") == 2, 1).otherwise(0))

result_df = temp_df.groupBy("id") \
    .agg(
        F.sum("increase").alias("increase"),
        F.sum("decrease").alias("decrease")
    )

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:02:22