Dynamic Allocation

Dynamic Allocation and Executor Pods

With dynamic allocation, the driver adds executor pods while tasks are waiting and removes executors that have been idle for executorIdleTimeout. YARN keeps shuffle files in an external shuffle service; Kubernetes 5,150 has none, so Spark 129 must either track which executors hold shuffle data (shuffleTracking.enabled, used here, with a timeout so they can go eventually) or migrate it off with graceful decommissioning.

The same job with dynamic allocation, counting executor podsShell
export KUBECONFIG=/home/dev/v7-l3/k8s/kubeconfig
watch_executors() {                                         # print the count when it changes
  last=-1
  while sleep 2; do
    n=$(kubectl get pods -l spark-role=executor --no-headers 2>/dev/null | grep -c Running)
    [ "$n" != "$last" ] && echo "$(date +%T) running executor pods: $n" && last=$n
  done
}
watch_executors & WATCH=$!
spark-submit --properties-file k8s/booknest-k8s.conf --name booknest-dynamic \
  --conf spark.dynamicAllocation.enabled=true \
  --conf spark.dynamicAllocation.shuffleTracking.enabled=true \
  --conf spark.dynamicAllocation.shuffleTracking.timeout=20s \
  --conf spark.dynamicAllocation.executorIdleTimeout=10s \
  --conf spark.dynamicAllocation.minExecutors=0 --conf spark.dynamicAllocation.maxExecutors=3 \
  local:///booknest/jobs/k8s_genre.py 60 2> logs/k8s-dynamic.log   # then idles 60 s
kill $WATCH
Output
04:21:48 running executor pods: 0
04:22:19 running executor pods: 1
04:22:33 running executor pods: 2
04:22:48 running executor pods: 1
04:22:57 running executor pods: 0

The driver started with no executors, asked for one when the first tasks queued, and for a second as the backlog persisted (the four-task scan never needed the third allowed). After the query finished, the job's 60-second idle period let the idle and shuffle-tracking timeouts expire, and both pods were deleted while the driver was still running. On a shared cluster this is how a nightly Spark job returns capacity to other workloads, and with a cluster autoscaler the freed nodes themselves disappear.