【问题标题】:how to load a word2vec model and call its function into the mapper如何加载 word2vec 模型并将其函数调用到映射器中
【发布时间】:2017-04-24 03:08:27
【问题描述】:

我想加载一个 word2vec 模型并通过执行单词类比任务来评估它(例如,a is to b as c is to something?)。为此,首先我加载我的 w2v 模型:

model = Word2VecModel.load(spark.sparkContext, str(sys.argv[1]))

然后我调用映射器来评估模型:

rdd_lines = spark.read.text("questions-words.txt").rdd.map(getAnswers)

getAnswers 函数每次从 questions-words.txt 中读取一行,其中每一行包含评估我的模型的问题和答案(例如,雅典希腊巴格达伊拉克,其中=雅典,b=希腊,c=巴格达和某事=伊拉克)。阅读该行后,我创建了current_questionactual_answer(例如:current_question=Athens Greece Baghdadactual_answer=Iraq)。之后,我调用用于计算类比的getAnalogy 函数(基本上,给定它计算答案的问题)。最后,计算类比后,我返回答案并将其写入文本文件。

问题是我得到以下异常:

Exception: It appears that you are attempting to reference SparkContext from a broadcast variable, action, or transformation. SparkContext can only be used on the driver, not in code that it run on workers.

我认为它被抛出是因为我在 map 函数中使用模型。这个question 与我的问题类似,但我不知道如何将该答案应用于我的代码。我怎么解决这个问题?以下为完整代码:

def getAnalogy(s, model):
    try:
        qry = model.transform(s[0]) - model.transform(s[1]) - model.transform(s[2])    
        res = model.findSynonyms((-1)*qry,5) # return 5 "synonyms"
        res = [x[0] for x in res]
        for k in range(0,3):
            if s[k] in res:
                res.remove(s[k])
        return res[0]
    except ValueError:
        return "NOT FOUND"

def getAnswers (text):
    tmp = text[0].split(' ', 3)
    answer_list = []
    current_question = " ".join(str(x) for x in tmp[:3])
    actual_answer = tmp[-1]

    model_answer = getAnalogy(current_question, model)
    if model_answer is "NOT FOUND":
        answer_list.append("NOT FOUND\n")
    elif model_answer is actual_answer:
        answer_list.append("TRUE\n")
    else:
        answer_list.append("FALSE:\n")
    return answer_list.append


if __name__ == "__main__":

    if len(sys.argv) != 3:
        print("Usage: my_test <file>", file=sys.stderr)
        exit(-1)


    spark = SparkSession\
    .builder\
    .appName("my_test")\
    .getOrCreate()


    model = Word2VecModel.load(spark.sparkContext, str(sys.argv[1]))

    rdd_lines = spark.read.text("questions-words.txt").rdd.map(getAnswers)

    dataframe = rdd_lines.toDF()

    dataframe.write.text(str(sys.argv[2]))

    spark.stop()

【问题讨论】:

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


    【解决方案1】:

    正如您已经怀疑的那样,您不能在地图功能中使用模型。另一方面,questions-answers.txt 文件并没有那么大(约 20K 行),因此您最好使用普通 Python 列表推导进行评估(它本质上是您链接的问题中的第一个建议答案);它并不快,但这只是一次性的任务。这是一种使用my getAnalogy function 的方法,因为您已对其进行了增强以进行错误处理(请注意,我已经从questions-answers.txt 中删除了“注释”行,并且您应该将其转换为小写,而您似乎没有在你的代码中做):

    from pyspark.mllib.feature import Word2Vec, Word2VecModel
    model = Word2VecModel.load(sc, "word2vec/demo_200") # model built with k=200
    with open('/home/ctsats/word2vec/questions-words.txt') as f:
        lines = f.readlines()
    lines2 = [x.lower() for x in lines] # all to lowercase
    lines3 = [x.strip('\n') for x in lines2] # remove end-of-line characters
    lines4 = [x.split(' ',3) for x in lines3]
    lines4[0] # check:
    # ['Athens', 'Greece', 'Baghdad', 'Iraq']
    
    def getAnswers (text, model):
        actual_answer = text[-1]
        question = [text[0], text[1], text[2]]
        model_answer = getAnalogy(question, model)
        if model_answer == "NOT FOUND":
            correct_answer = "NOT FOUND"
        elif model_answer == actual_answer:
            correct_answer = "TRUE"
        else:
            correct_answer = "FALSE"
        return text, model_answer, correct_answer
    

    因此,您的评估列表现在可以构建为

    answer_list = [getAnswers(x, model) for x in lines4]    
    

    这是前 20 个条目的示例(模型为 k=200):

    [(['athens', 'greece', 'baghdad', 'iraq'], u'turkey', 'FALSE'),
     (['athens', 'greece', 'bangkok', 'thailand'], u'turkey', 'FALSE'),
     (['athens', 'greece', 'beijing', 'china'], u'albania', 'FALSE'),
     (['athens', 'greece', 'berlin', 'germany'], u'germany', 'TRUE'),
     (['athens', 'greece', 'bern', 'switzerland'], u'liechtenstein', 'FALSE'),
     (['athens', 'greece', 'cairo', 'egypt'], u'albania', 'FALSE'),
     (['athens', 'greece', 'canberra', 'australia'], u'liechtenstein', 'FALSE'),
     (['athens', 'greece', 'hanoi', 'vietnam'], u'turkey', 'FALSE'),
     (['athens', 'greece', 'havana', 'cuba'], u'turkey', 'FALSE'),
     (['athens', 'greece', 'helsinki', 'finland'], u'finland', 'TRUE'),
     (['athens', 'greece', 'islamabad', 'pakistan'], u'turkey', 'FALSE'),
     (['athens', 'greece', 'kabul', 'afghanistan'], u'albania', 'FALSE'),
     (['athens', 'greece', 'london', 'england'], u'italy', 'FALSE'),
     (['athens', 'greece', 'madrid', 'spain'], u'portugal', 'FALSE'),
     (['athens', 'greece', 'moscow', 'russia'], u'russia', 'TRUE'),
     (['athens', 'greece', 'oslo', 'norway'], u'albania', 'FALSE'),
     (['athens', 'greece', 'ottawa', 'canada'], u'moldova', 'FALSE'),
     (['athens', 'greece', 'paris', 'france'], u'france', 'TRUE'),
     (['athens', 'greece', 'rome', 'italy'], u'italy', 'TRUE'),
     (['athens', 'greece', 'stockholm', 'sweden'], u'norway', 'FALSE')]
    

    【讨论】:

    • 谢谢,这绝对是一个不错的选择,但问题是我必须评估许多模型,我想在 map 函数中调用 w2v 模型函数(例如 findSynonyms())。
    • 正如我所说,使用 Spark 根本无法做到这一点。或者,您可以尝试另一个 word2vec 实现(例如 gensim)并将其包含在您的地图函数中。
    • 好的,我明白了.....但是如何使用 pyspark 保存 w2v 模型,然后使用 gensim 加载它?问题是用 pyspark 保存的 w2v 模型是一组 parquet 文件,而 gensim 存储“.model”文件(而不是 parquet)...
    • @Juniorhpc 你不能那样做;您必须使用 gensim 构建并保存模型,然后将此 gensim 模型包含在您的地图功能中(虽然尚未测试)
    猜你喜欢
    • 2016-08-23
    • 1970-01-01
    • 2023-02-01
    • 2016-09-19
    • 1970-01-01
    • 1970-01-01
    • 2016-10-12
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多