Spark DataFrame中计算两列向量间欧氏距离的实现问题
在Spark DataFrame中计算两向量的欧氏距离
我尝试在Spark DataFrame的不同列中计算两个向量的欧氏距离,但折腾了很久都没搞定,以下是和我场景类似的示例代码:
from pyspark.sql import SparkSession from pyspark.ml.linalg import Vector, Vectors from pyspark.sql.functions import expr from pyspark.sql import SparkSession, types as T, functions as F # Initialize Spark session spark = SparkSession.builder \ .appName("euclidian distance") \ .getOrCreate() # Sample data data = [(Vectors.dense([1.0, 2.0]), Vectors.dense([3.0, 4.0])), (Vectors.dense([5.0, 6.0]), Vectors.dense([7.0, 8.0])), (Vectors.dense([9.0, 10.0]), Vectors.dense([11.0, 12.0]))] # Create DataFrame df = spark.createDataFrame(data, ["vector1", "vector2"]) # Define UDF for vector subtraction def ed_vectors_udf(v1, v2): return Vectors.dense(v1).squared_distance(Vectors.dense(v2)) # return Vectors.dense([x - y for x, y in zip(v1, v2)]) # Register UDF # ed_vectors_udf = fn.udf(lambda v1, v2: eq_vectors(v1, v2), T.DoubleType()) spark.udf.register("ed_vectors_udf", ed_vectors_udf, T.DoubleType()) # Subtract vectors using UDF df = df.withColumn("distance", ed_vectors_udf(fn.col('vector1'), fn.col('vector2'))) # Show DataFrame with subtraction result df.show(5)
问题分析
原代码存在几个关键问题:
- 重复导入
SparkSession - 使用了未定义的
fn(应为F) squared_distance返回的是平方欧氏距离,并非真实欧氏距离,需取平方根- 无需将已有的
Vector类型再次转换为Vectors.dense,输入本身就是Vector对象
解决方案
方法一:修正自定义UDF实现欧氏距离
from pyspark.sql import SparkSession from pyspark.ml.linalg import Vectors from pyspark.sql import functions as F, types as T # 初始化Spark会话 spark = SparkSession.builder \ .appName("euclidean distance") \ .getOrCreate() # 示例数据 data = [(Vectors.dense([1.0, 2.0]), Vectors.dense([3.0, 4.0])), (Vectors.dense([5.0, 6.0]), Vectors.dense([7.0, 8.0])), (Vectors.dense([9.0, 10.0]), Vectors.dense([11.0, 12.0]))] # 创建DataFrame df = spark.createDataFrame(data, ["vector1", "vector2"]) # 定义计算欧氏距离的UDF def euclidean_distance(v1, v2): # 计算平方距离后取平方根得到欧氏距离 return float(v1.squared_distance(v2)**0.5) # 注册UDF ed_udf = F.udf(euclidean_distance, T.DoubleType()) # 添加距离列 df = df.withColumn("distance", ed_udf(F.col('vector1'), F.col('vector2'))) # 展示结果 df.show(5)
方法二:使用Spark SQL内置函数(无需UDF)
直接通过Spark表达式实现,避免自定义UDF的性能开销:
# 使用expr计算欧氏距离 df = df.withColumn( "distance", F.expr("sqrt(sum(pow((vector1 - vector2)[i], 2) for i in range(size(vector1))))") ) df.show(5)
方法三:使用Spark ML内置距离工具
借助Spark ML提供的EuclideanDistance类简化计算:
from pyspark.ml.linalg.distances import EuclideanDistance distance_calculator = EuclideanDistance() ed_udf = F.udf(lambda v1, v2: distance_calculator.compute(v1, v2), T.DoubleType()) df = df.withColumn("distance", ed_udf(F.col('vector1'), F.col('vector2'))) df.show(5)
内容的提问来源于stack exchange,提问作者Prabin Shrestha
相关产品推荐
相关产品推荐

