Is there a hook for Executor Startup in Spark?

Viewed 1618

So, basically I want multiple tasks running on the same node/executor to read data from a shared memory. For that I need some initialization function that would load the data into the memory before the tasks are started. If Spark provides a hook for an Executor startup, I could put this initialization code in that callback function, with the tasks only running after this startup is completed.

So, my question is, does Spark provides such hooks? If not, with which other method, I can achieve the same?

2 Answers

Spark's solution for "shared data" is using broadcast - where you load the data once in the driver application and Spark serializes it and sends to each of the executors (once). If a task uses that data, Spark will make sure it's there before the task is executed. For example:

object MySparkTransformation {

  def transform(rdd: RDD[String], sc: SparkContext): RDD[Int] = {
    val mySharedData: Map[String, Int] = loadDataOnce()
    val broadcast = sc.broadcast(mySharedData)
    rdd.map(r => broadcast.value(r))
  }
}

Alternatively, if you want to avoid reading the data into driver memory and sending it over to the executors, you can use lazy values in a Scala object to create a value that gets populated once per JVM, which in Spark's case is once per executor. For example:

// must be an object, otherwise will be serialized and sent from driver
object MySharedResource {
  lazy val mySharedData: Map[String, Int] = loadDataOnce()
}

// If you use mySharedData in a Spark transformation, 
// the "local" copy in each executor will be used:
object MySparkTransformation {
  def transform(rdd: RDD[String]): RDD[Int] = {
    // Spark won't include MySharedResource.mySharedData in the 
    // serialized task sent from driver, since it's "static"
    rdd.map(r => MySharedResource.mySharedData(r))
  }
}

In practice, you'll have one copy of mySharedData in each executor.

You don't have to run multiple instances of the app to be able to run multiple tasks (i.e. one app instance, one Spark task). The same SparkSession object can be used by multiple threads to submit Spark tasks in parallel.

So it may work like this:

  • The application starts up and runs an initialization function to load shared data in memory. Say, into a SharedData class object.
  • SparkSession is created
  • A thread pool is created, each thread has access to (SparkSession, SharedData) objects
  • Each thread creates Spark task using shared SparkSession and SharedData objects.
  • Depending on your use case, the application then does one of the following:
    • waits for all tasks to complete and then closes Spark Session
    • waits in a loop for new requests to arrive and creates new Spark tasks as necessary using threads from the thread pool.

SparkContext (sparkSession.sparkContext) is useful when you want to do per-thread things like assigning a task description using setJobDescription or assigning a group to the task using setJobGroup so related tasks can be cancelled simultaneously using cancelJobGroup. You can also tweak priority for the tasks that use the same pool, see https://spark.apache.org/docs/latest/job-scheduling.html#scheduling-within-an-application for details.

Related