PySpark如何实现行与行之间的pairwise distance成对距离计算
PySpark实现行成对相关距离的方案
PySpark的Spark ML库有内置功能可以实现和sklearn pairwise_distances(metric='correlation')等价的计算,具体实现方案如下:
前置说明
sklearn中metric='correlation'的成对距离本质是 1 - 两个向量的皮尔逊相关系数,PySpark的实现完全对齐该逻辑。
实现步骤
步骤1:特征向量化
Spark ML的所有数值计算都要求输入为向量格式,首先用VectorAssembler把你数据中的数值列合并为特征向量:
from pyspark.ml.feature import VectorAssembler # 替换inputCols为你实际的数值列名 assembler = VectorAssembler( inputCols=["Mitsubishi", "Toyota", "Tesla", "Honda"], outputCol="features" ) # 得到带特征向量的DataFrame df_vec = assembler.transform(你的原始DataFrame变量名)
步骤2:计算成对距离
分两种情况适配不同Spark版本:
情况1:Spark 3.3及以上版本(推荐)
直接使用内置的PairwiseDistance工具,一行代码即可完成计算:
from pyspark.ml.stat import PairwiseDistance # 指定metric为correlation,和sklearn逻辑完全一致 pairwise_calculator = PairwiseDistance( metric="correlation", inputCol="features", outputCol="pairwise_distance" ) # 得到计算结果 distance_result = pairwise_calculator.transform(df_vec)
你可以直接从结果的pairwise_distance列取出对应距离矩阵,也可以根据需求进一步处理。
情况2:Spark 3.2及以下版本
用RowMatrix做分布式矩阵计算实现等价逻辑:
from pyspark.mllib.linalg.distributed import RowMatrix from pyspark.mllib.linalg import Vectors # 提取特征向量转成RDD vec_rdd = df_vec.select("features").rdd.map(lambda row: Vectors.dense(row.features.toArray())) # 构造分布式行矩阵 row_mat = RowMatrix(vec_rdd) # 转置矩阵后计算列相似性,等价于原矩阵的行相似性 similarity_mat = row_mat.transpose().columnSimilarities(method="pearson") # 转换成相关距离:1 - 相关系数 corr_distance = similarity_mat.map(lambda x: (x.i, x.j, 1 - x.value))
计算得到的corr_distance是包含行索引、对比行索引、对应距离三个字段的RDD,你可以按需转成DataFrame或者收集到本地生成矩阵。
注意事项
- 如果你的数据集规模很大,不要直接调用
collect()把全量距离矩阵拉取到本地,避免 driver 内存溢出 - 除了correlation之外,
PairwiseDistance还支持cosine、euclidean等常用距离指标,直接修改metric参数即可
内容的提问来源于stack exchange,提问作者madfrostx
相关产品推荐
相关产品推荐

