【问题标题】:Spark Dataframe Nested Case When StatementSpark Dataframe 嵌套 Case When 语句
【发布时间】:2018-03-20 08:00:04
【问题描述】:

我需要在 Spark DataFrame 中实现以下 SQL 逻辑

SELECT KEY,
    CASE WHEN tc in ('a','b') THEN 'Y'
         WHEN tc in ('a') AND amt > 0 THEN 'N'
         ELSE NULL END REASON,
FROM dataset1;

我的输入DataFrame如下:

val dataset1 = Seq((66, "a", "4"), (67, "a", "0"), (70, "b", "4"), (71, "d", "4")).toDF("KEY", "tc", "amt")

dataset1.show()
+---+---+---+
|KEY| tc|amt|
+---+---+---+
| 66|  a|  4|
| 67|  a|  0|
| 70|  b|  4|
| 71|  d|  4|
+---+---+---+

我已经实现了嵌套 case when 语句为:

dataset1.withColumn("REASON", when(col("tc").isin("a", "b"), "Y")
  .otherwise(when(col("tc").equalTo("a") && col("amt").geq(0), "N")
    .otherwise(null))).show()
+---+---+---+------+
|KEY| tc|amt|REASON|
+---+---+---+------+
| 66|  a|  4|     Y|
| 67|  a|  0|     Y|
| 70|  b|  4|     Y|
| 71|  d|  4|  null|
+---+---+---+------+

如果嵌套的when语句走得更远,上述逻辑与“otherwise”语句的可读性会有点混乱。

在 Spark DataFrames 中,有没有更好的方法来实现嵌套 case when 语句?

【问题讨论】:

    标签: sql apache-spark dataframe apache-spark-sql


    【解决方案1】:

    这里没有嵌套,因此不需要otherwise。您需要的只是链接when

    import spark.implicits._
    
    when($"tc" isin ("a", "b"), "Y")
      .when($"tc" === "a" && $"amt" >= 0, "N")
    

    ELSE NULL 是隐含的,因此您可以完全省略它。

    您使用的模式更适用于folding,而不是数据结构:

    val cases = Seq(
      ($"tc" isin ("a", "b"), "Y"),
      ($"tc" === "a" && $"amt" >= 0, "N")
    )
    

    其中when - otherwise 自然遵循递归模式,null 提供基本情况。

    cases.foldLeft(lit(null)) {
      case (acc, (expr, value)) => when(expr, value).otherwise(acc)
    }
    

    请注意,在这一系列条件下,不可能达到“N”结果。如果tc 等于“a”,它将被第一个子句捕获。如果不是,它将无法同时满足谓词并默认为NULL。你应该:

    when($"tc" === "a" && $"amt" >= 0, "N")
     .when($"tc" isin ("a", "b"), "Y")
    

    【讨论】:

      【解决方案2】:

      对于更复杂的逻辑,我更喜欢使用 UDF 以获得更好的可读性:

      val selectCase = udf((tc: String, amt: String) =>
        if (Seq("a", "b").contains(tc)) "Y"
        else if (tc == "a" && amt.toInt <= 0) "N"
        else null
      )
      
      
      dataset1.withColumn("REASON", selectCase(col("tc"), col("amt")))
        .show
      

      【讨论】:

      • 但应该提到 udf 可能会降低性能,因为它们可能会阻止过滤器下推。当然,情况并非总是如此,但尽可能坚持使用 spark 的本机功能是一个好习惯。
      【解决方案3】:

      你可以简单地在你的数据集上使用 selectExpr

      dataset1.selectExpr("*", "CASE WHEN tc in ('a') AND amt > 0 THEN 'N' WHEN tc in ('a','b') THEN 'Y' ELSE NULL END
      REASON").show()
      
      +---+---+---+------+
      |KEY| tc|amt|REASON|
      +---+---+---+------+
      | 66|  a|  4|     N|
      | 67|  a|  0|     Y|
      | 70|  b|  4|     Y|
      | 71|  d|  4|  null|
      +---+---+---+------+
      

      第二个条件应该放在第一个之前,因为第一个条件更通用。

      当 tc in ('a') AND amt > 0 THEN 'N'

      【讨论】:

        猜你喜欢
        • 2016-12-10
        • 2014-07-12
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2011-03-12
        相关资源
        最近更新 更多