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

Spark DataFrame按双值排序并按指定大小划分分类桶

程序化划分Spark DataFrame为等大小分类桶

现有包含val和id列的Spark DataFrame,id值唯一、val值可重复。需先按val、id排序生成行号rn,再将数据划分为4个尽可能等大的分类桶。手动划分易产生null值且不够灵活,以下是两种程序化实现方案:

原始数据与行号生成

先准备数据并生成排序后的行号:

from pyspark.sql import SparkSession
from pyspark.sql.functions import row_number
from pyspark.sql.window import Window

spark = SparkSession.builder.appName("BucketSplit").getOrCreate()

data = [ 
  (0, 'A'), 
  (0, 'B'), 
  (0, 'C'), 
  (1, 'D'), 
  (1, 'E'), 
  (2, 'F'), 
  (3, 'G'), 
  (4, 'H'), 
  (5, 'I'),
] 
schm = ['val', 'id'] 
df = spark.createDataFrame(data, schema=schm)

# 按val、id排序生成行号rn
window_spec = Window().orderBy('val', 'id')
df = df.withColumn("rn", row_number().over(window_spec))

方案1:基于行号的计算式划分

通过行号与平均桶大小的比值向上取整,自动处理无法整除的情况,无null值问题:

from pyspark.sql.functions import ceil, col

n_splits = 4
total_rows = df.count()
bucket_size = total_rows / n_splits  # 用浮点数保证除法精度

df = df.withColumn("cat", ceil(col("rn") / bucket_size).cast("integer"))
df.show()

方案2:使用Spark内置ntile函数

Spark原生ntile函数可直接将有序数据划分为指定数量的桶,自动均分,代码更简洁:

from pyspark.sql.functions import ntile

n_splits = 4
# ntile依赖有序窗口,复用之前的排序窗口即可
df = df.withColumn("cat", ntile(n_splits).over(window_spec))
df.show()

最终输出

两种方案均会得到符合预期的结果:

+---+---+---+---+
|val| id| rn|cat|
+---+---+---+---+
|  0|  A|  1|  1|
|  0|  B|  2|  1|
|  0|  C|  3|  2|
|  1|  D|  4|  2|
|  1|  E|  5|  3|
|  2|  F|  6|  3|
|  3|  G|  7|  4|
|  4|  H|  8|  4|
|  5|  I|  9|  4|
+---+---+---+---+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 10:26:04