Menu

Apache Spark course · Lesson 9 of 9

Deploying and Monitoring Spark: Kubernetes and the History Server

Run Spark on Kubernetes: images, service accounts, pod resources and dependencies, then keep every job inspectable with event logs and the History Server.

  • Advanced
  • 14 min read
  • Updated Oct 2026
On this page
  1. Spark on Kubernetes
  2. What it is
  3. How a submission works
  4. A cluster-mode submission
  5. The configuration you need to know
  6. Images
  7. RBAC: the service account
  8. Dependencies
  9. Scaling, shuffle and spot nodes
  10. spark-submit or an operator
  11. Debugging on Kubernetes
  12. Pitfalls
  13. In interviews
  14. Monitoring with the Spark History Server
  15. What it is
  16. Turning on event logs
  17. Seeing it locally
  18. Running a History Server for a platform
  19. What to look for when reviewing a finished application
  20. Pitfalls
  21. In interviews
  22. Practice questions
  23. Key takeaways

Writing a good Spark job is half the work; the other half is running it reliably and being able to explain afterwards why it was slow or failed. This lesson covers the two operational topics interviewers most often probe for modern platforms: running Spark natively on Kubernetes, and keeping a record of every application with event logs and the Spark History Server.

Spark on Kubernetes

What it is

Since Spark 2.3 (generally available since 3.1), Kubernetes can act as Spark’s cluster manager. There is no Spark master or YARN: spark-submit talks to the Kubernetes API server, which runs the driver and executors as pods built from a container image. It suits teams that already run services on Kubernetes and want Spark jobs to share the same cluster, tooling, autoscaling and security model.

How a submission works

  1. spark-submit --master k8s://https://<api-server>:<port> --deploy-mode cluster creates a driver pod (plus a ConfigMap with the Spark configuration and a headless service so executors can reach the driver).
  2. The driver, using its service account, asks the API server to create executor pods.
  3. Executors connect back to the driver and run tasks as on any other cluster manager.
  4. When the application finishes, executor pods are deleted (spark.kubernetes.executor.deleteOnTermination, default true). The driver pod stays in Completed state so you can read its logs, until you or a cleanup policy delete it.

In client mode the driver runs where you launched it (for example a notebook pod), and you must make it reachable from executors, typically with a headless service and spark.driver.host.

A cluster-mode submission

spark-submit \
  --master k8s://https://k8s-api.example.internal:6443 \
  --deploy-mode cluster \
  --name daily-sales \
  --conf spark.kubernetes.namespace=data-jobs \
  --conf spark.kubernetes.container.image=registry.example.internal/data/spark-py:4.2.0 \
  --conf spark.kubernetes.container.image.pullPolicy=IfNotPresent \
  --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
  --conf spark.executor.instances=10 \
  --conf spark.executor.cores=4 \
  --conf spark.executor.memory=12g \
  --conf spark.kubernetes.executor.request.cores=3.5 \
  --conf spark.kubernetes.executor.limit.cores=4 \
  --conf spark.driver.memory=4g \
  --conf spark.eventLog.enabled=true \
  --conf spark.eventLog.dir=s3a://data-platform-logs/spark-events/ \
  local:///opt/app/jobs/daily_sales.py

local:// means the file is already inside the image. Other dependency options are covered below.

The configuration you need to know

Setting Default (4.2) Purpose
spark.kubernetes.namespace default Namespace for driver and executor pods
spark.kubernetes.container.image none (required) Image for driver and executors; spark.kubernetes.driver.container.image and spark.kubernetes.executor.container.image override it per role
spark.kubernetes.container.image.pullPolicy IfNotPresent Image pull policy
spark.kubernetes.container.image.pullSecrets empty Secrets for private registries
spark.kubernetes.authenticate.driver.serviceAccountName not set Service account the driver uses to create executor pods
spark.kubernetes.driver.request.cores / spark.kubernetes.executor.request.cores not set CPU request per pod; can be fractional and differ from spark.executor.cores
spark.kubernetes.driver.limit.cores / spark.kubernetes.executor.limit.cores not set CPU limit per pod
spark.kubernetes.driver.podTemplateFile / spark.kubernetes.executor.podTemplateFile not set Pod template for anything without a dedicated setting (node selectors, tolerations, sidecars, volumes)
spark.kubernetes.file.upload.path not set Remote location (for example S3) where spark-submit uploads local dependencies in cluster mode
spark.kubernetes.submission.waitAppCompletion true Whether spark-submit waits for the application to finish
spark.kubernetes.executor.deleteOnTermination true Delete executor pods when the application ends
spark.executor.memoryOverheadFactor 0.10 (0.40 for non-JVM jobs per the docs) Pod memory = executor memory + overhead

