PySpark如何按条件筛选DataFrame并计算连续符合条件组的最大值
PySpark实现连续条件分组取最大值方案
实现思路
- 第一步:给原始数据添加顺序行号,保证遍历顺序和原始数据一致
- 第二步:用累加窗口函数生成连续Condition=1的分组标记:遇到Condition=0时分组标记累加1,所有连续的Condition=1行将归属同一分组
- 第三步:过滤掉Condition=0的无效行,按生成的分组标记聚合
- 第四步:每个分组内计算最大值、拼接计算依据字符串,匹配最大值对应的Identifiant
- 第五步:按要求输出结果字段
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("ContinuousGroupMax").getOrCreate() # 构造原始数据 raw_data = [ ("ID1", 16, 1), ("ID2", 8, 1), ("ID3", 4, 0), ("ID4", 5, 0), ("ID5", 6, 1), ("ID6", 10, 1), ("ID7", 9, 1), ("ID8", 8, 1), ("ID9", 9, 0), ("ID10", 11, 0), ("ID11", 6, 1), ("ID12", 8, 1), ("ID13", 10, 0), ("ID13", 12, 1), ("ID14", 15, 0), ("ID15", 14, 1), ("ID16", 8, 1), ("ID17", 9, 1) ] df = spark.createDataFrame(raw_data, schema=["Identifiant", "Value", "Condition"]) # 生成全局顺序行号,匹配原始数据遍历顺序 df_with_order = df.withColumn("order_seq", F.regexp_extract("Identifiant", r"ID(\d+)", 1).cast("int")) # 生成连续Condition=1的分组标记 window_order = Window.orderBy("order_seq") df_with_group = df_with_order.withColumn( "group_flag", F.sum(F.when(F.col("Condition") == 0, 1).otherwise(0)).over(window_order) ) # 过滤无效行后按分组聚合 df_grouped = df_with_group.filter(F.col("Condition") == 1).groupBy("group_flag").agg( F.max("Value").alias("max_value"), F.concat_ws(",", F.collect_list("Value")).alias("value_list"), F.collect_list(F.struct("Value", "Identifiant")).alias("id_value_mapping") ) # 匹配最大值对应ID、拼接计算依据 df_result = df_grouped.withColumn( "Identifiant", F.expr("filter(id_value_mapping, x -> x.Value = max_value)[0].Identifiant") ).withColumn( "计算依据", F.concat(F.lit("max("), F.col("value_list"), F.lit(")")) ).select("Identifiant", F.col("max_value").alias("Value"), "计算依据") # 按分组顺序输出结果 df_result.orderBy("group_flag").show(truncate=False)
输出结果
| Identifiant | Value | 计算依据 |
|---|---|---|
| ID1 | 16 | max(16,8) |
| ID6 | 10 | max(6,10,9,8) |
| ID12 | 8 | max(6,8) |
| ID13 | 12 | max(12) |
| ID15 | 14 | max(14,8,9) |
内容的提问来源于stack exchange,提问作者elokema
相关产品推荐
相关产品推荐

