microsoft / microsoft/SynapseML
Could not recover from a failed barrier ResultStage
- Dominant language
- Scala
- Stars
- 5.2k
- Forks
- 868
- Avg merge
- 22h 9m
- Merged PRs (30d)
- 45
Description
**Describe the bug**
I am training a LightGBM Classifier with the following dataset sizes and feature
- Train: 2,017,289 samples
- Valid: 200,000 samples
- Test: 200,000 samples
The feature vector size is 316 with boolean values. For each data split, I am having 30-70% for my binary class labels
However, I am getting a `java._io` error
**To Reproduce**
Just run the following snippet with appropriate imports. Train/Valid/Test are spark data frames read just before the following snippet
```
def getR2(predictions, label):
evaluator = (
ComputeModelStatistics()
.setScoredLabelsCol("prediction")
.setLabelCol("label")
.setEvaluationMetric("accuracy"))
result = evaluator.transform(predictions)
return result.select("accuracy").collect()[0][0]
def train_evaluate(train,valid,params):
lgbm = LightGBMClassifier(
labelCol="label",
featuresCol='features',
objective="binary",
# timeout=1200.0,
isUnbalance=True,
useBarrierExecutionMode=True,
baggingSeed=42)
model = lgbm.fit(train.drop('userId'),params)
pred = model.transform(valid)
return getR2(pred, valid.select('label'))
def objective(params):
return -1.0 * train_evaluate(train, valid, params)
best = fmin(objective, space, algo=tpe.suggest, max_evals=200)
space_eval(space, best)
```
**Expected behavior**
Dataset is small & it should run easily with my huge config without any OOM errors
**Info (please complete the following information):**
- MMLSpark Version: mmlspark_2.11:1.0.0-rc3
- Spark Version 2.4.2
- driver.maxResultSize: 53488M
- executor.instances: 15
- spark.driver.memory: 53488M
- Executor memory: 104G
- System: Ubuntu
**Stacktrace**
```
job exception: An error occurred while calling o8456.fit.
: org.apache.spark.SparkException: Job aborted due to stage failure: Could not recover from a failed barrier ResultStage. Most recent failure reason: Stage failed because barrier task ResultTask(246, 14) finished unsuccessfully.
java.lang.UnsatisfiedLinkError: Could not load the native libraries because we encountered the following problems: no _lightgbm in java.library.path and No space left on device
at com.microsoft.ml.spark.core.env.NativeLoader.loadLibraryByName(NativeLoader.java:62)
at com.microsoft.ml.spark.lightgbm.LightGBMUtils$.initializeNativeLibrary(LightGBMUtils.scala:46)
at com.microsoft.ml.spark.lightgbm.TrainUtils$$anonfun$16.apply(TrainUtils.scala:547)
at com.microsoft.ml.spark.lightgbm.TrainUtils$$anonfun$16.apply(TrainUtils.scala:544)
at com.microsoft.ml.spark.core.env.StreamUtilities$.using(StreamUtilities.scala:29)
at com.microsoft.ml.spark.lightgbm.TrainUtils$.trainLightGBM(TrainUtils.scala:543)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$$anonfun$7.apply(LightGBMBase.scala:225)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$$anonfun$7.apply(LightGBMBase.scala:225)
at org.apache.spark.rdd.RDDBarrier$$anonfun$mapPartitions$1$$anonfun$apply$1.apply(RDDBarrier.scala:51)
at org.apache.spark.rdd.RDDBarrier$$anonfun$mapPartitions$1$$anonfun$apply$1.apply(RDDBarrier.scala:51)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
at org.apache.spark.scheduler.Task.run(Task.scala:121)
at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1889)
at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1877)
at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1876)
at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1876)
at org.apache.spark.scheduler.DAGScheduler.handleTaskCompletion(DAGScheduler.scala:1665)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2107)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2059)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2048)
at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:737)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2061)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2158)
at org.apache.spark.rdd.RDD$$anonfun$reduce$1.apply(RDD.scala:1035)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
at org.apache.spark.rdd.RDD.reduce(RDD.scala:1017)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$class.innerTrain(LightGBMBase.scala:228)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.innerTrain(LightGBMClassifier.scala:24)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$class.train(LightGBMBase.scala:48)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.train(LightGBMClassifier.scala:24)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.train(LightGBMClassifier.scala:24)
at org.apache.spark.ml.Predictor.fit(Predictor.scala:118)
at sun.reflect.GeneratedMethodAccessor111.invoke(Unknown Source)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
at py4j.Gateway.invoke(Gateway.java:282)
at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79)
at py4j.GatewayConnection.run(GatewayConnection.java:238)
at java.lang.Thread.run(Thread.java:748)
15%|█▌ | 30/200 [1:16:23<7:12:55, 152.80s/trial, best loss: -0.7321714624740433]
---------------------------------------------------------------------------
Py4JJavaError Traceback (most recent call last)
in
24 return -1.0 * train_evaluate(train, valid, params)
25
---> 26 best = fmin(objective, space, algo=tpe.suggest, max_evals=200)
27 space_eval(space, best)
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/hyperopt/fmin.py in fmin(fn, space, algo, max_evals, timeout, loss_threshold, trials, rstate, allow_trials_fmin, pass_expr_memo_ctrl, catch_eval_exceptions, verbose, return_argmin, points_to_evaluate, max_queue_len, show_progressbar, early_stop_fn, trials_save_file)
551
552 # next line is where the fmin is actually executed
--> 553 rval.exhaust()
554
555 if return_argmin:
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/hyperopt/fmin.py in exhaust(self)
354 def exhaust(self):
355 n_done = len(self.trials)
--> 356 self.run(self.max_evals - n_done, block_until_done=self.asynchronous)
357 self.trials.refresh()
358 return self
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/hyperopt/fmin.py in run(self, N, block_until_done)
290 else:
291 # -- loop over trials and do the jobs directly
--> 292 self.serial_evaluate()
293
294 self.trials.refresh()
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/hyperopt/fmin.py in serial_evaluate(self, N)
168 ctrl = base.Ctrl(self.trials, current_trial=trial)
169 try:
--> 170 result = self.domain.evaluate(spec, ctrl)
171 except Exception as e:
172 logger.error("job exception: %s" % str(e))
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/hyperopt/base.py in evaluate(self, config, ctrl, attach_attachments)
905 print_node_on_error=self.rec_eval_print_node_on_error,
906 )
--> 907 rval = self.fn(pyll_rval)
908
909 if isinstance(rval, (float, int, np.number)):
in objective(params)
22
23 def objective(params):
---> 24 return -1.0 * train_evaluate(train, valid, params)
25
26 best = fmin(objective, space, algo=tpe.suggest, max_evals=200)
in train_evaluate(train, valid, params)
17 useBarrierExecutionMode=True,
18 baggingSeed=42)
---> 19 model = lgbm.fit(train.drop('userId'),params)
20 pred = model.transform(valid)
21 return getR2(pred, valid.select('label'))
/usr/local/spark/python/pyspark/ml/base.py in fit(self, dataset, params)
128 elif isinstance(params, dict):
129 if params:
--> 130 return self.copy(params)._fit(dataset)
131 else:
132 return self._fit(dataset)
/usr/local/spark/python/pyspark/ml/wrapper.py in _fit(self, dataset)
293
294 def _fit(self, dataset):
--> 295 java_model = self._fit_java(dataset)
296 model = self._create_model(java_model)
297 return self._copyValues(model)
/usr/local/spark/python/pyspark/ml/wrapper.py in _fit_java(self, dataset)
290 """
291 self._transfer_params_to_java()
--> 292 return self._java_obj.fit(dataset._jdf)
293
294 def _fit(self, dataset):
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/py4j/java_gateway.py in __call__(self, *args)
1255 answer = self.gateway_client.send_command(command)
1256 return_value = get_return_value(
-> 1257 answer, self.gateway_client, self.target_id, self.name)
1258
1259 for temp_arg in temp_args:
/usr/local/spark/python/pyspark/sql/utils.py in deco(*a, **kw)
61 def deco(*a, **kw):
62 try:
---> 63 return f(*a, **kw)
64 except py4j.protocol.Py4JJavaError as e:
65 s = e.java_exception.toString()
~/miniconda3/envs/reco_pyspark/lib/python3.6/site-packages/py4j/protocol.py in get_return_value(answer, gateway_client, target_id, name)
326 raise Py4JJavaError(
327 "An error occurred while calling {0}{1}{2}.\n".
--> 328 format(target_id, ".", name), value)
329 else:
330 raise Py4JError(
Py4JJavaError: An error occurred while calling o8456.fit.
: org.apache.spark.SparkException: Job aborted due to stage failure: Could not recover from a failed barrier ResultStage. Most recent failure reason: Stage failed because barrier task ResultTask(246, 14) finished unsuccessfully.
java.lang.UnsatisfiedLinkError: Could not load the native libraries because we encountered the following problems: no _lightgbm in java.library.path and No space left on device
at com.microsoft.ml.spark.core.env.NativeLoader.loadLibraryByName(NativeLoader.java:62)
at com.microsoft.ml.spark.lightgbm.LightGBMUtils$.initializeNativeLibrary(LightGBMUtils.scala:46)
at com.microsoft.ml.spark.lightgbm.TrainUtils$$anonfun$16.apply(TrainUtils.scala:547)
at com.microsoft.ml.spark.lightgbm.TrainUtils$$anonfun$16.apply(TrainUtils.scala:544)
at com.microsoft.ml.spark.core.env.StreamUtilities$.using(StreamUtilities.scala:29)
at com.microsoft.ml.spark.lightgbm.TrainUtils$.trainLightGBM(TrainUtils.scala:543)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$$anonfun$7.apply(LightGBMBase.scala:225)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$$anonfun$7.apply(LightGBMBase.scala:225)
at org.apache.spark.rdd.RDDBarrier$$anonfun$mapPartitions$1$$anonfun$apply$1.apply(RDDBarrier.scala:51)
at org.apache.spark.rdd.RDDBarrier$$anonfun$mapPartitions$1$$anonfun$apply$1.apply(RDDBarrier.scala:51)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
at org.apache.spark.scheduler.Task.run(Task.scala:121)
at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
at org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1889)
at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1877)
at org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1876)
at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1876)
at org.apache.spark.scheduler.DAGScheduler.handleTaskCompletion(DAGScheduler.scala:1665)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2107)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2059)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2048)
at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:737)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2061)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2158)
at org.apache.spark.rdd.RDD$$anonfun$reduce$1.apply(RDD.scala:1035)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
at org.apache.spark.rdd.RDD.reduce(RDD.scala:1017)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$class.innerTrain(LightGBMBase.scala:228)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.innerTrain(LightGBMClassifier.scala:24)
at com.microsoft.ml.spark.lightgbm.LightGBMBase$class.train(LightGBMBase.scala:48)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.train(LightGBMClassifier.scala:24)
at com.microsoft.ml.spark.lightgbm.LightGBMClassifier.train(LightGBMClassifier.scala:24)
at org.apache.spark.ml.Predictor.fit(Predictor.scala:118)
at sun.reflect.GeneratedMethodAccessor111.invoke(Unknown Source)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
at py4j.Gateway.invoke(Gateway.java:282)
at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79)
at py4j.GatewayConnection.run(GatewayConnection.java:238)
at java.lang.Thread.run(Thread.java:748)
```
Contributor guide
Research direction
Start with NativeLoader.scala and LightGBMUtils.scala, then trace the native-library initialization from TrainUtils.scala and LightGBMBase.scala. Reproduce the failure using the reported Spark and MMLSpark versions with useBarrierExecutionMode enabled, focusing on the reported No space left on device condition. Done means training no longer aborts with the native-library loading error, with coverage for the failure path if the repository has relevant tests.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java, python, scala, spark
- Domain
- distributed-systems, machine-learning
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 35/100