【问题标题】:UDF with Dynamic Data Type具有动态数据类型的 UDF
【发布时间】:2021-08-07 23:43:44
【问题描述】:

我正在尝试编写可以从 Map 中删除几个键的 udf。但是 Map 的 key 和 value 类型并不固定,可以是 String 或者 Array 或者其他。我应该如何定义这样的udf。我使用的是 Spark 2.4.4 版。

下面是我对 Map[String, string] 的 udf:

val mapKeys = //Seq[String]
val mapFilterUdf = udf[Map[String, String], Map[String, String]] {
    map => map.filter{case (key, _) => mapKeys.contains(key)}
}
mapFilterUdf(dataFrame.col("column_name")).as(column.name)

【问题讨论】:

    标签: scala apache-spark apache-spark-sql user-defined-functions


    【解决方案1】:

    你可以为 udf 做一个通用的工厂方法:

    import scala.reflect.runtime.universe._
    
    def filterUdfFactory[T](mapKeys:Seq[T])(implicit tag:TypeTag[T]) = udf((map:Map[T,T]) => map.filter{case (k,v) => mapKeys.contains(k)})
    

    然后用作例如对于字符串:

    val mapKeys = Seq("k1")
    
    val tt = typeTag[String]
    val filterUdf = filterUdfFactory[String](mapKeys)
    
     val df = Seq(
        Map("k1" -> "v1","k2" -> "v2")
     ).toDF("map")
    
     df.select(filterUdf($"map"))
    .show()
    

    给予:

    +----------+
    |  UDF(map)|
    +----------+
    |[k1 -> v1]|
    +----------+
    

    【讨论】:

    • 在上面的代码中,您正在调用类型为“String”的“filterUdfFactory”。但就我而言,Map 的键类型和值类型仅在运行时可用。
    【解决方案2】:

    仅当您将列的运行时模式作为第二个参数提供给 udf 时,您才能在 UDF 中使用 Any

    val mapKeys : Seq[Any] = Seq("k1")
    
    val df = Seq(
        Map("k1" -> "v1","k2" -> "v2")
    ).toDF("map")
    
    val colSchema = df.select($"map").schema.head.dataType
    
    val filterUdf = udf((map:Map[Any,Any]) => map.filter{case (k:Any,v:Any) => mapKeys.contains(k)},colSchema)
    
    df
    .select(filterUdf($"map"))
    .show()
    

    给予

    +----------+
    |  UDF(map)|
    +----------+
    |[k1 -> v1]|
    +----------+
    

    这项工作也适用于Row,请参阅:https://stackoverflow.com/a/49714640/1138523

    【讨论】:

      猜你喜欢
      • 2017-06-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-02-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多