【问题标题】:Apply comparison operations defined in dataset应用数据集中定义的比较操作
【发布时间】:2022-10-25 06:16:36
【问题描述】:

我有一个包含几个字段的表,我需要在这些字段上进行数据质量检查。

数据质量检查被定义为第二个表中的规则。

数据表:

ID Name1 Name2 Zip1 Zip2
001 John John 123 123
002 Sara Sarah 234 234
003 Bill William 999 111
004 Lisa Lisa 888 333
005 Martin Martin 345 345
006 Margaret Margaret 456 456
007 Oscar Oscar 678 678
008 Peter Peter 789 789

规则表:

ID FieldLeft FieldRight ComparisonOperation
R001 Name1 Name2 EQUALS
R002 Zip1 Zip2 EQUALS

所以规则本质上是说:Name1=Name2 and Zip1=Zip2

预期的输出是不符合规则的记录。 它应该为每个规则违规产生一行(参见记录 003,名称和 zip 都不一致 -> 因此记录 003 的结果中有两行)。

输出:

Rule ID FieldLeft FieldRight
R001 002 Sara Sarah
R001 003 Bill William
R002 003 999 111
R002 004 888 333

【问题讨论】:

    标签: pyspark


    【解决方案1】:

    这是我的实现

    from pyspark.sql import functions as F
    from pyspark.sql.types import *
    
    df = spark.createDataFrame(
        [
            ("001", "John", "John", "123", "123"),
            ("002", "Sara", "Sarah", "234", "234"),
            ("003", "Bill", "William", "999", "111"),
            ("004", "Lisa", "Lisa", "888", "333"),
            ("005", "Martin", "Martin", "345", "345"),
            ("006", "Margaret", "Margaret", "456", "456"),
            ("007", "Oscar", "Oscar", "678", "678"),
            ("008", "Peter", "Peter", "789", "789"),
        ],
        ["ID", "Name1", "Name2", "Zip1", "Zip2"],
    )
    #df.show()
    
    rule_df = spark.createDataFrame(
        [
            ("R001", "Name1", "Name2", "EQUALS"),
            ("R002", "Zip1", "Zip2", "EQUALS"),
        ],
        ["ID", "FieldLeft", "FieldRight", "ComparisonOperation"],
    )
    #rule_df.show()
    
    final_rule_df = (rule_df
        .withColumn(
            "operator",
            F.when(
                F.lower(F.col("ComparisonOperation")) == "equals",
                F.lit(" == "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "not equals",
                F.lit(" != "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "greater than",
                F.lit(" > "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "less than",
                F.lit(" < "),
            )
            .otherwise(F.lit("operator_na")),
        )
        .filter(F.col("operator") != "operator_na" )
        .withColumn("expression", concat(F.col("FieldLeft"),F.col("operator"), F.col("FieldRight"))  )
        .drop("operator")
        #.withColumn(
        #    "select_clause", 
        #    F.concat(
        #        F.lit('"'),
        #        F.lit( F.col("FieldLeft") ),
        #        F.lit(" as " + F.col("FieldLeft")._jc.toString()),
        #        F.lit('"'),
        #        F.lit(", "),
        #        F.lit('"'),
        #        F.col("FieldRight"),
        #        F.lit(" as " + F.col("FieldRight")._jc.toString()),
        #        F.lit('"'),
        #    )
        #)                      
    )
    final_rule_df.show(truncate=False)
    
    schema = StructType(
        [
            StructField("Rule", StringType(), True),
            StructField("ID", StringType(), True),
            StructField("FieldLeft", StringType(), True),
            StructField("FieldRight", StringType(), True),
        ]
    )
    
    final_non_compliant_df = spark.createDataFrame(
        spark.sparkContext.emptyRDD(), schema
    )
    
    rule_df_rows = final_rule_df.select("*").collect()
    for row in rule_df_rows:
        rule_id = row.ID
        print(f"rule_id: {rule_id}")
        
        expression = row.expression
        print(f"expression: {expression}")
        
        #select_clause = row.select_clause
        #print(f"select_clause: {select_clause}")
        
        rule_df = df.filter(expr(expression))
        #rule_df.show()
        
        non_compliant_df = (df.subtract(rule_df)
            .withColumn("Rule", F.lit(rule_id))
            .withColumn("FieldLeft", F.col(row.FieldLeft))
            .withColumn("FieldRight", F.col(row.FieldRight))
            .selectExpr("Rule", "ID", "FieldLeft", "FieldRight")
        )
        non_compliant_df.show()
        final_non_compliant_df = final_non_compliant_df.union(non_compliant_df)
    
    final_non_compliant_df.show()
    

    输出:

    +----+---------+----------+-------------------+--------------+
    |ID  |FieldLeft|FieldRight|ComparisonOperation|expression    |
    +----+---------+----------+-------------------+--------------+
    |R001|Name1    |Name2     |EQUALS             |Name1 == Name2|
    |R002|Zip1     |Zip2      |EQUALS             |Zip1 == Zip2  |
    +----+---------+----------+-------------------+--------------+
    
    rule_id: R001
    expression: Name1 == Name2
    +----+---+---------+----------+
    |Rule| ID|FieldLeft|FieldRight|
    +----+---+---------+----------+
    |R001|003|     Bill|   William|
    |R001|002|     Sara|     Sarah|
    +----+---+---------+----------+
    
    rule_id: R002
    expression: Zip1 == Zip2
    +----+---+---------+----------+
    |Rule| ID|FieldLeft|FieldRight|
    +----+---+---------+----------+
    |R002|004|      888|       333|
    |R002|003|      999|       111|
    +----+---+---------+----------+
    

    最终输出:

    +----+---+---------+----------+
    |Rule| ID|FieldLeft|FieldRight|
    +----+---+---------+----------+
    |R001|003|     Bill|   William|
    |R001|002|     Sara|     Sarah|
    |R002|004|      888|       333|
    |R002|003|      999|       111|
    +----+---+---------+----------+
    

    【讨论】:

    • 太感谢了!一个问题:有没有一种方法可以并行化遍历规则的 for 循环?在我的场景中,我正在查看数百万行数据集中的数百条规则。所以我想知道是否可以优化。
    【解决方案2】:

    @hbit 我不确定在没有显式循环的情况下执行此操作的完整解决方案。我尽可能使用交叉连接为创建笛卡尔结果集的每条记录添加规则。我无法弄清楚如何让表达式列评估为布尔值

    from pyspark.sql import functions as F
    from pyspark.sql.types import *
    
    df = spark.createDataFrame(
        [
            ("001", "John", "John", "123", "123"),
            ("002", "Sara", "Sarah", "234", "234"),
            ("003", "Bill", "William", "999", "111"),
            ("004", "Lisa", "Lisa", "888", "333"),
            ("005", "Martin", "Martin", "345", "345"),
            ("006", "Margaret", "Margaret", "456", "456"),
            ("007", "Oscar", "Oscar", "678", "678"),
            ("008", "Peter", "Peter", "789", "789"),
        ],
        ["ID", "Name1", "Name2", "Zip1", "Zip2"],
    )
    
    
    rule_df = spark.createDataFrame(
        [
            ("R001", "Name1", "Name2", "EQUALS"),
            ("R002", "Zip1", "Zip2", "EQUALS"),
        ],
        ["ID", "FieldLeft", "FieldRight", "ComparisonOperation"],
    )
    #rule_df.show()
    
    final_rule_df = (rule_df
        .withColumn(
            "operator",
            F.when(
                F.lower(F.col("ComparisonOperation")) == "equals",
                F.lit(" == "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "not equals",
                F.lit(" != "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "greater than",
                F.lit(" > "),
            )
            .when(
                F.lower(F.col("ComparisonOperation")) == "less than",
                F.lit(" < "),
            )
            .otherwise(F.lit("operator_na")),
        )
        .filter(F.col("operator") != "operator_na" )
        .withColumn("expression", concat(F.lit("("), F.col("FieldLeft"),F.col("operator"), F.col("FieldRight"), F.lit(")"))  )
        .drop("operator")        
    )
    
    final_df = (
        df.crossJoin(final_rule_df) 
    )
    
    final_df.show()
    
    +----+---------+----------+-------------------+----------------+
    |ID  |FieldLeft|FieldRight|ComparisonOperation|expression      |
    +----+---------+----------+-------------------+----------------+
    |R001|Name1    |Name2     |EQUALS             |(Name1 == Name2)|
    |R002|Zip1     |Zip2      |EQUALS             |(Zip1 == Zip2)  |
    +----+---------+----------+-------------------+----------------+
    
    +---+--------+--------+----+----+----+---------+----------+-------------------+----------------+
    | ID|   Name1|   Name2|Zip1|Zip2|  ID|FieldLeft|FieldRight|ComparisonOperation|      expression|
    +---+--------+--------+----+----+----+---------+----------+-------------------+----------------+
    |001|    John|    John| 123| 123|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |001|    John|    John| 123| 123|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |002|    Sara|   Sarah| 234| 234|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |002|    Sara|   Sarah| 234| 234|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |003|    Bill| William| 999| 111|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |003|    Bill| William| 999| 111|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |004|    Lisa|    Lisa| 888| 333|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |004|    Lisa|    Lisa| 888| 333|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |005|  Martin|  Martin| 345| 345|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |005|  Martin|  Martin| 345| 345|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |006|Margaret|Margaret| 456| 456|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |006|Margaret|Margaret| 456| 456|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |007|   Oscar|   Oscar| 678| 678|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |007|   Oscar|   Oscar| 678| 678|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    |008|   Peter|   Peter| 789| 789|R001|    Name1|     Name2|             EQUALS|(Name1 == Name2)|
    |008|   Peter|   Peter| 789| 789|R002|     Zip1|      Zip2|             EQUALS|  (Zip1 == Zip2)|
    +---+--------+--------+----+----+----+---------+----------+-------------------+----------------+
    

    【讨论】:

      猜你喜欢
      • 2020-09-28
      • 1970-01-01
      • 1970-01-01
      • 2023-02-06
      • 1970-01-01
      • 2012-02-10
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多