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.
相关产品推荐
相关产品推荐

