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

如何在Spark Scala环境下计算DataFrame每行与常量参考数组的欧氏距离?

计算DataFrame每行与常量数组的欧氏距离(适配Spark 2.1 + Scala + Zeppelin)

针对你这个512列浮点型DataFrame的需求,我整理了完全适配你环境的解决方案,步骤清晰且高效:

完整代码实现

import org.apache.spark.ml.feature.VectorAssembler
import org.apache.spark.ml.linalg.Vectors
import org.apache.spark.sql.functions.udf
import scala.math.sqrt

// 从Parquet文件加载DataFrame
val filePath = "/tmp/vector.parquet/*.parquet"
val df = spark.read.parquet(filePath)

// 1. 自动获取所有512个特征列的列名
val featureCols = df.columns

// 2. 将分散的列合并成单个向量列(Spark ML处理向量数据的标准操作)
val assembler = new VectorAssembler()
  .setInputCols(featureCols)
  .setOutputCol("features")

val dfWithVector = assembler.transform(df)

// 3. 定义你的常量参考数组(这里用全1数组示例,替换成你实际的参考值即可)
val referenceVector = Vectors.dense(Array.fill(512)(1.0))

// 4. 写一个UDF来计算欧氏距离——用Spark内置的sqdist高效计算平方距离再开根号
val calculateEuclideanDistance = udf((rowVector: org.apache.spark.ml.linalg.Vector) => {
  sqrt(Vectors.sqdist(rowVector, referenceVector))
})

// 5. 把UDF应用到DataFrame上,生成距离列
val dfWithDistance = dfWithVector.withColumn("euclidean_distance", calculateEuclideanDistance($"features"))

// 可选:查看前5条结果(数据量大的话记得加limit,别全量show)
dfWithDistance.select("features", "euclidean_distance").show(5)

关键细节说明

  • 为什么用VectorAssembler?:把512列合并成一个向量列是最高效的方式,比手动遍历每列计算快得多,而且符合Spark ML的向量数据格式规范。
  • 欧氏距离的高效计算:Spark的Vectors.sqdist是底层优化过的方法,专门用来计算两个向量的平方欧氏距离,我们只需要对结果取平方根就能得到标准欧氏距离,避免自己写循环带来的性能损耗。
  • Spark 2.1的类型注意点:这个版本里ml.linalg.Vector是强类型,所以UDF的输入参数必须明确指定org.apache.spark.ml.linalg.Vector,不然容易出现类型推断错误。

小优化建议

  • 如果你的数据量特别大,可以先给DataFrame调整分区数(比如df.repartition(200)),让集群资源利用更充分。
  • 最后如果不需要保留features列,记得用drop("features")删掉,节省内存空间。

内容的提问来源于stack exchange,提问作者ozlemg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:16:33