PySpark CrossValidator调优ALS模型报错:Nothing added to summarizer
ALS矩阵分解模型CrossValidator调参报错排查
问题背景
- 正在开展ALS矩阵分解模型参数调优工作,采用
pyspark.ml.tuning.CrossValidator遍历参数网格选取最优模型,但运行CrossValidator调参流程时持续触发报错 - 初步排查怀疑报错由ALS模型对验证集中存在、但未收录于训练集的用户做推理触发;查阅Spark官方JIRA记录SPARK-35271可知,CrossValidator在多节点并行执行任务时,若子任务出现拟合错误会中断主线程抛出异常,但暂未找到对应解决方案
- 手动通过for循环实现网格搜索逻辑时无任何报错,无法定位错误仅在CrossValidator场景下触发的原因
- 需要确认ALS超参数
coldStartStrategy = "drop"是否仅对物品维度生效,该配置已知可规避训练集未覆盖数据的冷启动报错
复现代码
from pyspark.ml.recommendation import ALS als = ALS( userCol = "user" # 用户列 , itemCol = "item" # 物品列 , ratingCol = 'read' # 交互日志(read取值为0或1) , coldStartStrategy = "drop" # 冷启动处理策略 , implicitPrefs = True # 隐式反馈场景 , nonnegative = True # 非负矩阵分解 ) from pyspark.ml.tuning import TrainValidationSplit, CrossValidator, ParamGridBuilder from pyspark.ml.evaluation import BinaryClassificationEvaluator from pyspark.ml.evaluation import RegressionEvaluator ParamMaps = (ParamGridBuilder() .addGrid(als.rank, [20, 40, 80, 120]) .addGrid(als.maxIter, [10, 30, 50]) .addGrid(als.regParam, [0.001, 0.01, 0.05, 0.1]) ).build() evaluator = RegressionEvaluator( metricName = "rmse" , labelCol = "read" , predictionCol = "prediction") cv = CrossValidator( estimator = als , estimatorParamMaps = ParamMaps , evaluator = evaluator , parallelism = 5 , numFolds = 5 , seed = 42 ) cvmodel = cv.fit(train)
报错日志
MLlib will automatically track trials in MLflow. After your tuning fit() call has completed, view the MLflow UI to see logged runs. /databricks/spark/python/pyspark/ml/util.py:90: UserWarning: CrossValidator_3986dadc628f fit call failed but some spark jobs may still running for unfinished trials. To address this issue, you should enable pyspark pinned thread mode. .format(uid)) IllegalArgumentException: requirement failed: Nothing has been added to this summarizer. --------------------------------------------------------------------------- IllegalArgumentException Traceback (most recent call last) <command-3029966> in <module> 10 if one of the trial tasks failed, the CrossValidator/TrainValidationSplit fit will raise error and break the main thread, but other backgroud threads running other trial tasks will continue to run, and trial tasks which are pending to run in thread queue will also continue to launch. 11 ''' ---> 12 cvmodel = cv.fit(train) /databricks/spark/python/pyspark/ml/base.py in fit(self, dataset, params) 127 return self.copy(params)._fit(dataset) 128 else: ---> 129 return self._fit(dataset) 130 else: 131 raise ValueError("Params must be either a dict or a list/tuple of param maps, ") /databricks/spark/python/pyspark/ml/tuning.py in _fit(self, dataset) 458 subModels[i][j] = subModel 459 ---> 460 _cancel_on_failure(dataset._sc, self.uid, sub_task_failed, calculate_metrics) 461 validation.unpersist() 462 train.unpersist() /databricks/spark/python/pyspark/ml/util.py in _cancel_on_failure(sc, uid, sub_task_failed, f) 89 "issue, you should enable pyspark pinned thread mode." 90 .format(uid)) ---> 91 raise e 92 93 old_job_group = sc.getLocalProperty("spark.jobGroup.id") /databricks/spark/python/pyspark/ml/util.py in _cancel_on_failure(sc, uid, sub_task_failed, f) 83 if os.environ.get("PYSPARK_PIN_THREAD", "false").lower() != "true": 84 try: ---> 85 return f() 86 except Exception as e: 87 warnings.warn("{} fit call failed but some spark jobs ") /databricks/spark/python/pyspark/ml/tuning.py in calculate_metrics() 452 return task() 453 ---> 454 for j, metric, subModel in pool.imap_unordered(run_task, tasks): 455 metrics[j] += (metric / nFolds) 456 metrics_all[i][j] = metric /usr/lib/python3.7/multiprocessing/pool.py in next(self, timeout) 746 if success: 747 return value ---> 748 raise value 749 750 __next__ = next # XXX /usr/lib/python3.7/multiprocessing/pool.py in worker(inqueue, outqueue, initializer, initargs, maxtasks, wrap_exception) 119 job, i, func, args, kwds = task 120 try: ---> 121 result = (True, func(*args, **kwds)) 122 except Exception as e: 123 if wrap_exception and func is not _helper_reraises_exception: /databricks/spark/python/pyspark/ml/tuning.py in run_task(task) 450 if sub_task_failed[0]: 451 raise RuntimeError("Terminate this task because one of other task failed.") ---> 452 return task() 453 454 for j, metric, subModel in pool.imap_unordered(run_task, tasks): /databricks/spark/python/pyspark/ml/tuning.py in singleTask() 57 # `MetaAlgorithmReadWrite.getAllNestedStages`, make it return 58 # all nested stages and evaluators ---> 59 metric = eva.evaluate(model.transform(validation, epm[index])) 60 return index, metric, model if collectSubModel else None 61 /databricks/spark/python/pyspark/ml/evaluation.py in evaluate(self, dataset, params) 70 return self.copy(params)._evaluate(dataset) 71 else: ---> 72 return self._evaluate(dataset) 73 else: 74 raise ValueError("Params must be a param map but got %s." % type(params)) /databricks/spark/python/pyspark/ml/evaluation.py in _evaluate(self, dataset) 100 """ 101 self._transfer_params_to_java() ---> 102 return self._java_obj.evaluate(dataset._jdf) 103 104 def isLargerBetter(self): /databricks/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py in __call__(self, *args) 1303 answer = self.gateway_client.send_command(command) 1304 return_value = get_return_value( -> 1305 answer, self.gateway_client, self.target_id, self.name) 1306 1307 for temp_arg in temp_args: /databricks/spark/python/pyspark/sql/utils.py in deco(*a, **kw) 131 # Hide where the exception came from that shows a non-Pythonic 132 # JVM exception message. ---> 133 raise_from(converted) 134 else: 135 raise /databricks/spark/python/pyspark/sql/utils.py in raise_from(e) IllegalArgumentException: requirement failed: Nothing has been added to this summarizer.
问题解答
- 关于
coldStartStrategy="drop"的生效范围:该配置同时对用户、物品两个维度的冷启动ID生效,并非仅针对物品维度。预测时只要遇到训练集中不存在的用户ID或物品ID,对应样本行都会被直接丢弃,不参与后续预测计算。 - 报错根因:CrossValidator做5折拆分时,存在某一折的验证集在经过
coldStartStrategy="drop"过滤冷启动样本后,变成了空数据集,此时原生RegressionEvaluator对空数据集计算RMSE时,没有任何样本输入统计汇总器,就会抛出Nothing has been added to this summarizer异常。手写for循环网格搜索不报错,通常是因为手动拆分的验证集不存在全量样本都是冷启动ID的极端情况,过滤冷启动样本后仍有剩余数据可计算指标。 - 可落地解决方案:
- 数据预处理阶段过滤交互次数过少的低频用户、低频物品,从根源上减少冷启动ID占比,避免K折拆分后某一折验证集的ID完全未出现在对应训练折中的极端情况
- 自定义评估器,在计算指标前先判断输入数据集是否为空,如果为空直接返回最差指标值(比如RMSE返回双精度浮点数最大值),跳过空数据集的指标计算逻辑
- 隐式反馈场景不建议使用RMSE作为评估指标,替换为排序类指标(如NDCG、MAP、AUC)更符合业务逻辑,这类指标本身内置了空预测集的边界处理逻辑,不易触发该类报错
- 临时规避方案:将CrossValidator的
parallelism设为1关闭并行执行,同时配置环境变量PYSPARK_PIN_THREAD=true开启PySpark固定线程模式,可避免子任务异常直接中断主线程,但该方案不能解决空数据集计算的本质问题,仅能拿到部分折的计算结果
内容的提问来源于stack exchange,提问作者오진석
相关产品推荐
相关产品推荐

