When running "local-cluster" model in Apache Spark, how to prevent executor from dissociating prematurely?

Viewed 85

I have a Spark application that should be tested in both local mode & local-cluster mode, using scalatest.

The local-cluster mode is submitted using this method:

How to scala-test a Spark program under "local-cluster" mode?

The test run successfully, but when terminating the test I got the following error in the log:


22/05/16 17:45:25 ERROR TaskSchedulerImpl: Lost executor 0 on 172.16.224.18: Remote RPC client disassociated. Likely due to containers exceeding thresholds, or network issues. Check driver logs for WARN messages.
22/05/16 17:45:25 ERROR Worker: Failed to launch executor app-20220516174449-0000/2 for Test.
java.lang.IllegalStateException: Shutdown hooks cannot be modified during shutdown.
    at org.apache.spark.util.SparkShutdownHookManager.add(ShutdownHookManager.scala:195)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:153)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:142)
    at org.apache.spark.deploy.worker.ExecutorRunner.start(ExecutorRunner.scala:77)
    at org.apache.spark.deploy.worker.Worker$$anonfun$receive$1.applyOrElse(Worker.scala:547)
    at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:117)
    at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:215)
    at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:102)
    at org.apache.spark.rpc.netty.Dispatcher$MessageLoop.run(Dispatcher.scala:221)
    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)
22/05/16 17:45:25 ERROR Worker: Failed to launch executor app-20220516174449-0000/3 for Test.
java.lang.IllegalStateException: Shutdown hooks cannot be modified during shutdown.
    at org.apache.spark.util.SparkShutdownHookManager.add(ShutdownHookManager.scala:195)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:153)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:142)
    at org.apache.spark.deploy.worker.ExecutorRunner.start(ExecutorRunner.scala:77)
    at org.apache.spark.deploy.worker.Worker$$anonfun$receive$1.applyOrElse(Worker.scala:547)
    at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:117)
    at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:215)
    at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:102)
    at org.apache.spark.rpc.netty.Dispatcher$MessageLoop.run(Dispatcher.scala:221)
    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)
22/05/16 17:45:25 ERROR Worker: Failed to launch executor app-20220516174449-0000/4 for Test.
java.lang.IllegalStateException: Shutdown hooks cannot be modified during shutdown.
    at org.apache.spark.util.SparkShutdownHookManager.add(ShutdownHookManager.scala:195)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:153)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:142)
    at org.apache.spark.deploy.worker.ExecutorRunner.start(ExecutorRunner.scala:77)
    at org.apache.spark.deploy.worker.Worker$$anonfun$receive$1.applyOrElse(Worker.scala:547)
    at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:117)
    at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:215)
    at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:102)
    at org.apache.spark.rpc.netty.Dispatcher$MessageLoop.run(Dispatcher.scala:221)
    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)
22/05/16 17:45:25 ERROR Worker: Failed to launch executor app-20220516174449-0000/5 for Test.
java.lang.IllegalStateException: Shutdown hooks cannot be modified during shutdown.
    at org.apache.spark.util.SparkShutdownHookManager.add(ShutdownHookManager.scala:195)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:153)
    at org.apache.spark.util.ShutdownHookManager$.addShutdownHook(ShutdownHookManager.scala:142)
    at org.apache.spark.deploy.worker.ExecutorRunner.start(ExecutorRunner.scala:77)
    at org.apache.spark.deploy.worker.Worker$$anonfun$receive$1.applyOrElse(Worker.scala:547)
    at org.apache.spark.rpc.netty.Inbox.$anonfun$process$1(Inbox.scala:117)
    at org.apache.spark.rpc.netty.Inbox.safelyCall(Inbox.scala:215)
    at org.apache.spark.rpc.netty.Inbox.process(Inbox.scala:102)
    at org.apache.spark.rpc.netty.Dispatcher$MessageLoop.run(Dis

...

It turns out executor 0 was dropped before the SparkContext is stopped, this triggered a violent self-healing reaction from Spark master that tries to repeatedly launch new executors to compensate for the loss. How do I prevent this from happening?

1 Answers

Spark attempts to recover from failed tasks by attempting to run them again. What you can do to avoid this is to set some properties to 1 in

  • spark.task.maxFailures (default is 4)
  • spark.stage.maxConsecutiveAttempts (default is 4)

These properties can be set in $SPARK_HOME/conf/spark-defaults.conf or given as options to spark-submit:

spark-submit --conf spark.task.maxFailures=1 --conf spark.stage.maxConsecutiveAttempts=1

or in the Spark context/session configuration before starting the session.

EDIT:

It looks like your executors are lost due to insufficient memory. You could try to increase:

  • spark.executor.memory
  • spark.executor.memoryOverhead
  • spark.memory.offHeap.size with (spark.memory.offHeap.enabled=true)

(see Spark configuration)

The maximum memory size of container to running executor is determined by the sum of spark.executor.memoryOverhead, spark.executor.memory, spark.memory.offHeap.size and spark.executor.pyspark.memory.

Related