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.

Advertisement

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.

Driver: app Aevent log writerDriver: app Brolling filesDriver: app Cin progressLog directoryHDFS / S3 / GCS / ABFSappendScan loop (FsHistoryProvider)every fs.update.intervalListing DB + app storesstore.path on local diskReplay into KV storememory / disk / hybridApp UI cacheretainedApplicationsBrowser / REST:18080 /api/v1listread logpagesThe server is a cache in front of the log directory: everything it shows can be rebuilt from the event logs.
Drivers write event logs to shared storage; the History Server scans the directory, replays logs on demand into a local store, and serves pages and the REST API from a small cache of application UIs.

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 true

The 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.

Advertisement

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:

StoreHow to enableBehaviour
In memoryDefault when spark.history.store.path is unsetFast pages, but every cached application lives on the heap; large apps cause long GC pauses and OutOfMemoryError
DiskSet spark.history.store.pathReplayed data is written to a local LevelDB or RocksDB store and survives restarts; capped by store.maxDiskUsage (10g default)
HybridAlso set store.hybridStore.enabled=trueReplays 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=200g holds 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

SymptomLikely causeFix
Application missing from the listEvent logging off, wrong directory, or the scan has not reached it yetCheck the app's spark.eventLog.* settings; check the server log for parse errors
Finished app shown as incompleteDriver 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 outFull replay of a huge log with the in-memory storeDisk or hybrid store on SSD, more heap, PROTOBUF serializer
Server OOMs or pauses under loadDefault 1 GB heap, many cached UIsRaise SPARK_DAEMON_MEMORY; lower retainedApplications
Scans get slower every weekLog directory never cleanedEnable the cleaner, raise the update interval
Old micro-batches vanishedCompaction discarded themWorking as designed; lower maxFilesToRetain only for streaming apps
Permission denied reading logsServer identity lacks read access, or apps write with restrictive ACLsGrant 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

  1. Set event logging, compression and rolling in cluster-wide spark-defaults, and verify one job of each type appears on the server.
  2. 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.
  3. Enable the cleaner with a retention your debugging and audit needs actually require.
  4. Enable compaction only for long-running streaming applications, after testing on a copy of their logs.
  5. Put authentication in front of the server and turn on history UI ACLs.
  6. Pin the server to the same Spark version as your jobs and upgrade them together.
  7. Alert when the daily count of logged applications falls below what the scheduler launched.
Key takeaway: The Spark History Server is a cache in front of a directory of event logs: it scans the directory for a listing, replays a log into a key-value store when someone opens an application, and keeps a few replayed UIs in memory. Most operational problems follow from that design. Big applications are slow because replay is expensive, so use a disk or hybrid store and a real heap. Directories grow without the cleaner, and streaming logs need rolling and, carefully, lossy compaction. Because every piece of state can be rebuilt from the logs, the server is easy to redeploy, but only if every application logs to the right place.