What would be the most efficient way to embed sentences in a distributed Spark system?

Viewed 1463

I have a file with word embeddings (defining word embedding as the vector representation of a word), with the following format:

a | [0.23, 0.04, ..., -0.22]
aaron | [0.21, 0.08, ..., -0.41]
... | ...
zebra | [0.97, 0.01, ..., -0.34]

This file is about 2.5 GB. I also have a large amount of sentences I want to convert into vectors, for example:

Yes sir, today is a great day.
Would you want to buy that blue shirt?
...
Is there anything else I can help you with?

My sentence embedding strategy is simple for now:

For each sentence:
  For each word:
    Obtain the vector representation of the word using the word embedding file.
  End
  Calculate the average of the word vectors of the sentence.
End

I figured that since I have a large amount of sentences I want to embed, I could use Spark for this task; storing the word embeddings as a file in the HDFS and using Spark SQL to query the sentences from a Hive table, but since each node would likely need to have access to the entire word embedding file that would imply collecting the entire word embedding RDD in each node, making communication between nodes very expensive.

Anyone has any ideas on how could this problem be efficiently solved? Please also let me know if the problem is not clear or you think I've misunderstood something about the way Spark works. I'm still learning and would really appreciate your help!

Thanks in advance.

2 Answers

You could do the following:

  1. Convert your word-embedding file to a Spark DataFrame,
    1. it looks you could use something like my_embeddings = spark.read.csv(path="path/to/your_file.csv", sep="|") pyspark api docs
  2. Change the DataFrame schema (my_embeddings.schema) to match the following:

    1. StructType(List(StructField(word,StringType,true),StructField(vector,ArrayType(FloatType,true),true)))
  3. Create a small-and-simple placeholder Spark ML Word2Vec model and save to hdfs. pyspark api docs
    1. e.g. model_name.write().overwrite().save("your_hdfs_path_to/model_name")
  4. Overwrite the small-and-simple Word2Vec model data with your embedding DataFrame you created above in the your_hdfs_path_to/model_name/data/ directory.
    1. my_embeddings.write.parquet("your_hdfs_path_to/model_name/data/", mode='overwrite')
  5. Load the Word2Vec model with Word2VecModel.load("your_hdfs_path_to/model_name") pyspark api docs
  6. Create a Spark DataFrame where each of your sentences is in a separate row.
  7. Tokenize your sentences with RegexTokenizer pyspark api docs
  8. Use the model to transform the Spark DataFrame that contains your tokenized sentences. The output column will contain a single vector with the same dimensions as the word-embedding vectors, that will be an average of all of the word vectors in the sentence.
    1. "The Word2VecModel transforms each document into a vector using the average of all words in the document" docs. In your case "each document" will be each of your sentences. pyspark api docs

All together (guessing on certain parameters, and using pySpark):

import pyspark
from pyspark.sql import SparkSession
from pyspark.ml.feature import RegexTokenizer
from pyspark.ml.feature import Word2Vec, Word2VecModel
from pyspark.ml import Pipeline, PipelineModel


spark = (
    SparkSession
    .builder
    .master('yarn')
    .appName('my_embeddings')
    .getOrCreate()
)

my_embeddings = spark.read.csv(path="path/to/your_embeddings.csv", sep="|")

my_embeddings.schema
# needs to be
# StructType(List(StructField(word,StringType,true),StructField(vector,ArrayType(FloatType,true),true)))

my_sentences = spark.read.csv(path="path/to/your_sentences.csv", sep="|")

tokenizer = (
    RegexTokenizer()
    .setInputCol("sentences")
    .setOutputCol("tokens") 
    .setPattern("\w+")
)

words2vecs = (
    Word2Vec()
    .setInputCol("tokens")
    .setOutputCol("vecs")
    .setMinCount(1)
    .setNumPartitions(5)
    .setStepSize(0.1)
    .setWindowSize(5)
    .setVectorSize(200)
    .setMaxSentenceLength(1)
)


pipeline = (
    Pipeline()
    .setStages([tokenizer, words2vecs])
)

pipe_model = pipeline.fit(my_sentences.limit(100))

pipe_model.stages[1].write().overwrite().save("your_hdfs_path_to/model_name")

my_embeddings.write.parquet("your_hdfs_path_to/model_name/data/", mode='overwrite')

my_embedding_model = Word2VecModel.load("your_hdfs_path_to/model_name")

df_final = my_embedding_model.transform(tokenizer.transform(my_sentences))

First of all, you word is immutable and you are worried about network efficiency, in your case. I think you could make word a broadcast parameter, so word will be stored in each node locally, and you just transferred the whole word just one time(total N times, N is the count of executors). Then if you will embed word and sentence at the same time which means a network transfer is necessary, it's better to do local reduce before final aggregation.

Related