You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark中如何基于TF-IDF计算余弦相似度实现新闻推荐?

PySpark实现新闻推荐:基于TF-IDF与余弦相似度

先修正你的TF-IDF计算错误

你之前的PySpark代码分别对两个数据集训练IDF模型,这会导致两个数据集的TF-IDF特征空间不一致(词汇表、IDF权重不同),这是结果和Python版差异大的核心原因。必须用同一个CountVectorizer和IDF模型处理两个数据集,和你原来的Python逻辑对齐:

from pyspark.sql import SparkSession
from pyspark.ml.feature import Tokenizer, CountVectorizer, IDF, StopWordsRemover
from pyspark.sql.functions import col, monotonically_increasing_id
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, udf
from pyspark.ml.linalg import VectorUDT

# 初始化Spark会话
spark = SparkSession.builder.appName("NewsRec").getOrCreate()

# 读取数据
df1 = spark.read.csv('bbcclear.csv', header=True, inferSchema=True)  # 对应原Python中的det数据集(用于训练TF-IDF)
df2 = spark.read.csv('yenisafakcategorypredict.csv', header=True, inferSchema=True)  # 对应原Python中的ted数据集(候选推荐集)

# 文本预处理:分词+去停用词(和sklearn的stop_words='english'对齐)
tokenizer = Tokenizer(inputCol="News Data", outputCol="raw_words")
stop_remover = StopWordsRemover(inputCol="raw_words", outputCol="filtered_words", stopWords=StopWordsRemover.loadDefaultStopWords("english"))

# 处理两个数据集的文本
df1_processed = tokenizer.transform(df1)
df1_processed = stop_remover.transform(df1_processed)

df2_processed = tokenizer.transform(df2)
df2_processed = stop_remover.transform(df2_processed)

# 训练CountVectorizer(基于df1的文本,和原Python逻辑一致)
vectorizer = CountVectorizer(inputCol="filtered_words", outputCol="count_vec").fit(df1_processed)
df1_vec = vectorizer.transform(df1_processed)
df2_vec = vectorizer.transform(df2_processed)

# 训练IDF模型(同样基于df1的向量)
idf = IDF(inputCol="count_vec", outputCol="tfidf").fit(df1_vec)
df1_tfidf = idf.transform(df1_vec)
df2_tfidf = idf.transform(df2_vec)

# 给数据集添加唯一标识(如果原数据有类似Unnamed:0的ID列,直接用col("Unnamed:0")即可)
df1_tfidf = df1_tfidf.withColumn("target_id", monotonically_increasing_id())
df2_tfidf = df2_tfidf.withColumn("candidate_id", monotonically_increasing_id()).withColumnRenamed("News Data", "candidate_news")

计算余弦相似度并生成推荐

1. 定义余弦相似度UDF

PySpark没有内置的跨数据集余弦相似度计算函数,我们自己实现:

def calc_cosine(vec1, vec2):
    if vec1 is None or vec2 is None:
        return 0.0
    dot = vec1.dot(vec2)
    norm1 = vec1.norm(2)
    norm2 = vec2.norm(2)
    if norm1 == 0 or norm2 == 0:
        return 0.0
    return dot / (norm1 * norm2)

# 注册为UDF
cosine_udf = udf(calc_cosine, returnType="double")

2. 计算两两相似度并取Top10推荐

通过笛卡尔积得到所有目标新闻和候选新闻的相似度,再用窗口函数筛选每个目标的Top10:

# 重命名TF-IDF列避免冲突
df1_tfidf = df1_tfidf.withColumnRenamed("tfidf", "target_tfidf")
df2_tfidf = df2_tfidf.withColumnRenamed("tfidf", "candidate_tfidf")

# 计算两两相似度
cross_df = df1_tfidf.crossJoin(df2_tfidf).withColumn("similarity", cosine_udf(col("target_tfidf"), col("candidate_tfidf")))

# 用窗口函数对每个目标新闻的候选结果排序,取Top10
window = Window.partitionBy("target_id").orderBy(col("similarity").desc())
top10_recs = cross_df.withColumn("rank", row_number().over(window))\
    .filter(col("rank") <= 10)\
    .select("target_id", "candidate_id", "candidate_news", "similarity", "rank")

# 示例:获取target_id=0的推荐结果
top10_recs.filter(col("target_id") == 0).show(truncate=False)

关键说明

  • 必须保证两个数据集的TF-IDF在同一特征空间:所有预处理、向量化、IDF训练都基于同一个基准数据集(这里是df1),和你原来的Python逻辑完全对齐。
  • 性能优化:如果数据量很大,笛卡尔积会导致性能问题,可以考虑用近似相似度算法(比如LSH)来减少计算量,不过对于新手来说,先实现基础版本再优化。

内容的提问来源于stack exchange,提问作者Alp Buğra Aker

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 14:50:38