PySpark中如何高效计算每个样本的Silhouette coefficient
PySpark 单样本轮廓系数实现方案
Spark MLlib 原生未提供等价于 sklearn silhouette_samples 的单样本轮廓系数计算接口,可通过分布式矢量化方式高效实现,无需将全量数据拉取到本地计算,避免驱动节点内存溢出。
实现前提
完成流水线训练后,首先提取带预测结果的数据集、聚类模型参数,注意计算距离必须使用KMeans实际训练用的特征列,也就是流水线中的pca_output列,不能用原始特征或标准化后未降维的特征,否则结果和聚类逻辑不匹配:
# 完成流水线训练 model = pipeline.fit(train_data) # 得到带聚类预测标签的数据集 pred_df = model.transform(train_data) # 提取训练好的KMeans模型,拿到聚类中心、簇数等参数 kmeans_model = model.stages[-1] cluster_centers = kmeans_model.clusterCenters() k = kmeans_model.getK()
高效分布式实现步骤
轮廓系数的核心计算公式为:s(i) = (b(i) - a(i)) / max(a(i), b(i)),其中:
- a(i):样本i到同簇所有其他样本的平均距离,衡量簇内凝聚度
- b(i):样本i到所有非所属簇中,平均距离最小的那个簇的平均距离,衡量簇间分离度
1. 广播公共变量
将聚类中心广播到所有计算节点,减少节点间重复数据传输开销:
from pyspark.sql import functions as F from pyspark.sql.types import FloatType import numpy as np import pandas as pd sc = spark.sparkContext # 广播聚类中心到所有Executor bc_centers = sc.broadcast(np.array(cluster_centers))
2. 编写矢量化计算函数
用Pandas UDF(applyInPandas)实现单分区内的批量轮廓系数计算,比逐行Python UDF性能高5~10倍:
def calc_silhouette_sample(pdf: pd.DataFrame) -> pd.DataFrame: centers = bc_centers.value # 提取分区内的特征、标签转为numpy数组加速计算 X = np.array([vec.toArray() for vec in pdf["pca_output"]]) labels = pdf["prediction"].values n_samples = len(X) silhouette_scores = np.zeros(n_samples) for idx in range(n_samples): current_label = labels[idx] current_x = X[idx] # 计算a(i):同簇其他样本的平均距离 same_cluster_mask = labels == current_label same_cluster_cnt = same_cluster_mask.sum() # 簇内仅当前样本1个点时,和sklearn逻辑保持一致记为0 if same_cluster_cnt <= 1: silhouette_scores[idx] = 0.0 continue dist_same = np.linalg.norm(X[same_cluster_mask] - current_x, axis=1) a_i = dist_same.sum() / (same_cluster_cnt - 1) # 计算b(i):最近异簇的平均距离 b_i = np.inf for cluster_id in range(k): if cluster_id == current_label: continue other_cluster_mask = labels == cluster_id other_cluster_cnt = other_cluster_mask.sum() if other_cluster_cnt == 0: continue dist_other = np.linalg.norm(X[other_cluster_mask] - current_x, axis=1) avg_dist = dist_other.mean() if avg_dist < b_i: b_i = avg_dist # 计算当前样本轮廓系数 silhouette_scores[idx] = (b_i - a_i) / max(a_i, b_i) pdf["silhouette_sample"] = silhouette_scores return pdf
3. 重分区后执行计算
按聚类预测标签对数据重分区,保证同个簇的样本尽量分布在同一个分区,避免全量数据笛卡尔积的O(n²)开销:
# 按簇标签重分区,分区数可设为簇数或簇数的整数倍 repartitioned_df = pred_df.repartition(k, "prediction") # 按分区批量计算单样本轮廓系数 result_df = repartitioned_df.groupBy(F.spark_partition_id()).applyInPandas( calc_silhouette_sample, schema=pred_df.schema.add("silhouette_sample", FloatType()) )
执行完成后,result_df的silhouette_sample列就是每个样本对应的轮廓系数,和sklearnsilhouette_samples输出逻辑完全一致。
性能优化提示
- 如果单个簇的样本量极大(单簇超100万样本)单分区内存无法承载,可将大簇拆分为多个分区,先计算每个分块的样本距离和、样本数,再二次聚合得到全局平均距离,避免OOM
- 如果业务允许近似结果,可直接用样本到各簇中心的距离代替到簇内所有样本的平均距离计算a(i)、b(i),计算速度可提升两个数量级,大部分场景下结果误差可接受
- 不要直接对全量数据做笛卡尔积计算两两距离,数据量超过10万时会产生严重的性能问题甚至任务失败
内容的提问来源于stack exchange,提问作者Sreeram TP
相关产品推荐
相关产品推荐

