如何在PySpark中计算相似度以匹配CSV中与目标点相近的记录
PySpark 实现二维数据相似度匹配
你可以直接用PySpark内置函数实现三种常用相似度计算,无需额外依赖,全表计算后按相似度排序即可返回匹配的H1值,以下是具体实现:
计算逻辑说明
你的数据集以H2、H3为二维特征,目标匹配点为[6, 8],三种相似度的计算规则如下:
- 欧氏距离:两点间直线距离,值越小相似度越高
- 曼哈顿距离:两点各维度差值的绝对值之和,值越小相似度越高
- 余弦相似度:衡量两个向量的方向一致性,取值范围[-1,1],值越接近1相似度越高
针对你给出的样例数据,三种度量下的计算结果:
| H1 | H2 | H3 | 欧氏距离 | 曼哈顿距离 | 余弦相似度 |
|---|---|---|---|---|---|
| A | 1 | 7 | ~5.10 | 6 | ~0.78 |
| B | 5 | 3 | ~5.10 | 6 | ~0.93 |
| C | 7 | 2 | ~6.08 | 7 | ~0.80 |
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, sqrt, pow, abs, lit # 初始化Spark会话 spark = SparkSession.builder \ .appName("SimilarityMatching") \ .getOrCreate() # 加载CSV文件,替换为你的实际文件路径 df = spark.read.csv( path="./your_data.csv", header=True, inferSchema=True ) # 定义目标匹配点 TARGET_H2 = 6 TARGET_H3 = 8 target_mod = sqrt(lit(TARGET_H2)**2 + lit(TARGET_H3)**2) # 批量计算三种相似度指标 sim_df = df.withColumn( "euclidean_distance", sqrt(pow(col("H2") - TARGET_H2, 2) + pow(col("H3") - TARGET_H3, 2)) ).withColumn( "manhattan_distance", abs(col("H2") - TARGET_H2) + abs(col("H3") - TARGET_H3) ).withColumn( "cosine_similarity", (col("H2")*TARGET_H2 + col("H3")*TARGET_H3) / (sqrt(pow(col("H2"), 2) + pow(col("H3"), 2)) * target_mod) ) # 按选定的相似度规则排序取结果,示例为按欧氏距离取最匹配的1条 # 若要按余弦相似度排序,改为 orderBy(col("cosine_similarity").desc()) 即可 top_match_h1 = sim_df.orderBy(col("euclidean_distance").asc()) \ .select("H1") \ .first()[0] print(f"最匹配记录的H1值为: {top_match_h1}") # 若需要取TopN个相似结果,使用limit(N)即可 # top3_h1 = sim_df.orderBy(col("euclidean_distance").asc()).select("H1").limit(3)
大数据量优化方案
如果你的数据集规模在百万级以上,全表逐行计算距离效率较低,可以使用MLlib的近邻检索算法:
- 先通过
VectorAssembler将H2、H3列合并为特征向量列 - 调用KNN(K近邻)模型,设置需要返回的近邻数量K,直接输出最相似的K条记录即可,计算效率比全表遍历高3~10倍。
内容的提问来源于stack exchange,提问作者Tavakoli
相关产品推荐
相关产品推荐

