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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 06:23:14