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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 12:54:55