Pod memory is spark.executor.memory plus overhead, and it is set as both the request and the limit, so a pod that exceeds it is killed by Kubernetes (OOMKilled). The memory and tuning lesson explains the calculation.

Images

Spark ships bin/docker-image-tool.sh to build JVM, Python and R images from a Spark distribution, and the Apache Spark project publishes official images. Most teams build their own image on top, adding:

  • connector JARs (for example hadoop-aws for S3, a JDBC driver, Delta Lake or Iceberg),
  • Python packages the job needs, pinned to versions matching the driver,
  • the job code itself, referenced as local:///....

Keep the image Spark version identical to the version your code was tested on; mixing client and image versions causes serialization and protocol errors.

RBAC: the service account

The driver needs permission to create, watch and delete pods, services and ConfigMaps in its namespace. A minimal setup:

apiVersion: v1
kind: ServiceAccount
metadata:
  name: spark
  namespace: data-jobs
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: spark-driver
  namespace: data-jobs
rules:
  - apiGroups: [""]
    resources: ["pods", "services", "configmaps", "persistentvolumeclaims"]
    verbs: ["create", "get", "list", "watch", "delete", "deletecollection", "patch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: spark-driver
  namespace: data-jobs
subjects:
  - kind: ServiceAccount
    name: spark
    namespace: data-jobs
roleRef:
  kind: Role
  name: spark-driver
  apiGroup: rbac.authorization.k8s.io

A missing or wrong service account shows up as a driver that starts and then fails with a Forbidden error when creating executor pods.

Dependencies

Option How When
Bake into the image local:///opt/app/... Production: reproducible and fast to start
Remote URLs s3a://, https:// paths for --py-files, --jars, --files Shared artefacts in object storage
Upload from the submitting machine Local paths plus spark.kubernetes.file.upload.path=s3a://bucket/path Development, ad-hoc jobs

Scaling, shuffle and spot nodes

  • Dynamic allocation works on Kubernetes without an external shuffle service, using shuffle tracking (spark.dynamicAllocation.shuffleTracking.enabled, on by default) so executors holding live shuffle data are kept.
  • Executor loss is common on preemptible or spot nodes. Graceful decommissioning (spark.decommission.enabled, spark.storage.decommission.enabled, both false by default) lets an executor that is being shut down migrate its shuffle and cached blocks to others first.
  • Spark’s local directories default to emptyDir volumes in the pod; for heavy shuffles, mount fast local disks or volumes and point Spark at them.

spark-submit or an operator

Plain spark-submit is a one-shot client. Many teams instead use a Kubernetes operator, which adds a SparkApplication custom resource you apply with kubectl or GitOps tools; the operator submits the job, restarts it according to a policy, and reports status. Two operators are common: the Apache Spark project’s own spark-kubernetes-operator, and the older Kubeflow spark-operator. Both are separate projects with their own release cycles, so check which Spark versions the release you choose supports.

Debugging on Kubernetes

kubectl get pods -n data-jobs -l spark-app-name=daily-sales      # driver and executor pods
kubectl logs -n data-jobs daily-sales-driver                     # driver log (pod name varies)
kubectl describe pod -n data-jobs <executor-pod>                 # OOMKilled, Pending, image pull errors
kubectl port-forward -n data-jobs <driver-pod> 4040:4040         # live Spark UI while the job runs
Symptom Usual cause
Driver pod Pending Not enough CPU or memory in the namespace quota or on nodes
Executor pods Pending forever Requests too large for any node; reduce request.cores or memory
ImagePullBackOff Wrong image name or missing pullSecrets
Executors OOMKilled Pod memory limit exceeded: raise memory overhead (Python workers)
Driver fails with Forbidden Service account lacks RBAC permissions

Pitfalls

  • No event logs. The live UI dies with the driver pod; without event logs in durable storage, a failed job leaves only container logs behind.
  • Setting spark.executor.cores higher than the CPU limit throttles tasks; keep requests, limits and cores consistent.
  • Dependency drift between the image and what the code expects; pin versions in the image.
  • Leaving completed driver pods around indefinitely; add cleanup (a TTL controller or the operator’s policy).

In interviews

“How does Spark run on Kubernetes?” Describe spark-submit creating the driver pod, the driver creating executor pods through the API with its service account, images carrying Spark and dependencies, pod memory being heap plus overhead (and OOMKilled when exceeded), dynamic allocation via shuffle tracking, and operators for declarative management. Comparing it with YARN (queues versus namespaces and quotas, NodeManagers versus pods, external shuffle service versus shuffle tracking) is a strong finish.

Monitoring with the Spark History Server

What it is

The Spark UI is served by the driver and disappears when the application ends. To inspect finished applications, Spark writes an event log: a stream of JSON events (jobs, stages, tasks, executors, SQL plans and metrics) that encodes everything the UI shows. The History Server reads event logs from a shared directory and rebuilds the same UI, plus the same REST API, for every past application.

Turning on event logs

On every application (usually in spark-defaults.conf or the platform’s defaults):

Setting Default (4.2) Purpose
spark.eventLog.enabled false Write event logs
spark.eventLog.dir file:///tmp/spark-events Base directory; use HDFS or object storage on clusters
spark.eventLog.compress true (since 4.0) Compress event logs
spark.eventLog.rolling.enabled true (since 4.0) Roll logs into several files instead of one growing file
spark.eventLog.rolling.maxFileSize 128m Size at which a rolled file is closed

The migration guide lists both 4.0 changes: event logs are now rolled and compressed by default; set spark.eventLog.rolling.enabled=false and spark.eventLog.compress=false to restore the earlier behaviour.

Seeing it locally

This writes an event log for a small application, then starts a real History Server from the PySpark installation (on port 18099 to avoid clashing with anything on the default 18080) and queries its REST API:

import json, os, re, subprocess, tempfile, time, urllib.request
import pyspark
from pyspark.sql import SparkSession, functions as F

log_dir = tempfile.mkdtemp(prefix="spark-events-")
spark = (SparkSession.builder.master("local[2]").appName("history-demo")
         .config("spark.eventLog.enabled", "true")
         .config("spark.eventLog.dir", "file://" + log_dir)
         .getOrCreate())
sc = spark.sparkContext

spark.range(0, 100_000).groupBy((F.col("id") % 10).alias("k")).count().collect()
app_id = sc.applicationId
spark.stop()

for root, _, files in os.walk(log_dir):
    for name in sorted(files):
        rel = os.path.relpath(os.path.join(root, name), log_dir).replace(app_id, "<app-id>")
        print(rel)
eventlog_v2_<app-id>/.appstatus_<app-id>.crc
eventlog_v2_<app-id>/appstatus_<app-id>
eventlog_v2_<app-id>/events_1_<app-id>.zstd

Because rolling is on by default, the log is a directory (eventlog_v2_...) holding numbered events_N_... files compressed with zstd, plus an appstatus marker file (and its checksum) that tells the History Server whether the application is still running.

spark_home = os.path.dirname(pyspark.__file__)
env = dict(os.environ, SPARK_HOME=spark_home,
           SPARK_HISTORY_OPTS=f"-Dspark.history.fs.logDirectory=file://{log_dir} -Dspark.history.ui.port=18099")
server = subprocess.Popen([os.path.join(spark_home, "bin", "spark-class"),
                           "org.apache.spark.deploy.history.HistoryServer"],
                          env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)

opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
apps = []
try:
    for _ in range(90):
        try:
            apps = json.load(opener.open("http://localhost:18099/api/v1/applications"))
            if apps:
                break
        except OSError:
            pass
        time.sleep(1)
    for a in apps:
        print("app found:", a["id"] == app_id, "| name:", a["name"], "| completed:", a["attempts"][0]["completed"])
    jobs = json.load(opener.open(f"http://localhost:18099/api/v1/applications/{app_id}/jobs"))
    print("jobs:", [j["status"] for j in jobs])
finally:
    server.terminate()
    server.wait()
app found: True | name: history-demo | completed: True
jobs: ['SUCCEEDED', 'SUCCEEDED']

On a server you would start it with sbin/start-history-server.sh and open http://<host>:18080.

Running a History Server for a platform

Setting Default (4.2) Purpose
spark.history.fs.logDirectory file:/tmp/spark-events Where to read event logs: the same location applications write to (hdfs://, s3a://, …)
spark.history.ui.port 18080 Web UI and REST API port
spark.history.fs.update.interval 10s How often new or updated logs are scanned
spark.history.retainedApplications 50 Applications whose UI data is kept in the server’s cache (others are reloaded on demand)
spark.history.fs.cleaner.enabled false Periodically delete old event logs
spark.history.fs.cleaner.maxAge 7d Age at which logs are deleted when the cleaner is on

Operational notes:

  • Retention. With the cleaner off (the default), event logs accumulate forever. Turn it on, or use object-storage lifecycle rules, and choose a retention long enough for incident reviews.
  • Large applications. Long streaming or very task-heavy applications produce large logs; rolling plus compaction (spark.history.fs.eventLog.rolling.maxFilesToRetain) keeps them manageable.
  • Security. Event logs contain configuration and query text; restrict access to the log directory and enable the History Server’s ACLs.
  • On Kubernetes, run the History Server as a Deployment pointed at the shared bucket; on managed platforms, use the provider’s persistent history UI.

What to look for when reviewing a finished application

The History Server shows the same tabs as the live UI, frozen at the end of the run:

  1. Jobs and timeline: which job took the time, executors added and lost, and failed jobs with their error messages.
  2. Stages: retries (a stage attempt number above 0 means fetch failures or lost executors), spill, shuffle sizes, and skew from the task-metric percentiles.
  3. SQL: the final AQE plan, join strategies, rows per operator, files and partitions read.
  4. Executors: GC time, failed tasks per executor, and executors removed (and why, if the cluster manager reported it).
  5. Environment: the exact configuration that ran, which is what you compare between a good run and a bad run.

The REST API (/api/v1/applications/...) exposes the same data as JSON, which lets you build alerts or dashboards, for example tracking spill or shuffle bytes per job over time. For live metrics, Spark also has a metrics system with sinks (JMX, Prometheus servlet, Graphite and others) configured through metrics.properties.

Pitfalls

  • Writing logs to local disk on ephemeral machines (pods, autoscaled nodes): they vanish with the machine. Use shared durable storage.
  • Mismatched directories: applications write to spark.eventLog.dir, the server reads spark.history.fs.logDirectory; they must point to the same place.
  • Incomplete applications: a driver killed abruptly leaves its log marked as in progress; the History Server lists it as incomplete.
  • Assuming the History Server collects logs: it only reads what applications wrote. If an application ran without event logging, nothing can be recovered.

In interviews

“How do you debug a job that failed last night?” Open it in the History Server (event logs enabled in durable storage), read the failed job’s error and the stage with retries, check the final SQL plan and task metrics for skew or spill, and compare the Environment tab with a good run. Mention that Spark 4 rolls and compresses event logs by default and that cleanup is off unless you enable it.

Practice questions

Walk through what happens when you spark-submit to Kubernetes in cluster mode.

spark-submit sends a request to the Kubernetes API server to create a driver pod from the configured image, along with a ConfigMap of Spark settings and a headless service for the driver. The driver starts, uses its service account to create executor pods, and executors connect back to it. Tasks run as usual. When the application ends, executor pods are deleted by default and the driver pod remains in Completed or Error state with its logs. With waitAppCompletion=true (default) spark-submit reports the final status.

Executor pods are repeatedly OOMKilled although the Spark UI shows the heap far from full. Why?

Kubernetes enforces the pod memory limit, which is executor memory plus overhead. Off-heap usage (Python workers, Arrow buffers, native libraries, network buffers) counts against the limit but not the JVM heap. Increase spark.executor.memoryOverhead (or spark.executor.memoryOverheadFactor), set spark.executor.pyspark.memory for Python-heavy jobs, and reduce per-task Python memory.

How do you make dynamic allocation work on Kubernetes without an external shuffle service?

Enable spark.dynamicAllocation.enabled and rely on shuffle tracking (spark.dynamicAllocation.shuffleTracking.enabled, on by default), which keeps executors that hold shuffle data needed by later stages. Set minExecutors and maxExecutors, and consider graceful decommissioning so blocks migrate when executors are removed.

What do you need for the History Server to show yesterday’s applications?

Every application must have written event logs (spark.eventLog.enabled=true) to a durable shared location (spark.eventLog.dir), and the History Server’s spark.history.fs.logDirectory must point to that same location. The logs must still exist, so check the cleaner settings or storage lifecycle rules.

What changed about event logs in Spark 4.0?

They are rolled by default (spark.eventLog.rolling.enabled=true), so each application gets a directory of numbered event files, and compressed by default (spark.eventLog.compress=true). Setting both to false restores the Spark 3 behaviour of one uncompressed file.

When would you use a Kubernetes operator instead of plain spark-submit?

When you want jobs declared as Kubernetes resources (a SparkApplication manifest applied through GitOps or kubectl), automatic restarts, status reporting through the Kubernetes API, and lifecycle management without keeping a spark-submit client process. Plain spark-submit is simpler for ad-hoc runs or when an orchestrator such as Airflow already manages submission and retries.

Key takeaways

  • On Kubernetes, spark-submit creates a driver pod, and the driver creates executor pods using its service account; there is no Spark master or YARN.
  • Images carry Spark, connectors, Python packages and often the job code; pin them to the Spark version you tested.
  • Pod memory is heap plus overhead and is enforced as a limit; OOMKilled usually means too little overhead.
  • Dynamic allocation on Kubernetes uses shuffle tracking; decommissioning helps with spot and preemptible nodes.
  • Enable event logs to durable storage on every application; the History Server rebuilds the UI and REST API from them.
  • Spark 4.0 rolls and compresses event logs by default; the History Server’s cleaner is off unless you enable it.

By DataDank Editorial · Last reviewed Oct 2026 · The event-log and History Server example runs on PySpark 4.2.0 in local mode, starting a real History Server from the PySpark installation on port 18099. Kubernetes commands and manifests are not executed (no cluster available); configuration names and defaults were checked against the Spark 4.2 runtime and the Running on Kubernetes documentation.

Progress is saved in this browser only. No account needed.

Search
Filter by type