Я разрабатываю это на 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.
Мне нужна помощь, чтобы определить, в каком сценарии я нахожусь (и, если второй, как улучшить код).< /p>
Позвольте мне поделиться кодом.
Это код 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 neccesary 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]