Apache Spark courseLesson 9 of 9
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.
On this page
- Spark on Kubernetes
- What it is
- How a submission works
- A cluster-mode submission
- The configuration you need to know
- Images
- RBAC: the service account
- Dependencies
- Scaling, shuffle and spot nodes
- spark-submit or an operator
- Debugging on Kubernetes
- Pitfalls
- In interviews
- Monitoring with the Spark History Server
- What it is
- Turning on event logs
- Seeing it locally
- Running a History Server for a platform
- What to look for when reviewing a finished application
- Pitfalls
- In interviews
- Practice questions
- 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
spark-submit --master k8s://https://<api-server>:<port> --deploy-mode clustercreates a driver pod (plus a ConfigMap with the Spark configuration and a headless service so executors can reach the driver).- The driver, using its service account, asks the API server to create executor pods.
- Executors connect back to the driver and run tasks as on any other cluster manager.
- When the application finishes, executor pods are deleted (
spark.kubernetes.executor.deleteOnTermination, defaulttrue). The driver pod stays inCompletedstate 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-awsfor 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, bothfalseby default) lets an executor that is being shut down migrate its shuffle and cached blocks to others first. - Spark’s local directories default to
emptyDirvolumes 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.coreshigher 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:
- Jobs and timeline: which job took the time, executors added and lost, and failed jobs with their error messages.
- Stages: retries (a stage attempt number above 0 means fetch failures or lost executors), spill, shuffle sizes, and skew from the task-metric percentiles.
- SQL: the final AQE plan, join strategies, rows per operator, files and partitions read.
- Executors: GC time, failed tasks per executor, and executors removed (and why, if the cluster manager reported it).
- 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 readsspark.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;
OOMKilledusually 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.
Progress is saved in this browser only. No account needed.

