基于PySpark在大规模基因数据集上实现K-Means聚类遇故障
用PySpark对大规模基因数据集执行K-Means聚类的问题
问题描述
需处理维度为(6883, 9995)的基因数据集,已成功加载为PySpark DataFrame,但在特征转换及K-Means拟合环节失败,核心障碍是无法将特征列转换为K-Means所需格式,运行时出现任务失败报错。
复现代码
from pyspark.sql import SparkSession spark = SparkSession.getActiveSession() file_location = "/user/dorwayne/bgc_features_part0001.tsv" file_type = "tsv" # CSV options infer_schema = "false" first_row_is_header = "false" delimiter = "," bgc = spark.read.format("csv").option("inferSchema", "False").option("delimiter", "\\t").option("header","true").load("dbfs:/user/dorwayne/bgc_features_part0001.tsv") assembler = VectorAssembler(inputCols = bgc.columns, outputCol = "features") assembled_data = assembler.transform(bgc) kmeans = KMeans().setK(2).setSeed(1) model = kmeans.fit(assembled_data) predictions = model.transform(assembled_data)
报错信息
Py4JJavaError: An error occurred while calling o8473.fit. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 2 in stage 43.0 failed 4 times, most recent failure: Lost task 2.3 in stage 43.0 (TID 147) (ip-10-131-129-26.ec2.internal executor driver): java.lang.AssertionError: assertion failed
问题分析与解决方法
核心问题点
- Schema未正确推断:读取数据时
inferSchema设为False,导致所有列均为字符串类型,而VectorAssembler仅支持数值类型(int、float等),无法将字符串列转换为特征向量。 - 包含非特征列:直接将
bgc.columns全部作为输入列,若数据集包含样本ID等非数值标识列,会引发类型不兼容问题。 - 缺失/异常值未处理:基因数据集可能存在缺失值或无穷值,导致K-Means算法计算时失败。
修正步骤
- 正确读取数据并推断Schema:读取TSV时开启
inferSchema自动推断数值类型,或手动定义Schema保证列类型正确。 - 筛选数值型特征列:排除非数值列,仅保留用于聚类的数值特征列。
- 处理缺失值:使用
Imputer填充缺失值,避免算法因空值报错。
修正后的完整代码
from pyspark.sql import SparkSession from pyspark.ml.feature import VectorAssembler, Imputer from pyspark.ml.clustering import KMeans # 初始化SparkSession(若未存在则创建) spark = SparkSession.builder.appName("GeneKMeans").getOrCreate() # 读取TSV数据,自动推断Schema bgc = spark.read.format("csv")\ .option("inferSchema", "True")\ .option("delimiter", "\\t")\ .option("header", "true")\ .load("dbfs:/user/dorwayne/bgc_features_part0001.tsv") # 筛选数值型特征列(排除非数值列,比如样本ID列,可根据实际列名调整) numeric_cols = [col for col, dtype in bgc.dtypes if dtype in ("int", "double", "float")] # 处理缺失值 imputer = Imputer(inputCols=numeric_cols, outputCols=[f"{col}_imputed" for col in numeric_cols]) imputed_data = imputer.fit(bgc).transform(bgc) # 组装特征向量 assembler = VectorAssembler( inputCols=[f"{col}_imputed" for col in numeric_cols], outputCol="features" ) assembled_data = assembler.transform(imputed_data) # 训练K-Means模型 kmeans = KMeans().setK(2).setSeed(1) model = kmeans.fit(assembled_data) # 生成聚类预测结果 predictions = model.transform(assembled_data) # 查看结果 predictions.select("features", "prediction").show(5)
内容的提问来源于stack exchange,提问作者Spidey1610
相关产品推荐
相关产品推荐

