代码之家  ›  专栏  ›  技术社区  ›  Markus

在Spark中使用UDF时出现任务序列化错误

  •  1
  • Markus  · 技术社区  · 7 年前

    当我创建如上所示的UDF函数时,我得到了任务序列化错误。只有在使用在集群部署模式下运行代码时,才会出现此错误 spark-submit . 然而,它在火花壳中工作良好。

    import org.apache.spark.sql.expressions.Window
    import org.apache.spark.sql.SparkSession
    import org.apache.spark.sql.functions._
    import scala.collection.mutable.WrappedArray
    
    def mfnURL(arr: WrappedArray[String]): String = {
      val filterArr = arr.filterNot(_ == null)
      if (filterArr.length == 0)
        return null
      else {
        filterArr.groupBy(identity).maxBy(_._2.size)._1
      }
    }
    
    val mfnURLUDF = udf(mfnURL _)
    
    def windowSpec = Window.partitionBy("nodeId", "url", "typology")                                                     
    val result = df.withColumn("count", count("url").over(windowSpec))
      .orderBy($"count".desc)                                                                                            
      .groupBy("nodeId","typology")                                                                                      
      .agg(
      first("url"),
      mfnURLUDF(collect_list("source_url")),
      min("minTimestamp"),
      max("maxTimestamp")
    )
    

    spark.udf.register("mfnURLUDF",mfnURLUDF) ,但并没有解决问题。

    1 回复  |  直到 7 年前
        1
  •  2
  •   merenptah    7 年前

    val mfnURL = udf { arr: WrappedArray[String] =>
      val filterArr = arr.filterNot(_ == null)
      if (filterArr.length == 0)
        return null
      else {
        filterArr.groupBy(identity).maxBy(_._2.size)._1
      }
    }
    
    推荐文章