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

如何用PySpark的K-means聚类为含分类值的大数据框生成聚类列?

在PySpark中实现带缺失值的分类列K-means聚类

完全可以在PySpark中实现你的需求,以下是针对新手的分步实现方案,包含代码和关键说明:

步骤1:初始化SparkSession

先完成PySpark环境的基础初始化:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, coalesce
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.clustering import KMeans

spark = SparkSession.builder.appName("KMeansSurveyClustering").getOrCreate()

步骤2:处理缺失值

K-means无法直接处理NaN值,针对分类列,最常用的方式是用众数填充(也可根据业务需求选择删除含缺失值的行):

# 替换成你实际的8个目标列名
target_cols = ["column1", "column2", "column3", "column4", "column5", "column6", "column7", "column8"]

# 计算每个目标列的众数
mode_values = {}
for col_name in target_cols:
    mode_val = df.groupBy(col_name).count().orderBy("count", ascending=False).first()[0]
    mode_values[col_name] = mode_val

# 用众数填充各列的NaN值
df_filled = df
for col_name in target_cols:
    df_filled = df_filled.withColumn(col_name, coalesce(col(col_name), mode_values[col_name]))

步骤3:分类标签转数值

K-means基于数值特征计算,需将文本标签转为有序数值(比如strongly agree=2、agree=1、disagree=0):

def map_category(col_name):
    return when(col(col_name) == "strongly agree", 2) \
           .when(col(col_name) == "agree", 1) \
           .when(col(col_name) == "disagree", 0) \
           .alias(f"{col_name}_num")

# 生成8个数值列
df_num = df_filled
for col_name in target_cols:
    df_num = df_num.withColumn(f"{col_name}_num", map_category(col_name))

num_cols = [f"{col}_num" for col in target_cols]

步骤4:组装特征向量

PySpark机器学习模型要求输入为单一特征向量列,用VectorAssembler合并数值列:

assembler = VectorAssembler(inputCols=num_cols, outputCol="features")
df_features = assembler.transform(df_num)

步骤5:训练K-means模型并预测

设置聚类数为8,训练后生成聚类标签:

# 初始化模型,固定随机种子保证结果可复现
kmeans = KMeans(k=8, seed=123, featuresCol="features", predictionCol="cluster_label")
model = kmeans.fit(df_features)

# 生成聚类结果
df_clustered = model.transform(df_features)

步骤6:调整聚类标签为1-8格式

PySpark默认输出的标签是0-7,加1转为1-8:

df_final = df_clustered.withColumn("new_column_required", col("cluster_label") + 1)

# 可选:只保留原始列和聚类标签列
df_final = df_final.select(*target_cols, "new_column_required")

关键提示

  • 若业务允许丢失数据,可直接用df.na.drop(subset=target_cols)替代众数填充
  • 若无需考虑分类的有序性,也可使用StringIndexer转换标签,但有序映射更贴合调研数据的逻辑
  • 聚类效果不佳时,可尝试调整k值或增加maxIter(最大迭代次数)参数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 09:57:38