PySpark处理50M行DataFrame:提取每本书前10位非空所有者及对应数量
解决方案
问题分析
你当前的代码存在两个核心问题:
- 窗口分区错误:
Window.partitionBy('book','owner')会将每个(book, owner)作为独立分组,导致collect_list只能获取当前所有者的数据,无法聚合同一书籍下的所有所有者信息。 - 未提前过滤无效数据:直接在全量数据上计算窗口,导致
owner为Null的记录被纳入统计,同时全量窗口聚合会带来巨大的内存和性能开销(5000万条记录下尤为明显)。
最优实现方案
要保留原表所有行(包括owner为Null的记录),同时正确计算每本书的前10位有效所有者,可通过先处理有效数据聚合top10,再与原表左连接的方式实现,既保证结果正确性,又兼顾性能:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 定义窗口:按书籍分区,按副本数降序排序 top_window = Window.partitionBy("book").orderBy(F.desc("count")) # 2. 处理有效数据(过滤owner不为Null的行),取每本书的前10位所有者并聚合 top_owners_agg = df.filter(F.col("owner").isNotNull()) \ .withColumn("rank", F.row_number().over(top_window)) \ .filter(F.col("rank") <= 10) \ .groupBy("book") \ .agg( # 将所有者和对应副本数打包成struct列表,也可拆分为两个独立列表 F.collect_list(F.struct("owner", "count")).alias("top_owners_details") # 若需要分开的列表,替换为: # F.collect_list("owner").alias("top_owners"), # F.collect_list("count").alias("top_counts") ) # 3. 与原表左连接,保留所有原始行(包括owner为Null的记录) final_df = df.join(top_owners_agg, on="book", how="left")
方案优势
- 性能高效:先过滤无效数据,仅对有效行计算排名和聚合,大幅减少计算量;
row_number()取前10后再聚合,避免全量窗口聚合带来的内存浪费。 - 结果准确:完全忽略
owner为Null的记录,仅统计有效所有者的top10。 - 兼容下游需求:左连接后保留原表所有行,无需删除任何数据。
内容的提问来源于stack exchange,提问作者pnv
相关产品推荐
相关产品推荐

