【问题标题】:Conditional Spark map() function based on input columns基于输入列的条件 Spark map() 函数
【发布时间】:2020-03-16 00:40:41
【问题描述】:

我在这里尝试实现的是向 Spark SQL map 函数发送有条件生成的列,具体取决于它们是否具有 null0 或我可能想要的任何其他值。

以这个初始 DF 为例。

val initialDF = Seq(
  ("a", "b", 1), 
  ("a", "b", null), 
  ("a", null, 0)
).toDF("field1", "field2", "field3")

我想从最初的 DataFrame 生成另一列,这将是一个地图,就像这样。

initialDF.withColumn("thisMap", MY_FUNCTION)

我目前的处理方法基本上是在方法中使用Seq[String] flatMap Spark SQL 方法接收的键值对,就像这样。

def toMap(columns: String*): Column = {
  map(
    columns.flatMap(column => List(lit(column), col(column))): _*
  )
}

但是,过滤变成了 Scala 的事情,而且是一团糟。

处理后我想获得的是,对于这些行中的每一行,下一个 DataFrame。

val initialDF = Seq(
  ("a", "b", 1, Map("field1" -> "a", "field2" -> "b", "field3" -> 1)),
  ("a", "b", null, Map("field1" -> "a", "field2" -> "b")),
  ("a", null, 0, Map("field1" -> "a"))
)
  .toDF("field1", "field2", "field3", "thisMap")

我想知道这是否可以使用Column API 来实现,而.isNull.equalTo 更直观?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    这是对上面 Lamanus 的回答的一个小改进,它只在 df.columns 上循环一次:

    import org.apache.spark.sql._
    import org.apache.spark.sql.functions._
    
    case class Record(field1: String, field2: String, field3: java.lang.Integer)
    
    val df = Seq(
      Record("a", "b", 1),
      Record("a", "b", null),
      Record("a", null, 0)
    ).toDS
    
    df.show
    
    // +------+------+------+
    // |field1|field2|field3|
    // +------+------+------+
    // |     a|     b|     1|
    // |     a|     b|  null|
    // |     a|  null|     0|
    // +------+------+------+
    
    df.withColumn("thisMap", map_concat(
        df.columns.map { colName => 
            when(col(colName).isNull or col(colName) === 0, map())
            .otherwise(map(lit(colName), col(colName)))
        }: _*
    )).show(false)
    
    // +------+------+------+---------------------------------------+
    // |field1|field2|field3|thisMap                                |
    // +------+------+------+---------------------------------------+
    // |a     |b     |1     |[field1 -> a, field2 -> b, field3 -> 1]|
    // |a     |b     |null  |[field1 -> a, field2 -> b]             |
    // |a     |null  |0     |[field1 -> a]                          |
    // +------+------+------+---------------------------------------+
    

    【讨论】:

      【解决方案2】:

      更新

      我找到了实现预期结果的方法,但它有点脏。

      val df2 = df.columns.foldLeft(df) { (df, n) => df.withColumn(n + "_map", map(lit(n), col(n))) }
      val col_cond = df.columns.map(n => when(not(col(n + "_map").getItem(n).isNull || col(n + "_map").getItem(n) === lit("0")), col(n + "_map")).otherwise(map()))
      df2.withColumn("map", map_concat(col_cond: _*))
        .show(false)
      

      原创

      这是我对可以在 spark 2.4+ 中使用的函数 map_from_arrays 的尝试。

      df.withColumn("array", array(df.columns.map(col): _*))
        .withColumn("map", map_from_arrays(lit(df.columns), $"array")).show(false)
      

      那么,结果是:

      +------+------+------+---------+---------------------------------------+
      |field1|field2|field3|array    |map                                    |
      +------+------+------+---------+---------------------------------------+
      |a     |b     |1     |[a, b, 1]|[field1 -> a, field2 -> b, field3 -> 1]|
      |a     |b     |null  |[a, b,]  |[field1 -> a, field2 -> b, field3 ->]  |
      |a     |null  |0     |[a,, 0]  |[field1 -> a, field2 ->, field3 -> 0]  |
      +------+------+------+---------+---------------------------------------+
      

      【讨论】:

      • hm... 我注意到我有空值映射。我会更多地挖掘它。
      猜你喜欢
      • 2019-11-24
      • 2017-02-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-01-15
      • 1970-01-01
      • 2016-10-11
      • 2017-06-08
      相关资源
      最近更新 更多