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

PySpark计算两个DataFrame条目间euclidean distance报错如何解决

报错原因

你的UDF输入参数要求是DataFrame的列对象,但你直接传入了a、b两个完整的DataFrame实例,不符合参数类型要求。此外计算两个DataFrame的行向量距离,需要先把两个向量拉到同一个DataFrame的同一行中才能进行计算。

修复方案

步骤1:合并两个DataFrame

因为两个DataFrame都只有1行,直接用笛卡尔积拼接即可,先重命名列避免冲突:

# 假设你的两个DataFrame分别名为df_a、df_b
df_a = df_a.withColumnRenamed("entry", "entry_a")
df_b = df_b.withColumnRenamed("entry", "entry_b")

# 拼接为单行列的合并表
combined_df = df_a.crossJoin(df_b)

步骤2:调用UDF计算欧氏距离

将合并后的两个列作为参数传入UDF即可:

from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType
from scipy.spatial import distance

inference = udf(lambda x, y: float(distance.euclidean(x, y)), DoubleType())

# 传入列名计算距离
result_df = combined_df.withColumn("euclidean_distance", inference("entry_a", "entry_b"))

# 输出结果
result_df.select("euclidean_distance").show()
性能优化可选方案:原生Spark函数实现

如果不想引入第三方依赖,也不需要UDF的序列化开销,可以直接用Spark内置的数组高阶函数计算:

from pyspark.sql.functions import transform, aggregate, sqrt, lit

result_df = combined_df.withColumn("distance",
    sqrt(
        aggregate(
            # 逐位计算差的平方和
            transform("entry_a", lambda val, idx: (val - combined_df.entry_b[idx]) ** 2),
            lit(0.0),
            lambda acc, cur: acc + cur
        )
    )
)

内容的提问来源于stack exchange,提问作者A.M.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 20:24:02