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

Databricks中PySpark随机森林任务阶段失败问题求助

问题:Databricks PySpark随机森林任务阶段失败排查

在Databricks标准默认集群运行PySpark,导入700万行数据表,此前操作正常,但运行随机森林分类模型时出现任务阶段失败,怀疑数据溢出。已尝试调整PySpark配置、增加分区、重启集群,均无法解决。

引发失败的模型代码

from pyspark.sql import functions as F 
import pyspark.sql.types as t

# 新增avatar_index列,转换为Double类型
user_sessions_df = user_sessions_df.withColumn("avatar_index", F.col("avatar").cast("Double"))

# 增加分区数
user_sessions_df = user_sessions_df.repartition(131072)

# 删除features列
user_sessions_df = user_sessions_df.drop("features")

# 特征向量组装
vector_assembler = VectorAssembler(inputCols=['player_name', 'scene_id', 'date', 'year', 'month', 'day', 'hour', 'avatar_index'],outputCol="features")

# 将特征转换为数值类型
user_sessions_df = user_sessions_df.withColumn("player_name", F.when(condition=((user_sessions_df["player_name"] != "") & (user_sessions_df["player_name"].isNotNull()) & (user_sessions_df["player_name"] != None)),value=F.col("player_name").cast("Double")).otherwise(-1))
user_sessions_df = user_sessions_df.withColumn("scene_id", F.when(condition=((user_sessions_df["scene_id"] != "") & (user_sessions_df["scene_id"].isNotNull()) & (user_sessions_df["scene_id"] != None)),value=F.col("scene_id").cast("Double")).otherwise(-1))
user_sessions_df = user_sessions_df.withColumn("date", F.when(condition=((user_sessions_df["date"] != "") & (user_sessions_df["date"].isNotNull()) & (user_sessions_df["date"] != None)),value=F.col("date").cast("Double")).otherwise(-1))  

# 应用特征向量组装
user_sessions_df = vector_assembler.transform(user_sessions_df)

# 拆分训练集和测试集
(train_df, test_df) = user_sessions_df.randomSplit([0.7, 0.3])

# 创建随机森林分类器
rf_classifier = RandomForestClassifier(labelCol="session_count", featuresCol="features", numTrees=128)

# 训练模型并评估
model = rf_classifier.fit(train_df) 
predictions = model.transform(test_df) 
evaluator = MulticlassClassificationEvaluator(labelCol="session_count", predictionCol="prediction", metricName="accuracy")

accuracy = evaluator.evaluate(predictions) 
print("Accuracy = {}".format(accuracy))

当前PySpark配置

from pyspark.sql import SQLContext 
from pyspark import SparkContext 
from pyspark import SparkConf

conf = SparkConf()

conf.set("spark.executor.instances", "1024")
conf.set("spark.executor.cores", "1024")

conf.set("spark.executor.memory", "1024g")
conf.set("spark.driver.memory", "2048")

conf.set("spark.driver.maxResultSize", "512g")
conf.set("spark.executor.maxResultSize", "512")

conf.set("spark.kryoserializer.buffer.max.mb", "131072")
conf.set("spark.sql.shuffle.partitions", "262144")

# 调整任务调度器核心数(默认1,改为2或4)
conf.set("spark.scheduler.mode", "FAIR")
conf.set("spark.scheduler.allocation.file", "fairscheduler.xml")

报错信息

org.apache.spark.SparkException: Job aborted due to stage failure: Task 3 in stage 82.0 failed 1 times, most recent failure: Lost task
3.0 in stage 82.0 (TID 108) 

org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 173.0 failed 1 times, most recent failure: Lost task
0.0 in stage 173.0 (TID 238) (ip-10-156-234-195.ec2.internal executor driver): org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to execute user defined function (StringIndexerModel$$Lambda$7385/614759719: (string) => double)

集群GC日志

2023-02-25T16:14:35.992+0000: [GC (Allocation Failure) [PSYoungGen: 2221216K->128K(4408320K)] 11188558K 8983958K(17633280K), 0.0409329 secs] [Times: user=0.12 sys=0.01, real=0.04 secs]  

2023-02-25T16:14:37.480+0000: [GC (Allocation Failure) [PSYoungGen: 2204800K->160K(4408320K)] 11188630K->8984070K(17633280K), 0.0381262 secs] [Times: user=0.11 sys=0.01, real=0.04 secs]  

2023-02-25T16:14:38.966+0000: [GC (Allocation Failure) [PSYoungGen: 2204832K->128K(4408320K)] 11188742K->8984118K(17633280K), 0.0270294 secs] [Times: user=0.10 sys=0.00, real=0.03 secs]  

