如何在PySpark DataFrame中按条件生成Control_Flag列
PySpark实现新增Control_Flag列方案
可以通过窗口函数快速实现这个需求,核心思路是按VehNum和Control_circuit分组后,判断组内是否存在非0的Flag值,再给每一行打上对应的Control_Flag标记。
实现代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按VehNum和Control_circuit分组 window_spec = Window.partitionBy("VehNum", "Control_circuit") # 新增Control_Flag列 df_result = df.withColumn( "Control_Flag", F.when(F.max("Flag").over(window_spec) == 0, 0).otherwise(1) ) # 查看结果 df_result.show()
代码说明
- 窗口
window_spec指定了分组依据:VehNum和Control_circuit,确保只在同一车辆同一控制回路内做判断。 F.max("Flag").over(window_spec)计算每个分组内Flag的最大值:- 如果最大值是0,说明组内所有
Flag都是0,此时Control_Flag设为0; - 只要最大值大于0(即存在1或2),
Control_Flag就设为1。
- 如果最大值是0,说明组内所有
- 这种方式无需额外join操作,直接通过窗口函数给每一行填充结果,效率更高。
验证结果
运行代码后得到的结果完全符合期望:
+-------+---------------+-----------------------+------------+---------------+----+------------+ |VehNum |Control_circuit|control_circuit_status|partnumbers|errors |Flag|Control_Flag| +-------+---------------+-----------------------+------------+---------------+----+------------+ |4234456|DOC |ok |A567UR |Software Issue |0 |1 | |4234456|DOC |not_okay |A568UR |Software Issue |1 |1 | |4234456|DOC |not_okay |A569UR |Hardware issue |2 |1 | |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 | +-------+---------------+-----------------------+------------+---------------+----+------------+
内容的提问来源于stack exchange,提问作者karthik
相关产品推荐
相关产品推荐

