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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:51:58