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
相关产品推荐
相关产品推荐

