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

