Spark-Python:筛选datos_acumulados超阈值行并保留空值分组的问题
PySpark按分组筛选阈值记录保留全部分组的解决方案
问题根因
你的代码中grupoDistinctDF没有对grupo_edad字段去重,保留的是原始表所有明细行,左连接后会产生大量无效重复数据,最终输出时grupo_edad=4的行被隐式过滤。
修正后代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 提取所有唯一的年龄分组作为维度表 grupo_distinct_df = df_datos_acumulados.select("grupo_edad").distinct() # 定义窗口:按年龄分组、日期升序排序,给每个分组内的记录按时间编号 grupo_window = Window.partitionBy("grupo_edad").orderBy("fecha") # 筛选每个分组中第一个满足累计值大于等于20480的记录 matched_df = df_datos_acumulados \ .withColumn("row_num", F.row_number().over(grupo_window)) \ .filter(F.col("datos_acumulados") >= 20480) \ .withColumn("first_match_row", F.min("row_num").over(Window.partitionBy("grupo_edad"))) \ .filter(F.col("row_num") == F.col("first_match_row")) \ .drop("row_num", "first_match_row") # 用全部分组左关联匹配到的结果,未匹配到的字段自动为null result_df = grupo_distinct_df.join(matched_df, on="grupo_edad", how="left") # 输出结果 result_df.orderBy("grupo_edad").show()
输出验证
执行后输出结果和预期一致:
+----------+----------+------------+----------------+ |grupo_edad| fecha|acumuladosMB|datos_acumulados| +----------+----------+------------+----------------+ | 1|2020-08-04| 4864| 20921| | 4| null| null| null| +----------+----------+------------+----------------+
内容的提问来源于stack exchange,提问作者Gaby
相关产品推荐
相关产品推荐

