【问题标题】:using Word2VecModel.transform() does not work in map function使用 Word2VecModel.transform() 在地图功能中不起作用
【发布时间】:2018-11-05 05:31:11
【问题描述】:

我已经使用 Spark 构建了一个 Word2Vec 模型并将其保存为模型。现在,我想在另一个代码中使用它作为离线模型。我已经加载了模型并用它来呈现一个单词的向量(例如你好)并且效果很好。但是,我需要在 RDD 中使用 map 来调用它。

当我在 map 函数中调用 model.transform() 时,它会抛出这个错误:

“您似乎正在尝试从广播中引用 SparkContext” 例外:您似乎正试图从广播变量、操作或转换中引用 SparkContext。 SparkContext 只能在驱动程序上使用,不能在它在工作人员上运行的代码中使用。有关详细信息,请参阅 SPARK-5063。

代码:

from pyspark import SparkContext
from pyspark.mllib.feature import Word2Vec
from pyspark.mllib.feature import Word2VecModel

sc = SparkContext('local[4]',appName='Word2Vec')

model=Word2VecModel.load(sc, "word2vecModel")

x= model.transform("Hello")
print(x[0]) # it works fine and returns [0.234, 0.800,....]

y=sc.parallelize([['Hello'],['test']])
y.map(lambda w: model.transform(w[0])).collect() #it throws the error

非常感谢您的帮助。

【问题讨论】:

    标签: python apache-spark pyspark apache-spark-mllib word2vec


    【解决方案1】:

    这是一种预期的行为。与其他 MLlib 模型一样,Python 对象只是 Scala 模型的包装器,实际处理委托给其对应的 JVM。由于工作人员无法访问 Py4J 网关(请参阅 How to use Java/Scala function from an action or a transformation?),因此您无法从操作或转换中调用 Java / Scala 方法。

    通常 MLlib 模型提供了一个辅助方法,可以直接在 RDD 上工作,但这里不是这样。 Word2VecModel 提供了 getVectors 方法,该方法返回从单词到向量的映射,但不幸的是它是 JavaMap,因此它在转换中不起作用。你可以试试这样的:

    from pyspark.mllib.linalg import DenseVector
    
    vectors_ = model.getVectors() # py4j.java_collections.JavaMap
    vectors = {k: DenseVector([x for x in vectors_.get(k)])
        for k in vectors_.keys()}
    

    获取 Python 字典,但它会非常慢。另一种选择是以 Python 可以使用的形式将此对象转储到磁盘,但它需要对 Py4J 进行一些修补,最好避免这种情况。而是让我们将模型读取为 DataFrame:

    lookup = sqlContext.read.parquet("path_to_word2vec_model/data").alias("lookup")
    

    我们将得到以下结构:

    lookup.printSchema()
    ## root
    ## |-- word: string (nullable = true)
    ## |-- vector: array (nullable = true)
    ## |    |-- element: float (containsNull = true)
    

    可用于将单词映射到向量,例如通过join:

    from pyspark.sql.functions import col
    
    words = sc.parallelize([('hello', ), ('test', )]).toDF(["word"]).alias("words")
    
    words.join(lookup, col("words.word") == col("lookup.word"))
    
    ## +-----+-----+--------------------+
    ## | word| word|              vector|
    ## +-----+-----+--------------------+
    ## |hello|hello|[-0.030862354, -0...|
    ## | test| test|[-0.13154022, 0.2...|
    ## +-----+-----+--------------------+
    

    如果数据适合驱动程序/工作人员内存,您可以尝试使用广播进行收集和映射:

    lookup_bd = sc.broadcast(lookup.rdd.collectAsMap())
    rdd = sc.parallelize([['Hello'],['test']])
    rdd.map(lambda ws: [lookup_bd.value.get(w) for w in ws])
    

    【讨论】:

    • 您好 zero323。谢谢你回答我的问题。确实,我应该更具体地问我的问题。我的问题是我有一个文档的 RDD(标记化的推文),而不是一个单词的 RDD。我需要通过找到每个单词的向量并平均来对每个文档进行向量化。有什么想法吗?
    • 你可以尝试收藏和广播。
    猜你喜欢
    • 2017-02-21
    • 1970-01-01
    • 1970-01-01
    • 2019-07-31
    • 2018-04-13
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多