Матрица косинусного сходства с Pyspark MS Fabric SparkJob ⇐ Python

Программы на Python
Anonymous
Матрица косинусного сходства с Pyspark MS Fabric SparkJob

Сообщение Anonymous »

У меня возникли проблемы при вычислении некоторых косинусных сходств для рекомендации продукта. У меня есть база данных статей, содержащая 40 тысяч статей, каждая из которых имеет описание. Я пытаюсь вычислить матрицу косинусного сходства для этих элементов, чтобы при задании какой-либо статьи можно было получить первые N, наиболее похожие по описанию.
Я разрабатываю это на Python и код следует запускать в функции Microsoft Azure. Как вы можете себе представить, учитывая размер исходного df, его нельзя запустить на этой платформе (ему не хватает памяти, поскольку выводимое косинусное сходство имеет размер более 8 ГБ).
Имея это ввиду, я решил применить другой подход, который использует Spark, точнее, определение задания Spark из Microsoft Fabric. Я просматривал подобные вопросы (поскольку мои знания Pyspark крайне ограничены), например этот: Вычисление сходства косинусов в Pyspark Dataframe, но я изо всех сил пытаюсь заставить код работать в моей ситуации.
Пытаясь решать проблемы одну за другой, я понял, что мой код даже не доходит до расчета сходства косинусов (или, по крайней мере, не всего), поскольку моя искровая работа зависает при попытке перекрестного соединения матриц, и я не могу проверить в чем именно заключается проблема из журналов stderr.
Я знаю, что мой начальный df из 40 тыс. строк довольно велик, и это подразумевает матрицу cos sim размером 40 x 40 тыс., но я попробовал эту операцию с простой Python на ноутбуке Google Collab, и он смог вычислить его примерно за 10 минут (это правда, что мне пришлось использовать серверную часть TPU, поскольку обычный процессор также выходит из строя после использования стандартной памяти -12 ГБ-).
На данный момент я думаю либо:
-MS Fabric SparkJobs также не является подходящим инструментом для моих нужд, поскольку они предлагают «слабую/ограниченную» версию Spark.
-В моем коде Pyspark явно есть что-то неправильное, что приводит к ошибочному выводу из эксплуатации исполнителя Spark.
Мне нужна помощь, чтобы определить, в каком сценарии я нахожусь ( а если на втором, то как улучшить код).
Позвольте мне поделиться кодом.
Это код Python, который для той же задачи действительно работает на ноутбуке TPU Collab. (с использованием импорта sklearn):

Код: Выделить всё

    articulos = articulos[['article_id','product_code', 'detail_desc']]
articulos_unicos = articulos.drop_duplicates(subset=['product_code'])
articulos_final = articulos_unicos.dropna(subset=['detail_desc'])
articulos_final = articulos_final.reset_index(drop=True)

count = CountVectorizer(stop_words='english')
count_matrix = count.fit_transform(articulos_final['detail_desc'])
count_matrix = count_matrix.astype(np.float32)

similitud_coseno = cosine_similarity(count_matrix, count_matrix)

np.fill_diagonal(similitud_coseno, -1)

top_5_simart = []

for i in range(similitud_coseno.shape[0]):
top_indices = np.argpartition(similitud_coseno[i], -5)[-5:]
sorted_top_indices = top_indices[np.argsort(-similitud_coseno[i, top_indices])]
top_5_simart.append(sorted_top_indices.tolist())

with open('top_5_simart.json', 'w') as f:
json.dump(top_5_simart, f)
С другой стороны, это код Pyspark, который я пытаюсь реализовать:

Код: Выделить всё

    # Selecting the necessary columns
articulos = articulos.select("article_id", "product_code", "detail_desc")

# Removing duplicates
articulos = articulos.dropDuplicates(["product_code"]).dropna(subset=["detail_desc"])

# Tokenizing the description text
tokenizer = Tokenizer(inputCol="detail_desc", outputCol="words")
articulos_tokenized = tokenizer.transform(articulos)

# Removing stopwords
remover = StopWordsRemover(inputCol="words", outputCol="filtered_words")
articulos_clean = remover.transform(articulos_tokenized)

# Generating the feature column with CountVectorizer
vectorizer = CountVectorizer(inputCol="filtered_words", outputCol="features")
vectorizer_model = vectorizer.fit(articulos_clean)
articulos_final = vectorizer_model.transform(articulos_clean)
articulos_final = articulos_final.select("article_id","features")
print(articulos_final.dtypes)
print(articulos_final.schema)

# Using crossJoin to obatain all the pairwise values for the cosine sim
articulos_final2 = articulos_final.withColumnRenamed("article_id","article_id2").withColumnRenamed("features","features2")
articulos_final_cos_sim = articulos_final.crossJoin(articulos_final2)
articulos_final.unpersist()
articulos_final2.unpersist()
articulos_final_cos_sim.write.mode('overwrite').format('delta').save(cspath)

'''
# Realizar el crossJoin para obtener todas las combinaciones de pares de filas
articulos_final_cos_sim = (
articulos_final.alias("a")
.crossJoin(articulos_final.alias("b"))
.withColumn(
"dot_product",
F.col("a.features").dot(F.col("b.features"))  # Producto punto de los vectores
)
.withColumn(
"norm_a",
F.expr("a.features.norm(2)")  # Norma L2 del vector 'a'
)
.withColumn(
"norm_b",
F.expr("b.features.norm(2)")  # Norma L2 del vector 'b'
)
.withColumn(
"cosine_similarity",
F.col("dot_product") / (F.col("norm_a") * F.col("norm_b"))  # Similitud coseno
)
)

# Eliminar las columnas innecesarias
articulos_final_cos_sim = articulos_final_cos_sim.drop("dot_product", "norm_a", "norm_b")

# Agrupar para construir la matriz de similitud
articulos_final_cos_sim = articulos_final_cos_sim.groupBy("a.article_id").pivot("b.article_id").sum("cosine_similarity")

# Guardar los resultados
articulos_final_cos_sim.write.mode('overwrite').format('delta').save(testpath2)

# Filtrar para obtener los 5 artículos más similares para cada uno
windowSpec = Window.partitionBy("article_id").orderBy(F.col("cosine_similarity").desc())

top_5_simart = articulos_final_cos_sim.withColumn("rank", F.row_number().over(windowSpec)).filter(F.col("rank") 

Подробнее здесь: [url]https://stackoverflow.com/questions/79057851/cosine-similarity-matrix-with-pyspark-ms-fabric-sparkjob[/url]

Вернуться в «Python»