【问题标题】:how to segregate the column wrt certain conditions in pyspark dataframe如何在 pyspark 数据框中根据某些条件分离列
【发布时间】:2022-11-22 02:45:21
【问题描述】:

我有一个数据框 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 如果它的标志值仅为 0,则 Control_Flag 列将保持值 0,否则如果它具有 0、1 或 2,则 Control_Flag 列将保持值1.

结果应如下所示:

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

如何使用pyspark实现这一目标?

【问题讨论】:

    标签: python python-3.x pyspark


    【解决方案1】:

    使用带有 SUM() 的聚合窗口将有助于实现这一点

    from pyspark.sql import functions as F
    from pyspark.sql.types import *
    from pyspark.sql import Window
    
    df = spark.createDataFrame(
        [
            ("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", "A234TY", "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),
        ],
        ["VehNum", "Control_circuit", "control_circuit_status", "partnumbers", "errors", "Flag"],
    )
    
    df_agg_window = Window.partitionBy(
        "VehNum",
        "Control_circuit",
    )
    
    df = (
        df
        .withColumn(
            "flag_sum",
            F.sum("Flag").over(df_agg_window),
        )
        .withColumn(
            "Control_Flag",
            F.when(
                F.lower(F.col("flag_sum")) > 0,
                F.lit(1),
            )
            .otherwise(F.lit(0)),
        )
        #.drop(F.col("flag_sum"))
    )
    
    
    df.show()
    

    输出:

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

    【讨论】:

      猜你喜欢
      • 2022-09-23
      • 2020-10-01
      • 1970-01-01
      • 1970-01-01
      • 2021-11-28
      • 2012-10-16
      • 2016-11-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多