如何在PySpark DataFrame中基于Account和check列生成result新列?
PySpark实现按Account分组并按N分隔的Y值计数需求
需求说明
现有PySpark DataFrame结构如下:
| Account | check |
|---|---|
| 10 | Y |
| 10 | Y |
| 10 | N |
| 10 | Y |
| 12 | Y |
| 12 | Y |
| 12 | N |
需要生成result列,最终规则:
- 同一
Account分组下,连续的Y按N分隔成不同组,每组Y对应相同的序号(从1开始,每组递增) check为N时,result标记为na- 每个新
Account重新按规则初始化计数
最终结果示例:
| Account | check | result |
|---|---|---|
| 10 | Y | 1 |
| 10 | Y | 1 |
| 10 | N | na |
| 10 | Y | 2 |
| 12 | Y | 1 |
| 12 | Y | 1 |
| 12 | N | na |
实现代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("YGroupCount").getOrCreate() # 创建测试DataFrame(实际使用时替换为你的数据源) data = [ (10, "Y"), (10, "Y"), (10, "N"), (10, "Y"), (12, "Y"), (12, "Y"), (12, "N") ] df = spark.createDataFrame(data, ["Account", "check"]) # 1. 按Account分组,添加行号保证处理顺序(Spark DataFrame默认无序,需明确排序依据) # 若有业务有序字段(如时间戳),可替换monotonically_increasing_id()为该字段 window_account = Window.partitionBy("Account").orderBy(F.monotonically_increasing_id()) df_with_row = df.withColumn("row_num", F.row_number().over(window_account)) # 2. 计算每个Account下,截至当前行的累计N数量,以此划分Y的分组 window_n_cum = Window.partitionBy("Account").orderBy("row_num").rowsBetween(Window.unboundedPreceding, 0) df_with_group = df_with_row.withColumn( "cum_n_count", F.sum(F.when(F.col("check") == "N", 1).otherwise(0)).over(window_n_cum) ) # 3. 生成result列:N对应na,Y对应累计N数量+1(即分组序号) result_df = df_with_group.withColumn( "result", F.when(F.col("check") == "N", F.lit("na")).otherwise(F.col("cum_n_count") + 1) ).drop("row_num", "cum_n_count") # 查看最终结果 result_df.show()
关键步骤解释
- 添加行号:Spark DataFrame是分布式无序的,必须指定排序字段保证后续窗口函数的计算顺序,这里用
monotonically_increasing_id()生成全局唯一递增ID,确保原始顺序被保留。 - 累计N数量:通过窗口函数计算每个
Account下从第一行到当前行的N的累计数,每遇到一个N,后续行的累计数就会加1,从而将连续的Y划分到不同的分组中。 - 生成结果列:对
check为N的行直接标记na,Y的行用累计N数加1作为分组序号,第一个Y组累计N数为0,对应序号1;N之后的Y组累计N数为1,对应序号2,完全符合需求。
内容的提问来源于stack exchange,提问作者baqm
相关产品推荐
相关产品推荐