2023-02-25T16:14:40.435+0000: [GC (Allocation Failure) [PSYoungGen: 2204800K->65696K(4408320K)] 11188790K->9049742K(17633280K), 0.0335045 secs] [Times: user=0.12 sys=0.00, real=0.04 secs]  

2023-02-25T16:14:41.958+0000: [GC (Allocation Failure) [PSYoungGen: 2270368K->160K(4527616K)] 11254414K->9180894K(17752576K), 0.1504603 secs] [Times: user=0.33 sys=0.15, real=0.16 secs]  

2023-02-25T16:14:43.671+0000: [GC (Allocation Failure) [PSYoungGen: 2324128K->65728K(4408320K)] 11504862K->9246550K(17633280K), 0.0338819 secs] [Times: user=0.13 sys=0.00, real=0.03 secs]  

2023-02-25T16:14:45.308+0000: [GC (Allocation Failure) [PSYoungGen: 2389696K->164000K(4813312K)] 11570518K->9410462K(18038272K), 0.0868608 secs] [Times: user=0.31 sys=0.03, real=0.09 secs]  
2023-02-25T16:14:47.348+0000: [GC (Allocation Failure) [PSYoungGen: 3031200K->131264K(4665856K)] 12277662K->9574398K(17890816K), 0.1293894 secs] [Times: user=0.26 sys=0.13, real=0.13 secs]  

2023-02-25T16:14:54.775+0000: [GC (Allocation Failure) [PSYoungGen: 2998464K->542977K(4977664K)] 12441598K->10117288K(18202624K),0.1547034 secs] [Times: user=0.42 sys=0.09, real=0.16 secs]  

2023-02-25T16:14:55.789+0000: [GC (Allocation Failure) [PSYoungGen: 3852033K->3392K(4943360K)] 13426344K->10248792K(18168320K), 0.3667542 secs] [Times: user=0.53 sys=0.52, real=0.36 secs]

疑问

是否需要额外配置jar包?


解决方案分析

1. 修正严重不合理的Spark配置

你的配置完全脱离实际集群资源逻辑,是核心问题之一:

  • executor规格超标:spark.executor.instances=1024、spark.executor.cores=1024、spark.executor.memory=1024g,没有任何Databricks集群能支持这种规格,直接导致资源分配失败。
  • 内存参数错误:spark.driver.memory=2048缺失单位,默认仅2KB,驱动内存严重不足;spark.executor.maxResultSize=512同样缺失单位,结果内存限制过小触发溢出。
  • 分区数过度设置:spark.sql.shuffle.partitions=262144远超700万数据所需,导致调度开销暴增。

修正示例(根据常规Databricks集群规格调整):

conf.set("spark.executor.instances", "8")
conf.set("spark.executor.cores", "4")
conf.set("spark.executor.memory", "16g")
conf.set("spark.driver.memory", "8g")
conf.set("spark.driver.maxResultSize", "16g")
conf.set("spark.executor.maxResultSize", "16g")
conf.set("spark.sql.shuffle.partitions", "200")

2. 修复数据处理逻辑错误

报错中的FAILED_EXECUTE_UDF是因为字符串转Double失败:

  • player_name、scene_id是字符串类型,直接强制转Double会因非数字字符报错,应使用Spark ML标准的分类特征处理流程:
from pyspark.ml.feature import StringIndexer, OneHotEncoder

# 处理player_name特征
indexer_player = StringIndexer(inputCol="player_name", outputCol="player_name_idx", handleInvalid="keep")
encoder_player = OneHotEncoder(inputCol="player_name_idx", outputCol="player_name_vec")

# 处理scene_id特征
indexer_scene = StringIndexer(inputCol="scene_id", outputCol="scene_id_idx", handleInvalid="keep")
encoder_scene = OneHotEncoder(inputCol="scene_id_idx", outputCol="scene_id_vec")

# 重新定义特征向量组装列
vector_assembler = VectorAssembler(
    inputCols=['player_name_vec', 'scene_id_vec', 'date', 'year', 'month', 'day', 'hour', 'avatar_index'],
    outputCol="features"
)
  • date若为日期字符串,不要直接转Double,建议用to_date解析后提取特征(如星期、季度),或用unix_timestamp转成时间戳数值。
  • 空值判断冗余:isNotNull()已包含!= None,无需重复判断;需注意字符串空白值(如空格)的处理。

3. 调整分区数量

repartition(131072)会将700万数据拆分为13万多个分区,每个分区仅约50条数据,导致任务调度开销暴增、GC频繁。建议设置为集群总核心数的2-3倍即可。

4. Jar包配置说明

Databricks默认集群已预装Spark ML所需的所有jar包,无需额外配置,你的问题根源不在jar包。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 06:29:56