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

基于PySpark的KMeans算法:为DataFrame添加样本至聚类中心的距离列

为PySpark DataFrame添加特征与所属聚类中心的距离列

嘿,这个需求我熟,给你一步步来实现,代码和逻辑都给你理清楚:

核心思路

我们需要先拿到KMeans模型的所有聚类中心,通过广播变量让集群节点高效访问,再定义一个UDF来计算每条数据的特征向量和对应聚类中心的距离,最后把这个距离作为新列添加到DataFrame里。

步骤1:获取并广播聚类中心

首先从训练好的KMeans模型中提取所有聚类中心,转换成numpy数组后用Spark的广播变量分发,避免每个任务重复加载中心数据,提升效率:

from pyspark.sql.functions import udf
from pyspark.ml.linalg import Vectors
import numpy as np

# 假设你的训练好的KMeans模型叫kmeans_model
cluster_centers = np.array([center.toArray() for center in kmeans_model.clusterCenters()])
# 广播聚类中心到所有worker节点
broadcast_centers = spark.sparkContext.broadcast(cluster_centers)

步骤2:定义计算距离的UDF

这里我们用KMeans默认的欧氏距离来计算,你也可以根据需求换成其他距离(比如曼哈顿距离):

from pyspark.sql.types import DoubleType

def calculate_distance(features_vec, cluster_id):
    # 将Spark的Vector类型转换为numpy数组
    feature_array = features_vec.toArray()
    # 根据聚类ID获取对应的中心
    target_center = broadcast_centers.value[cluster_id]
    # 计算欧氏距离并返回浮点型结果
    return float(np.linalg.norm(feature_array - target_center))

# 注册UDF,指定返回类型为DoubleType
distance_udf = udf(calculate_distance, DoubleType())

步骤3:添加距离列到DataFrame

假设你的原始DataFrame叫df,其中包含特征列features(Word2Vec生成的20维Vector)和聚类预测列prediction,直接用withColumn添加新列即可:

# 添加名为cluster_distance的新列
df_with_distance = df.withColumn("cluster_distance", distance_udf(df["features"], df["prediction"]))

# 查看结果
df_with_distance.select("features", "prediction", "cluster_distance").show(5, truncate=False)

额外说明

  • 广播变量的作用:因为聚类中心是全局数据,广播后每个worker节点只会加载一次,避免重复传输,大幅提升分布式计算的效率。
  • 距离自定义:如果需要其他距离,比如曼哈顿距离,可以把np.linalg.norm换成np.sum(np.abs(feature_array - target_center))即可。
  • 类型匹配:确保你的features列是Spark的Vector类型(Word2Vec的输出默认就是这个类型),prediction列是整数类型(KMeans预测的结果默认也是整数)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:07:26