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

