When a Spark application ends, its live UI disappears with the driver. The History Server is how you still see the jobs, stages, SQL plans and executors of a run that finished last night, or crashed an hour ago. The Spark UI article explains how to read those pages and query them through the REST API. This article is about the server itself: how it turns event logs into pages, why a large application can take minutes to open, how to size and configure it, how rolling logs and compaction work, and what breaks when you run it for a whole platform.
Defaults below are from the Spark 4.x documentation. Several changed between major versions, so check the configuration page for the version you run.
What the History Server actually is
A running driver keeps its UI state in an in-memory key-value store fed by listener events: job started, stage submitted, task ended, executor added, SQL execution finished. With event logging on, the driver also writes each of those events as a JSON line to a file. The History Server is a separate, long-running daemon that reads those files and replays the events into the same kind of store, which is why it shows the same pages as the live UI.
Two consequences follow. First, the History Server knows nothing the event log does not contain: if an application did not enable event logging, or wrote to a different directory, it is invisible. Second, the server holds no unique state. Its local store is a cache that can be deleted and rebuilt from the logs, which makes it simple to redeploy.
The application side: writing good event logs
Every application that should appear must write event logs to a directory the server can read. Set this cluster-wide in spark-defaults rather than trusting each job. The directory must exist; Spark does not create it, and a missing directory fails the application at start-up.
# spark-defaults.conf for applications (or --conf on spark-submit)
spark.eventLog.enabled true
spark.eventLog.dir s3a://acme-spark-logs/events
# compress: default in the 4.x docs; set it explicitly on 3.x
spark.eventLog.compress true
spark.eventLog.compression.codec zstd
# rolling: essential for long-running and streaming apps
spark.eventLog.rolling.enabled true
spark.eventLog.rolling.maxFileSize 128m
# Optional: peak executor memory per stage in the History Server's stage pages
spark.eventLog.logStageExecutorMetrics trueThe 4.x documentation defaults compression to on with zstd and rolling to on with 128 MB files; older 3.x releases defaulted both to off, so on those versions a long streaming job wrote one ever-growing uncompressed file. Event logs compress very well, because they are repetitive JSON. With rolling enabled, each application writes a directory whose name starts with eventlog_v2_, holding numbered event files plus a status marker that shows whether the application is still running.
The scan loop and the listing
The default provider, FsHistoryProvider, lists the log directory every spark.history.fs.update.interval (10 seconds by default). For each new or changed log it reads just enough to populate the application listing (the ID, name, user, attempt and start and end times) and stores that in a listing database. For logs still in progress it avoids re-reading whole files: spark.history.fs.inProgressOptimization.enabled (on by default) and spark.history.fs.endEventReparseChunkSize (1 MB) let it read only the tail to find whether the application has finished.
On object storage the listing itself is the cost. A directory with 400,000 application logs means a large paginated LIST call every scan, which is slow and, on some clouds, billed per request. Raise the interval to 30 or 60 seconds for big directories, and keep the directory small with the cleaner (below). spark.history.fs.numReplayThreads (25 percent of cores by default) controls how many logs are parsed in parallel during a scan.
Replay: why a big application is slow to open
The listing is cheap; the pages are not. When someone opens an application, the server replays its entire event log into a key-value store for that application, then serves pages from it. A job with 200,000 tasks produces at least 400,000 task start and end events, each carrying metrics, so replay means parsing hundreds of megabytes of JSON. The last spark.history.retainedApplications (50 by default) replayed UIs are kept in an in-memory cache; opening an application outside that set replays it again.
Where the per-application store lives decides the server's memory profile:
| Store | How to enable | Behaviour |
|---|---|---|
| In memory | Default when spark.history.store.path is unset | Fast pages, but every cached application lives on the heap; large apps cause long GC pauses and OutOfMemoryError |
| Disk | Set spark.history.store.path | Replayed data is written to a local LevelDB or RocksDB store and survives restarts; capped by store.maxDiskUsage (10g default) |
| Hybrid | Also set store.hybridStore.enabled=true | Replays into memory (up to hybridStore.maxMemoryUsage, 2g default) then dumps to disk; faster first load than pure disk |
For any shared server, set a store path on a local SSD and consider the hybrid store. spark.history.store.serializer=PROTOBUF makes the disk store smaller and faster to read than the default JSON serializer. Because stores survive restarts, a redeploy no longer means re-replaying yesterday's large jobs, provided the volume is persistent.
Worked example: sizing a platform History Server
A platform runs about 3,000 applications a day, most small, with a few hundred nightly ETL jobs of 50,000 to 300,000 tasks and a dozen streaming applications that run for weeks. Engineers open roughly 200 distinct applications a day. The numbers below are illustrative; measure your own logs.
- Log volume: compressed logs average 5 MB, and large ETL jobs reach 150 to 400 MB. With a 30-day retention that is about 90,000 applications and a few terabytes of object storage.
- Memory: the default 1 GB daemon heap (
SPARK_DAEMON_MEMORY) is far too small. With the hybrid store at 4 GB and 30 retained UIs, a 12 GB heap leaves room for the listing and concurrent replays. - Disk: replayed stores run from a few megabytes to a few gigabytes for the biggest jobs; 200 GB of local SSD with
maxDiskUsage=200gholds a working set of recent large applications. - CPU: replay is single-threaded per application and mostly JSON parsing. Four to eight cores let several people open large jobs at once.
# conf/spark-defaults.conf on the History Server host
spark.history.fs.logDirectory s3a://acme-spark-logs/events
spark.history.fs.update.interval 30s
spark.history.fs.numReplayThreads 8
spark.history.store.path /var/lib/spark-history/store
spark.history.store.maxDiskUsage 200g
spark.history.store.serializer PROTOBUF
spark.history.store.hybridStore.enabled true
spark.history.store.hybridStore.maxMemoryUsage 4g
spark.history.retainedApplications 30
spark.history.fs.cleaner.enabled true
spark.history.fs.cleaner.maxAge 30d
spark.history.fs.eventLog.rolling.maxFilesToRetain 4
spark.history.ui.acls.enable true
spark.history.ui.admin.acls.groups data-platform
# conf/spark-env.sh
export SPARK_DAEMON_MEMORY=12g
Rolling logs and compaction for long-running applications
A streaming query that runs for three weeks produces an event log that grows without end, and replaying it means parsing every micro-batch's jobs since the start. Rolling splits the log into files of spark.eventLog.rolling.maxFileSize. Setting spark.history.fs.eventLog.rolling.maxFilesToRetain on the server enables compaction: with five rolled files and a retain value of 2, the server rewrites the oldest three into one compact file and then tries to delete the originals.
Compaction is lossy. It drops finished jobs with their stages and tasks, terminated executors, and finished SQL executions with their jobs. After compaction the UI cannot show those old micro-batches at all. That is usually the right trade for streaming, where the recent state matters, and usually the wrong one for batch jobs, where the finished jobs are the evidence. Batch applications tend to compact little anyway, and the server skips compaction when it would not save much. The documentation flags the feature as one to use with care, so test it on a copy of real logs first.
Retention, deployment and security
The cleaner is off by default, so without it the log directory grows forever. Enable spark.history.fs.cleaner.enabled and set maxAge (7 days default) or maxNum; the server needs delete permission on the directory. If your bucket already has a lifecycle rule, pick one mechanism, since a bucket rule deleting files mid-scan produces confusing errors. Driver logs, if you ship them to the history directory, have their own cleaner settings that inherit these values.
On Spark on Kubernetes there is no resource manager to host the server, so run it as a Deployment with a persistent volume for the store and workload identity for bucket access. Run the same Spark major and minor version as your jobs, because newer event formats may not replay on an older server.
apiVersion: apps/v1
kind: Deployment
metadata: {name: spark-history, namespace: spark}
spec:
replicas: 1
selector: {matchLabels: {app: spark-history}}
template:
metadata: {labels: {app: spark-history}}
spec:
serviceAccountName: spark-history # bound to read/delete on the log bucket
containers:
- name: shs
image: registry.example.com/spark:4.0.1 # same major.minor as your jobs
command: ["/opt/spark/bin/spark-class",
"org.apache.spark.deploy.history.HistoryServer"]
env:
- {name: SPARK_DAEMON_MEMORY, value: "12g"}
ports: [{containerPort: 18080}]
resources:
requests: {cpu: "4", memory: 14Gi}
limits: {memory: 14Gi}
readinessProbe:
httpGet: {path: /api/v1/applications?limit=1, port: 18080}
periodSeconds: 30
volumeMounts:
- {name: conf, mountPath: /opt/spark/conf}
- {name: store, mountPath: /var/lib/spark-history/store}
volumes:
- {name: conf, configMap: {name: spark-history-conf}}
- {name: store, persistentVolumeClaim: {claimName: spark-history-store}}Logs contain SQL text, table paths, configuration values and user names, so protect the server. Put it behind an authenticating proxy or a servlet filter, then enable spark.history.ui.acls.enable so each application's view ACLs, recorded in its log, are enforced; admins are granted through spark.history.ui.admin.acls and admin.acls.groups. Configure clusters so secrets such as access keys are redacted from the environment page before they reach the log.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Application missing from the list | Event logging off, wrong directory, or the scan has not reached it yet | Check the app's spark.eventLog.* settings; check the server log for parse errors |
| Finished app shown as incomplete | Driver killed before writing the end event (OOM, preemption, kill -9) | Expected; treat stale in-progress apps as crashed, and let the cleaner remove them |
| Opening a large app takes minutes or times out | Full replay of a huge log with the in-memory store | Disk or hybrid store on SSD, more heap, PROTOBUF serializer |
| Server OOMs or pauses under load | Default 1 GB heap, many cached UIs | Raise SPARK_DAEMON_MEMORY; lower retainedApplications |
| Scans get slower every week | Log directory never cleaned | Enable the cleaner, raise the update interval |
| Old micro-batches vanished | Compaction discarded them | Working as designed; lower maxFilesToRetain only for streaming apps |
| Permission denied reading logs | Server identity lacks read access, or apps write with restrictive ACLs | Grant read and delete on the prefix to the server's identity |
A small script against the REST API catches the most common data-quality problem, crashed drivers that stay listed as running:
import requests, time
BASE = "http://spark-history.spark:18080/api/v1"
def stale_in_progress(max_age_h=24):
# In-progress apps whose last update is old: crashed drivers that never
# wrote an end event. They clutter the listing and skew dashboards.
apps = requests.get(f"{BASE}/applications",
params={"status": "running"}, timeout=30).json()
cutoff = (time.time() - max_age_h * 3600) * 1000
return [a["id"] for a in apps
if a["attempts"][-1]["lastUpdatedEpoch"] < cutoff]Monitor the server like any service: heap usage and GC time, readiness-probe latency, local store disk use, and the number of applications listed per day compared with what your scheduler launched. A gap between the two is the earliest sign that some team's jobs stopped logging. Note also that dynamic allocation creates many executor add and remove events, which inflates logs for long jobs.
What to do next
- Set event logging, compression and rolling in cluster-wide spark-defaults, and verify one job of each type appears on the server.
- Give the server a persistent local store path on SSD, enable the hybrid store and PROTOBUF serializer, and raise the daemon heap from 1 GB.
- Enable the cleaner with a retention your debugging and audit needs actually require.
- Enable compaction only for long-running streaming applications, after testing on a copy of their logs.
- Put authentication in front of the server and turn on history UI ACLs.
- Pin the server to the same Spark version as your jobs and upgrade them together.
- Alert when the daily count of logged applications falls below what the scheduler launched.