使用PySpark训练Decision Tree时遭遇label歧义错误求助
问题分析与解决
错误原因
报错AnalysisException: Reference 'label' is ambiguous, could be: label, label的核心原因是DataFrame中存在重复的label列。从你打印的schema能看到,数据集中出现了两个label字段:
|-- label: double (nullable = false) |-- features: vector (nullable = true) ... |-- label: double (nullable = false)
这是因为你执行列选择时,错误地将已存在的label和features重复添加:
selectedCols = ['label', 'features'] + cols df = df.select(selectedCols)
其中cols是df.columns,已经包含了Pipeline生成的label列,重复添加导致字段冲突。
另外代码还有一处逻辑错误:在执行pipeline.fit(trainingData)时,trainingData还未定义,应该先拆分原始数据,再用Pipeline处理训练集。
解决步骤
- 移除重复列的选择逻辑:训练模型只需要
label和features列,直接选择这两个字段即可,避免重复添加。 - 调整数据拆分顺序:先将原始数据集拆分为训练集和测试集,再用Pipeline处理训练数据。
修正后的完整代码
from pyspark.ml import Pipeline from pyspark.ml.classification import DecisionTreeClassifier from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder.appName('ML_algo').getOrCreate() # 加载CSV文件到DataFrame df = spark.read.csv('C:/Users/johnc/Downloads/hospital.csv', header=True, inferSchema=True) # 分类列列表 categoricalColumns = ['ethnicity', 'gender', 'apache_3j_bodysystem'] # 构建Pipeline阶段 stages = [] # 目标变量的StringIndexer label_stringIdx = StringIndexer(inputCol='isHospitalDeath', outputCol='label') stages.append(label_stringIdx) # 处理每个分类列:StringIndexer + OneHotEncoder for categoricalCol in categoricalColumns: stringIndexer = StringIndexer(inputCol=categoricalCol, outputCol=f"{categoricalCol}Index") encoder = OneHotEncoder(inputCols=[stringIndexer.getOutputCol()], outputCols=f"{categoricalCol}classVec") stages.extend([stringIndexer, encoder]) # 数值列列表 numericCols = ['age', 'bmi', 'gcs_eyes_apache', 'gcs_verbal_apache', 'heart_rate_apache', 'intubated_apache', 'resprate_apache', 'temp_apache', 'ventilated_apache', 'd1_mbp_min', 'd1_spo2_min', 'd1_sysbp_min', 'd1_temp_min', 'h1_diasbp_min', 'h1_mbp_min', 'h1_resprate_max', 'h1_sysbp_min', 'apache_4a_hospital_death_prob'] # 特征组装 assemblerInputs = [f"{c}classVec" for c in categoricalColumns] + numericCols assembler = VectorAssembler(inputCols=assemblerInputs, outputCol='features') stages.append(assembler) # 创建Pipeline pipeline = Pipeline(stages=stages) # 先拆分原始数据为训练集和测试集 (trainingData, testData) = df.randomSplit([0.7, 0.3], seed=100) # 拟合Pipeline到训练集 pipelineModel = pipeline.fit(trainingData) # 转换训练集 processedTrainingData = pipelineModel.transform(trainingData) # 只选择训练模型需要的列 processedTrainingData = processedTrainingData.select('label', 'features') processedTrainingData.printSchema() # 训练决策树模型 dt = DecisionTreeClassifier(featuresCol='features', labelCol='label', maxDepth=3) dtModel = dt.fit(processedTrainingData)
内容的提问来源于stack exchange,提问作者leonardo splinter
相关产品推荐
相关产品推荐

