【问题标题】:How to pass in a map into UDF in spark如何在火花中将地图传递给UDF
【发布时间】:2017-12-19 16:09:10
【问题描述】:

这是我的问题,我有一张 Map[Array[String],String] 的地图,我想将它传递给 UDF。

这是我的 UDF:

def lookup(lookupMap:Map[Array[String],String]) = 
  udf((input:Array[String]) => lookupMap.lift(input))

这是我的 Map 变量:

val srdd = df.rdd.map { row => (
  Array(row.getString(1),row.getString(5),row.getString(8)).map(_.toString),  
  row.getString(7)
)}

这是我调用函数的方式:

val combinedDF  = dftemp.withColumn("a",lookup(lookupMap))(Array($"b",$"c","d"))

我首先得到一个关于不可变数组的错误,所以我将我的数组更改为不可变类型,然后我得到一个关于类型不匹配的错误。我用谷歌搜索了一下,显然我不能将非列类型直接传递给 UDF。有人可以帮忙吗?荣誉。


更新:所以我确实将所有内容都转换为包装数组。这是我所做的:

val srdd = df.rdd.map{row => (WrappedArray.make[String](Array(row.getString(1),row.getString(5),row.getString(8))),row.getString(7))}

val lookupMap = srdd.collectAsMap()


def lookup(lookupMap:Map[collection.mutable.WrappedArray[String],String]) = udf((input:collection.mutable.WrappedArray[String]) => lookupMap.lift(input))


val combinedDF  = dftemp.withColumn("a",lookup(lookupMap))(Array($"b",$"c",$"d"))

现在我遇到这样的错误:

必需:Map[scala.collection.mutable.WrappedArray[String],String] -ksh: Map[scala.collection.mutable.WrappedArray[String],String]: not found [No such file or directory]

我试图做这样的事情:

val m = collection.immutable.Map(1->"one",2->"Two")
val n = collection.mutable.Map(m.toSeq: _*) 

但后来我又回到了列类型的错误。

【问题讨论】:

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


    【解决方案1】:

    首先,您必须将Column 作为UDF 的参数传递;由于您希望此参数是一个数组,因此您应该使用org.apache.spark.sql.functions 中的array 函数,该函数从一系列其他列创建一个数组列。所以 UDF 调用将是:

    lookup(lookupMap)(array($"b",$"c",$"d"))
    

    现在,由于数组列被反序列化为mutable.WrappedArray,为了使映射查找成功,您最好确保这是您的 UDF 使用的类型:

    def lookup(lookupMap: Map[mutable.WrappedArray[String],String]) =
      udf((input: mutable.WrappedArray[String]) => lookupMap.lift(input))
    

    完全是这样:

    import spark.implicits._
    import org.apache.spark.sql.functions._
    
    // Create an RDD[(mutable.WrappedArray[String], String)]:
    val srdd = df.rdd.map { row: Row => (
      mutable.WrappedArray.make[String](Array(row.getString(1), row.getString(5), row.getString(8))), 
      row.getString(7)
    )}
    
    // collect it into a map (I assume this is what you're doing with srdd...)
    val lookupMap: Map[mutable.WrappedArray[String], String] = srdd.collectAsMap()
    
    def lookup(lookupMap: Map[mutable.WrappedArray[String],String]) =
      udf((input: mutable.WrappedArray[String]) => lookupMap.lift(input))
    
    val combinedDF  = dftemp.withColumn("a",lookup(lookupMap)(array($"b",$"c",$"d")))
    

    【讨论】:

    • 如何创建一个包装数组?我做了 val srdd = df.rdd.map{row => (WrappedArray(row.getString(1),row.getString(5),row.getString(8)),row.getString(7))} 但它说明了我包装的数组不带参数
    • 您可以使用mutable.WrappedArray.make[String](Array(...)) - 查看更新的答案
    • 我收到另一个错误是错误:类型不匹配;找到:scala.collection.Map[scala.collection.mutable.WrappedArray[String],String] 需要:scala.collection.immutable.Map[scala.collection.mutable.WrappedArray[String],String] 请看我的更新跨度>
    • 你不需要使用 immutable.Map - 如果你有任何 immutable._immutable.Map 的导入,请删除
    • 哦,抱歉 - 删除 mutable 映射 (mutable.Map) 的导入/使用,您应该只使用不可变类型。
    【解决方案2】:

    安娜,您的 srdd/lookupmap 代码的类型为 org.apache.spark.rdd.RDD[(Array[String], String)]

    val srdd = df.rdd.map { row => (
    Array(row.getString(1),row.getString(5),row.getString(8)).map(_.toString),  
      row.getString(7)
    )}
    

    在查找方法中,您希望将 Map 作为参数

    def lookup(lookupMap:Map[Array[String],String]) = 
    udf((input:Array[String]) => lookupMap.lift(input))
    

    这就是您收到类型不匹配错误的原因。

    首先将 srdd 从 RDD[tuple] 转换为 RDD[Map],然后尝试将 RDD 转换为 Map 以解决此错误。

    val srdd = df.rdd.map { row => Map(
    Array(row.getString(1),row.getString(5),row.getString(8)).map(_.toString) ->
      row.getString(7)
    )}
    

    【讨论】:

    • 嗨,谢谢,但我是这样做的:srdd.collectAsMap(),所以我认为它已经是地图类型了?
    • 是的,我现在可以在您的帖子中编辑后看到您将其转换为地图,但我根据您没有将 srdd 作为地图的旧帖子添加了我的 cmets。
    • @Ram Ghadiyaram 先生任何建议如何处理此 UDF 中的地图stackoverflow.com/questions/63935600/…
    猜你喜欢
    • 2012-03-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多