如何在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()
代码说明
- 窗口定义:用
Window.partitionBy指定分组键,确保后续计算仅在同一VehNum+Control_circuit的组内进行。 - 状态判断:通过
count(when(...))统计组内"ok"状态的记录数,若数量大于0则判定组内存在"ok"。 - 生成目标列:利用
when函数根据判断结果赋值0或1,最后清理中间辅助列has_ok。
内容的提问来源于stack exchange,提问作者karthik
相关产品推荐
相关产品推荐

