基于Spark实现ID对应词列表间相似度计算的技术咨询
嘿,我来帮你搞定用Spark计算这些条目相似度的事儿!先看看你的输入数据:
- 101: ["a1","b2","c4","d2"]
- 102: ["a6","b1","c5","d3"]
- 103: ["a1","b4","c4","d2"]
- 104: ["a2","b2","c3","d2"]
- 105: ["a7","b6","c1","d3"]
下面给你三个实用的Spark实现方案,覆盖不同场景需求:
一、先搞定数据准备
不管用哪个方案,第一步都是把数据加载成Spark DataFrame,这里用Scala示例(Python写法类似,核心逻辑一致):
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark = SparkSession.builder() .appName("WordListSimilarity") .master("local[*]") // 生产环境请移除这个配置 .getOrCreate() // 加载你的数据 val data = Seq( (101, Array("a1","b2","c4","d2")), (102, Array("a6","b1","c5","d3")), (103, Array("a1","b4","c4","d2")), (104, Array("a2","b2","c3","d2")), (105, Array("a7","b6","c1","d3")) ).toDF("id", "word_list")
二、方案1:手动计算Jaccard相似度(简单场景首选)
Jaccard相似度是计算两个集合的交集大小除以并集大小,特别适合你这种词列表的场景,直观又准确,适合小数据量:
// 把词列表转去重后的集合(用Array模拟Set) val dfWithSet = data.withColumn("word_set", array_distinct(col("word_list"))) // 做自连接,过滤掉重复配对(比如101和103、103和101只算一次) val pairwiseDf = dfWithSet.as("df1") .join(dfWithSet.as("df2"), col("df1.id") < col("df2.id")) // 定义计算Jaccard的UDF val jaccardUdf = udf((set1: Seq[String], set2: Seq[String]) => { val intersection = set1.intersect(set2).size val union = set1.union(set2).distinct.size if (union == 0) 0.0 else intersection.toDouble / union }) // 计算相似度并筛选高相似结果 val similarityDf = pairwiseDf.withColumn("jaccard_similarity", jaccardUdf(col("df1.word_set"), col("df2.word_set"))) similarityDf.select("df1.id", "df2.id", "jaccard_similarity").filter(col("jaccard_similarity") > 0.3).show()
运行后你会看到101和103的相似度是0.6,101和104是≈0.33,完全符合你说的高相似情况。
三、方案2:TF-IDF + 余弦相似度(考虑词重要性)
如果需要区分词的重要性(比如某个词在很多条目中出现,权重自动降低),可以用TF-IDF把词列表转成加权向量,再计算余弦相似度:
import org.apache.spark.ml.feature.{HashingTF, IDF} import org.apache.spark.ml.linalg.Vector // 把词列表转成空格分隔的文本串(适配TF-IDF输入要求) val dfWithText = data.withColumn("text", concat_ws(" ", col("word_list"))) // 生成TF特征 val hashingTF = new HashingTF() .setInputCol("text") .setOutputCol("rawFeatures") .setNumFeatures(1000) // 根据你的词汇量调整大小 val tfDf = hashingTF.transform(dfWithText) // 生成IDF特征(给TF值加权) val idf = new IDF().setInputCol("rawFeatures").setOutputCol("features") val tfIdfDf = idf.fit(tfDf).transform(tfDf) // 定义余弦相似度UDF val cosineUdf = udf((vec1: Vector, vec2: Vector) => { val dot = vec1.dot(vec2) val norm1 = vec1.norm(2) val norm2 = vec2.norm(2) if (norm1 == 0 || norm2 == 0) 0.0 else dot / (norm1 * norm2) }) // 计算两两相似度 val pairwiseTfIdfDf = tfIdfDf.as("df1").join(tfIdfDf.as("df2"), col("df1.id") < col("df2.id")) val tfIdfSimilarityDf = pairwiseTfIdfDf.withColumn("cosine_similarity", cosineUdf(col("df1.features"), col("df2.features"))) tfIdfSimilarityDf.select("df1.id", "df2.id", "cosine_similarity").show()
这个方案会给每个词赋予差异化权重,适合更复杂的文本相似性分析需求。
四、方案3:LSH局部敏感哈希(大数据量高效方案)
如果你的数据量很大,两两计算会非常耗时,用LSH可以快速找到近似相似的条目,不用全量遍历所有配对:
import org.apache.spark.ml.feature.MinHashLSH // 基于之前的TF特征训练MinHashLSH模型 val minHash = new MinHashLSH() .setNumHashTables(5) // 哈希表数量越多,结果越准但速度越慢,按需调整 .setInputCol("rawFeatures") .setOutputCol("hashes") val lshModel = minHash.fit(tfDf) val lshDf = lshModel.transform(tfDf) // 查找近似相似的条目(这里设置距离阈值0.5,对应Jaccard相似度>0.5) val similarItems = lshModel.approxSimilarityJoin(lshDf, lshDf, 0.5, "jaccard_distance") .filter(col("datasetA.id") < col("datasetB.id")) // 过滤重复配对 .withColumn("jaccard_similarity", lit(1) - col("jaccard_distance")) similarItems.select("datasetA.id", "datasetB.id", "jaccard_similarity").show()
这个方案是近似算法,速度极快,适合百万级以上的大数据场景。
内容的提问来源于stack exchange,提问作者Charmy Garg
相关产品推荐
相关产品推荐

