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

PySpark处理50M行DataFrame:提取每本书前10位非空所有者及对应数量

解决方案

问题分析

你当前的代码存在两个核心问题:

  1. 窗口分区错误:Window.partitionBy('book','owner')会将每个(book, owner)作为独立分组,导致collect_list只能获取当前所有者的数据,无法聚合同一书籍下的所有所有者信息。
  2. 未提前过滤无效数据:直接在全量数据上计算窗口,导致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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 17:34:59