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

如何在PySpark DataFrame中基于Account和check列生成result新列?

PySpark实现按Account分组并按N分隔的Y值计数需求

需求说明

现有PySpark DataFrame结构如下:

Accountcheck
10Y
10Y
10N
10Y
12Y
12Y
12N

需要生成result列,最终规则:

  • 同一Account分组下,连续的Y按N分隔成不同组,每组Y对应相同的序号(从1开始,每组递增)
  • check为N时,result标记为na
  • 每个新Account重新按规则初始化计数

最终结果示例:

Accountcheckresult
10Y1
10Y1
10Nna
10Y2
12Y1
12Y1
12Nna

实现代码

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()

关键步骤解释

  1. 添加行号:Spark DataFrame是分布式无序的,必须指定排序字段保证后续窗口函数的计算顺序,这里用monotonically_increasing_id()生成全局唯一递增ID,确保原始顺序被保留。
  2. 累计N数量:通过窗口函数计算每个Account下从第一行到当前行的N的累计数,每遇到一个N,后续行的累计数就会加1,从而将连续的Y划分到不同的分组中。
  3. 生成结果列:对check为N的行直接标记na,Y的行用累计N数加1作为分组序号,第一个Y组累计N数为0,对应序号1;N之后的Y组累计N数为1,对应序号2,完全符合需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 21:35:35