如何在PySpark中按每组N条数据进行分区分组
PySpark实现按分组最多5条数据分组(类似Rails的in_groups_of方法)
需求说明
现有PySpark DataFrame,需按text、title、target_url、display_domain列分组后,每个分组内最多5条数据划分为一个子组,生成全局唯一递增的group_id,最终得到8个数据分组,行为类似Rails的in_groups_of方法。
原始数据与初始化代码
from pyspark.sql.types import StructType, StructField, StringType, IntegerType from pyspark.sql.window import Window import pyspark.sql.functions as F data = [ ( 1, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 2, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 3, "AAA", "BBB", "CCC", "DDD", "mobile"), ( 4, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 5, "AAA", "BBB", "CCC", "DDD", "mobile"), ( 6, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 7, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 8, "AAA", "BBB", "CCC", "DDD", "desktop"), ( 9, "AAA", "BBB", "CCC", "DDD", "desktop"), (10, "AAA", "BBB", "CCC", "DDD", "mobile"), (11, "AAA", "BBB", "CCC", "DDD", "desktop"), (12, "EEE", "FFF", "GGG", "HHH", "desktop"), (13, "EEE", "FFF", "GGG", "HHH", "mobile"), (14, "EEE", "FFF", "GGG", "HHH", "desktop"), (15, "EEE", "FFF", "GGG", "HHH", "mobile"), (16, "EEE", "FFF", "GGG", "HHH", "desktop"), (17, "EEE", "FFF", "GGG", "HHH", "desktop"), (18, "EEE", "FFF", "GGG", "HHH", "desktop"), (19, "III", "JJJ", "KKK", "LLL", "desktop"), (20, "III", "JJJ", "KKK", "LLL", "mobile"), (21, "III", "JJJ", "KKK", "LLL", "desktop"), (22, "III", "JJJ", "KKK", "LLL", "desktop"), (23, "III", "JJJ", "KKK", "LLL", "mobile"), (24, "III", "JJJ", "KKK", "LLL", "desktop"), (25, "III", "JJJ", "KKK", "LLL", "desktop"), (26, "III", "JJJ", "KKK", "LLL", "desktop"), (27, "III", "JJJ", "KKK", "LLL", "desktop"), (28, "III", "JJJ", "KKK", "LLL", "desktop"), (29, "III", "JJJ", "KKK", "LLL", "desktop"), (30, "III", "JJJ", "KKK", "LLL", "mobile") ] schema = StructType([ \ StructField("id", IntegerType(),True), StructField("text", StringType(),True), StructField("title", StringType(),True), StructField("target_url", StringType(), True), StructField("display_domain", StringType(), True), StructField("device", StringType(), True) ]) df = spark.createDataFrame(data=data,schema=schema) columns = [ "text", "title", "target_url", "display_domain" ] windowSpecByPartition = ( Window.partitionBy( columns ).orderBy("id") ) overall_row_number_df = df.withColumn( "overall_row_number", F.row_number().over(windowSpecByPartition) )
解决方案
分两步实现:
- 计算分区内子组号:基于已生成的
overall_row_number,通过整数除法计算每个分区内的子组标识partition_group_id。 - 生成全局唯一group_id:将原始分区列与
partition_group_id组合,使用全局排序的dense_rank()生成全局递增的分组ID。
完整实现代码:
# 步骤1:计算每个分区内的子组号 partition_group_df = overall_row_number_df.withColumn( "partition_group_id", F.floor((F.col("overall_row_number") - 1) / 5) + 1 ) # 步骤2:生成全局唯一的group_id windowSpecGlobal = Window.orderBy(*columns, "partition_group_id") final_df = partition_group_df.withColumn( "group_id", F.dense_rank().over(windowSpecGlobal) ).drop("overall_row_number", "partition_group_id") # 查看结果 final_df.orderBy("id").show(30, truncate=False)
预期结果
| id | text | title | target_url | display_domain | device | group_id |
|---|---|---|---|---|---|---|
| 1 | AAA | BBB | CCC | DDD | desktop | 1 |
| 2 | AAA | BBB | CCC | DDD | desktop | 1 |
| 3 | AAA | BBB | CCC | DDD | mobile | 1 |
| 4 | AAA | BBB | CCC | DDD | desktop | 1 |
| 5 | AAA | BBB | CCC | DDD | mobile | 1 |
| 6 | AAA | BBB | CCC | DDD | desktop | 2 |
| 7 | AAA | BBB | CCC | DDD | desktop | 2 |
| 8 | AAA | BBB | CCC | DDD | desktop | 2 |
| 9 | AAA | BBB | CCC | DDD | desktop | 2 |
| 10 | AAA | BBB | CCC | DDD | mobile | 2 |
| 11 | AAA | BBB | CCC | DDD | desktop | 3 |
| 12 | EEE | FFF | GGG | HHH | desktop | 4 |
| 13 | EEE | FFF | GGG | HHH | mobile | 4 |
| 14 | EEE | FFF | GGG | HHH | desktop | 4 |
| 15 | EEE | FFF | GGG | HHH | mobile | 4 |
| 16 | EEE | FFF | GGG | HHH | desktop | 4 |
| 17 | EEE | FFF | GGG | HHH | desktop | 5 |
| 18 | EEE | FFF | GGG | HHH | desktop | 5 |
| 19 | III | JJJ | KKK | LLL | desktop | 6 |
| 20 | III | JJJ | KKK | LLL | mobile | 6 |
| 21 | III | JJJ | KKK | LLL | desktop | 6 |
| 22 | III | JJJ | KKK | LLL | desktop | 6 |
| 23 | III | JJJ | KKK | LLL | mobile | 6 |
| 24 | III | JJJ | KKK | LLL | desktop | 7 |
| 25 | III | JJJ | KKK | LLL | desktop | 7 |
| 26 | III | JJJ | KKK | LLL | desktop | 7 |
| 27 | III | JJJ | KKK | LLL | desktop | 7 |
| 28 | III | JJJ | KKK | LLL | desktop | 7 |
| 29 | III | JJJ | KKK | LLL | desktop | 8 |
| 30 | III | JJJ | KKK | LLL | mobile | 8 |
内容的提问来源于stack exchange,提问作者Sergio Flores
相关产品推荐
相关产品推荐

