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

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的极端情况,过滤冷启动样本后仍有剩余数据可计算指标。
  • 可落地解决方案:
    1. 数据预处理阶段过滤交互次数过少的低频用户、低频物品,从根源上减少冷启动ID占比,避免K折拆分后某一折验证集的ID完全未出现在对应训练折中的极端情况
    2. 自定义评估器,在计算指标前先判断输入数据集是否为空,如果为空直接返回最差指标值(比如RMSE返回双精度浮点数最大值),跳过空数据集的指标计算逻辑
    3. 隐式反馈场景不建议使用RMSE作为评估指标,替换为排序类指标(如NDCG、MAP、AUC)更符合业务逻辑,这类指标本身内置了空预测集的边界处理逻辑,不易触发该类报错
    4. 临时规避方案:将CrossValidator的parallelism设为1关闭并行执行,同时配置环境变量PYSPARK_PIN_THREAD=true开启PySpark固定线程模式,可避免子任务异常直接中断主线程,但该方案不能解决空数据集计算的本质问题,仅能拿到部分折的计算结果

内容的提问来源于stack exchange,提问作者오진석

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 08:09:23