PySpark中仅为过滤后的数据行添加行号
问题描述
我有如下DataFrame:
| id_cnt | id_prd | type | price |
|---|---|---|---|
| 1 | A | SS | 10 |
| 2 | A | AA | 20 |
| 3 | A | AA | 25 |
| 1 | B | AA | 55 |
| 2 | B | SS | 50 |
| 3 | B | AA | 75 |
| 4 | B | AA | 80 |
需要添加新列rownumber:按id_prd分组,仅对type = "AA"的行按price降序分配行号,其余行该列值为null。预期输出如下:
| id_cnt | id_prd | type | price | rownumber |
|---|---|---|---|---|
| 1 | A | SS | 10 | null |
| 2 | A | AA | 20 | 2 |
| 3 | A | AA | 25 | 1 |
| 1 | B | AA | 55 | 3 |
| 2 | B | SS | 50 | null |
| 3 | B | AA | 75 | 2 |
| 4 | B | AA | 80 | 1 |
Pandas 实现方案
结合分组、排名和条件筛选,实现需求:
import pandas as pd # 构造原始DataFrame df = pd.DataFrame({ 'id_cnt': [1,2,3,1,2,3,4], 'id_prd': ['A','A','A','B','B','B','B'], 'type': ['SS','AA','AA','AA','SS','AA','AA'], 'price': [10,20,25,55,50,75,80] }) # 生成rownumber列 df['rownumber'] = df.groupby('id_prd').apply( lambda x: x['price'].rank(ascending=False, method='first').where(x['type'] == 'AA') ).reset_index(level=0, drop=True) # 转换为可空整数类型,使空值显示为null df['rownumber'] = df['rownumber'].astype('Int64') print(df)
关键逻辑说明
groupby('id_prd'):按产品ID分组处理rank(ascending=False, method='first'):对价格降序排名,method='first'避免同价行出现重复排名where(x['type'] == 'AA'):仅保留type为AA的行的排名结果,其余设为NaNastype('Int64'):转换为Pandas可空整数类型,让空值显示为null,匹配预期输出格式
PySpark 实现方案
使用窗口函数结合条件判断实现:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import row_number, when # 初始化Spark会话 spark = SparkSession.builder.appName("rownumber_example").getOrCreate() # 构造原始DataFrame data = [ (1, 'A', 'SS', 10), (2, 'A', 'AA', 20), (3, 'A', 'AA', 25), (1, 'B', 'AA', 55), (2, 'B', 'SS', 50), (3, 'B', 'AA', 75), (4, 'B', 'AA', 80) ] df = spark.createDataFrame(data, ['id_cnt', 'id_prd', 'type', 'price']) # 定义窗口规则:按id_prd分组,price降序排序 window_spec = Window.partitionBy('id_prd').orderBy(df['price'].desc()) # 添加rownumber列 df = df.withColumn( 'rownumber', when(df['type'] == 'AA', row_number().over(window_spec)).otherwise(None) ) df.show()
关键逻辑说明
Window.partitionBy('id_prd').orderBy(df['price'].desc()):定义分组排序的窗口规则when(df['type'] == 'AA', row_number().over(window_spec)).otherwise(None):仅对type为AA的行生成行号,其余行设为null
内容的提问来源于stack exchange,提问作者tchita
相关产品推荐
相关产品推荐

