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
相关产品推荐
相关产品推荐

