【问题标题】:NullPointerException when using Word2VecModel with UserDefinedFunction将 Word2VecModel 与 UserDefinedFunction 一起使用时出现 NullPointerException
【发布时间】:2018-10-07 17:16:57
【问题描述】:

我正在尝试将 word2vec 模型对象传递给我的 spark udf。基本上我有一个带有电影 ID 的测试集,我想将这些 ID 与模型对象一起传递,以获得每行的推荐电影数组。

def udfGetSynonyms(model: org.apache.spark.ml.feature.Word2VecModel) = 
     udf((col : String)  => {
          model.findSynonymsArray("20", 1)
})

但是这给了我一个空指针异常。当我在 udf 之外运行 model.findSynonymsArray("20", 1) 时,我得到了预期的答案。出于某种原因,它不了解 udf 中的函数,但可以在 udf 之外运行它。

注意:我在这里添加了“20”只是为了得到一个固定的答案,看看这是否可行。当我用 col 替换“20”时也是如此。

感谢您的帮助!

堆栈跟踪:

SparkException: Job aborted due to stage failure: Task 0 in stage 23127.0 failed 4 times, most recent failure: Lost task 0.3 in stage 23127.0 (TID 4646648, 10.56.243.178, executor 149): org.apache.spark.SparkException: Failed to execute user defined function($anonfun$udfGetSynonyms1$1: (string) => array<struct<_1:string,_2:double>>)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$10$$anon$1.hasNext(WholeStageCodegenExec.scala:614)
at org.apache.spark.sql.execution.collect.UnsafeRowBatchUtils$.encodeUnsafeRows(UnsafeRowBatchUtils.scala:49)
at org.apache.spark.sql.execution.collect.Collector$$anonfun$2.apply(Collector.scala:126)
at org.apache.spark.sql.execution.collect.Collector$$anonfun$2.apply(Collector.scala:125)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
at org.apache.spark.scheduler.Task.run(Task.scala:111)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:350)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)

Caused by: java.lang.NullPointerException
at org.apache.spark.ml.feature.Word2VecModel.findSynonymsArray(Word2Vec.scala:273)
at linebb57ebe901e04c40a4fba9fb7416f724554.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$anonfun$udfGetSynonyms1$1.apply(command-232354:7)
at linebb57ebe901e04c40a4fba9fb7416f724554.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$anonfun$udfGetSynonyms1$1.apply(command-232354:4)
... 12 more

【问题讨论】:

  • 编辑后,原始帖子变得几乎无关紧要。您能否编辑问题并清理它?编辑后 - 此代码将不起作用,因为 findSynonyms 在内部使用分布式操作。你必须找到另一种方法来解决这个问题。
  • 您确定使用分布式操作吗?我在这里看不到任何东西:github.com/apache/spark/blob/master/mllib/src/main/scala/org/…
  • @JoeK 实际上你的权利,它没有。我检查了这个,我也不能重现 NPE,你可以吗?
  • 我一定会做得更好@user6910411
  • 啊我认为问题是因为我有一个集群..我无法在单台机器上重现 NPE 问题。似乎 wordVectors 在多个节点上运行时可能不可用

标签: scala apache-spark machine-learning nlp word2vec


【解决方案1】:

SQL 和 udf API 有点受限,我不确定是否可以使用自定义类型作为列或作为 udfs 的输入。谷歌搜索并没有发现任何有用的东西。

相反,您可以使用 DataSetRDD API 并使用常规 Scala 函数而不是 udf,例如:

val model: Word2VecModel = ...
val inputs: DataSet[String] = ...
inputs.map(movieId => model.findSynonymsArray(movieId, 10))

或者,我想你可以将模型序列化为字符串,但这看起来更难看。

【讨论】:

  • 这里应该没什么区别。
【解决方案2】:

我认为出现这个问题是因为wordVectors 是一个瞬态变量

class Word2VecModel private[ml] (
    @Since("1.4.0") override val uid: String,
    @transient private val wordVectors: feature.Word2VecModel)
  extends Model[Word2VecModel] with Word2VecBase with MLWritable {

我通过广播 w2vModel.getVectors 并在每个分区内重新创建 Word2VecModel 模型解决了这个问题

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-17
    • 2018-11-22
    • 2018-07-24
    • 1970-01-01
    • 2021-02-28
    相关资源
    最近更新 更多