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

如何在PySpark DataFrame中按OK/非OK状态生成Control_Flag列

PySpark实现新增Control_Flag列的需求

原始数据

现有PySpark DataFrame df 数据如下:

VehNum  Control_circuit control_circuit_status  partnumbers     errors     Flag
4234456 DOC             ok                      A567UR      Software Issue  0
4234456 DOC             not_okay                A568UR      Software Issue  1
4234456 DOC             not_okay                A569UR      Hardware issue  2
4234457 ACR             ok                      A234TY      Hardware issue  0
4234457 ACR             ok                      A235TY      Hardware issue  0
4234457 ACR             ok                      A234TY      Hardware issue  0
4234487 QWR             ok                      A276TY      Hardware issue  0
4234487 QWR             not_okay                A872UR      Hardware issue  1
3423448 QWR             not_okay                A872UR      Hardware issue  1

需求说明

需新增Control_Flag列,规则为:对每个VehNum和Control_circuit分组,若组内control_circuit_status存在"ok"状态,则Control_Flag取值为0,否则为1。

期望结果

VehNum  Control_circuit control_circuit_status  partnumbers     errors     Flag Control_Flag
4234456 DOC             ok                      A567UR      Software Issue  0   0
4234456 DOC             not_okay                A568UR      Software Issue  1   0
4234456 DOC             not_okay                A569UR      Hardware issue  2   0
4234457 ACR             ok                      A234TY      Hardware issue  0   0
4234457 ACR             ok                      A235TY      Hardware issue  0   0
4234457 ACR             ok                      A234TY      Hardware issue  0   0
4234487 QWR             ok                      A276TY      Hardware issue  0   1
4234487 QWR             not_okay                A872UR      Hardware issue  1   1
3423448 QWR             not_okay                A872UR      Hardware issue  1   1

实现代码

通过窗口函数实现分组内的状态判断,将结果同步到组内每一行:

from pyspark.sql import Window
from pyspark.sql.functions import when, count, col

# 定义窗口:按VehNum和Control_circuit分组
window_spec = Window.partitionBy("VehNum", "Control_circuit")

# 生成Control_Flag列
df_result = df.withColumn(
    # 统计分组内"ok"状态的记录数,判断是否存在
    "has_ok",
    count(when(col("control_circuit_status") == "ok", True)).over(window_spec) > 0
).withColumn(
    # 根据存在性设置Control_Flag值
    "Control_Flag",
    when(col("has_ok"), 0).otherwise(1)
).drop("has_ok")  # 删除中间辅助列

# 查看最终结果
df_result.show()

代码说明

  1. 窗口定义:用Window.partitionBy指定分组键,确保后续计算仅在同一VehNum+Control_circuit的组内进行。
  2. 状态判断:通过count(when(...))统计组内"ok"状态的记录数,若数量大于0则判定组内存在"ok"。
  3. 生成目标列:利用when函数根据判断结果赋值0或1,最后清理中间辅助列has_ok。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 09:15:59