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
相关产品推荐
相关产品推荐

