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

Spark:如何用窗口函数生成拆分列实现大数据框均分?

问题描述

我有一个包含2亿行的DataFrame(DF),无法对其进行分组操作,需要将其拆分为8个各约3000万行的小DF。尝试过拆分方法但未成功:不缓存DF时,拆分后的DF行数总和与原DF不符;若使用缓存,则会出现磁盘空间不足的问题(我的配置为64GB内存+512GB SSD)。

为此我考虑了以下方案:

  • 加载完整的DF
  • 为该DF分配8个随机数
  • 使随机数在DF中均匀分布

示例原DataFrame

+------+--------+
| val1 |  val2  |
+------+--------+
|Paul  |    1.5 |
|Bostap|    1   |
|Anna  |    3   |
|Louis |    4   |
|Jack  |    2.5 |
|Rick  |    0   |
|Grimes|    null|
|Harv  |    2   |
|Johnny|    2   |
|John  |    1   |
|Neo   |    5   |
|Billy |    null|
|James |    2.5 |
|Euler |    null|
+------+--------+

该DF共有14行,我希望通过窗口函数生成如下带sep列的DF:

目标DataFrame

+------+--------+----+
| val1 |  val2  | sep|
+------+--------+----+
|Paul  |    1.5 |1   |
|Bostap|    1   |1   |
|Anna  |    3   |1   |
|Louis |    4   |1   |
|Jack  |    2.5 |1   |
|Rick  |    0   |1   |
|Grimes|    null|1   |
|Harv  |    2   |2   |
|Johnny|    2   |2   |
|John  |    1   |2   |
|Neo   |    5   |2   |
|Billy |    null|2   |
|James |    2.5 |2   |
|Euler |    null|2   |
+------+--------+----+

之后我会通过过滤sep列来拆分DF,我的疑问是:如何使用窗口函数生成上述DF中的sep列?

解决方案

可以用row_number()窗口函数结合整数除法来生成均匀分配的sep列,既不需要缓存全量DF,也能保证拆分后行数总和与原DF完全一致。

代码实现(PySpark)

from pyspark.sql import Window
from pyspark.sql.functions import row_number, ceil, col, rand

# 预先计算总行数(避免重复触发计算)
total_rows = df.count()

# 定义窗口:若不需要保留原顺序,用随机排序避免数据倾斜
window_spec = Window.orderBy(rand(seed=42))  # 需保留顺序则改为orderBy("val1")

# 添加行号并生成sep列
df_with_sep = df.withColumn("row_num", row_number().over(window_spec)) \
                .withColumn("sep", ceil((col("row_num") * 8) / total_rows)) \
                .drop("row_num")

关键说明

  1. row_number():为每一行生成唯一连续的行号,确保每行对应唯一标识,不会出现重复或遗漏的情况。
  2. ceil((row_num * 8)/total_rows):通过行号乘以拆分份数(8)再除以总行数,向上取整后得到均匀分配的分组标识sep,能让每个分组的行数尽可能接近,不会出现总和不符的问题。
  3. 避免缓存:整个流程无需缓存全量DF,仅count()会触发一次计算,后续生成sep列的操作是流式处理,不会占用大量磁盘空间。
  4. 随机排序优化:使用rand(seed=42)排序可以避免数据倾斜,同时设置种子保证结果可复现;若需要保留原DF的顺序,替换为具体字段排序即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 04:09:22