PySpark使用IndexToString函数时遇ClassCastException异常求助
问题分析与解决
错误原因
出现java.lang.ClassCastException: UnresolvedAttribute$ cannot be cast to NominalAttribute的核心原因有三点:
- 未正确获取StringIndexer的标签映射:代码中使用的
customerIndexer和productIndexer变量从未定义——你通过列表推导式批量创建StringIndexer并放入Pipeline,但未保存单个拟合后的StringIndexerModel实例,导致IndexToString无法获取正确的原始标签集合。 - 变量名错误:
(training, test) = transformedDF.randomSplit([0.8, 0.2])里的transformedDF应为之前生成的transformed,属于未定义变量的低级错误。 - 推荐结果列缺少属性元数据:从
recommendations拆分出的customer_id_index和product_id_index列,没有携带StringIndexer生成的NominalAttribute元数据,直接用IndexToString转换时无法识别属性类型。
修正方案
步骤1:保存拟合后的PipelineModel
在拟合Pipeline时保存模型实例,方便后续提取StringIndexer的标签映射:
indexer = [StringIndexer(inputCol=column, outputCol=column+"_index") for column in ['customer_id', 'product_id']] pipeline = Pipeline(stages=indexer) # 保存拟合后的模型实例 pipeline_model = pipeline.fit(df2) transformed = pipeline_model.transform(df2)
步骤2:修正变量名错误
将未定义的transformedDF改为正确的transformed:
(training, test) = transformed.randomSplit([0.8, 0.2])
步骤3:正确创建IndexToString转换器
从PipelineModel中提取每个拟合后的StringIndexer实例,用其标签集合创建转换器:
# 从PipelineModel中获取对应阶段的StringIndexerModel customer_indexer_model = pipeline_model.stages[0] product_indexer_model = pipeline_model.stages[1] # 显式传入标签集合创建转换器 customerConverter = IndexToString( inputCol="customer_id_index", outputCol="customer_id", labels=customer_indexer_model.labels ) productConverter = IndexToString( inputCol="product_id_index", outputCol="product_id", labels=product_indexer_model.labels )
步骤4:修正转换逻辑
recs数据仅包含编码后的索引列,无需再用transformed数据拟合Pipeline,直接转换即可:
results = Pipeline(stages=[customerConverter, productConverter]).transform(recs)
完整修正后代码片段
# 生成索引列并保存模型 indexer = [StringIndexer(inputCol=column, outputCol=column+"_index") for column in ['customer_id', 'product_id']] pipeline = Pipeline(stages=indexer) pipeline_model = pipeline.fit(df2) transformed = pipeline_model.transform(df2) # 拆分训练测试集 (training, test) = transformed.randomSplit([0.8, 0.2]) # 训练ALS模型 als= ALS(maxIter=5, rank=4, regParam=0.01,userCol='customer_id_index', itemCol='product_id_index', \ ratingCol='star_rating', coldStartStrategy='drop',\ nonnegative=True) model=als.fit(training) # 模型评估 predictions=model.transform(test) evaluator=RegressionEvaluator(metricName='rmse',labelCol='star_rating',predictionCol='prediction') rmse=evaluator.evaluate(predictions) print(rmse) # 生成用户推荐并拆分数据 recommendations=model.recommendForAllUsers(10) recs=recommendations.withColumn("itemAndRating",explode(recommendations.recommendations))\ .select("customer_id_index","itemAndRating.*") # 转换回原始ID customer_indexer_model = pipeline_model.stages[0] product_indexer_model = pipeline_model.stages[1] customerConverter=IndexToString(inputCol="customer_id_index",outputCol="customer_id",labels=customer_indexer_model.labels) productConverter=IndexToString(inputCol="product_id_index",outputCol="product_id",labels=product_indexer_model.labels) results=Pipeline(stages=[customerConverter,productConverter]).transform(recs)
内容的提问来源于stack exchange,提问作者Nick
相关产品推荐
相关产品推荐

