Apache Livy is a REST server that sits in front of a Spark cluster and lets programs that are not Spark drivers run Spark work. A notebook, a workflow scheduler or a web backend sends JSON over HTTP; Livy launches a Spark driver on the cluster on the caller's behalf, runs code in it, and hands back results and status. The caller never needs a Spark installation, cluster configuration files or network access to the executors.

It came out of Cloudera, entered the Apache Incubator, and still ships there: 0.8.0-incubating arrived in September 2023 and 0.9.0-incubating in early 2026, adding Kubernetes submission, Java 17 and Spark 3.5 support. It is the engine behind the Spark kernels in sparkmagic and is offered on managed platforms such as Amazon EMR. This article explains how it works from the request down to the driver, walks through the API with real calls, and covers the settings, security model and failure modes you need before putting it in front of users.

Advertisement

Why a REST server in front of Spark

Spark's native contract is that your program is the driver. spark-submit ships your code to a cluster manager, and a SparkSession in that process plans jobs and talks to executors. That is fine for packaged jobs, but awkward in three common situations:

  • Shared notebooks. Fifty analysts on laptops cannot each hold YARN client configs, Kerberos keytabs and a route to every executor port.
  • Services that trigger Spark. A web app or scheduler wants to start a job and poll its state, not embed a JVM driver.
  • Long-lived context. Interactive work wants a warm session with cached DataFrames across many small requests, not a cold start per query.

Livy solves all three by moving the driver onto the cluster and putting an HTTP API in front of it. The price is an extra stateful service in the request path, which is what most of the operational advice below is about.

Architecture: interactive and batch sessions

Livy has two kinds of session, and they work differently enough that you should think of them as two products sharing a server.

Clientsnotebooks, schedulers, appsLivy server :8998REST, auth, session managerState storefilesystem or ZooKeeperspark-submit launcherper session or batchInteractive sessionRSC driver + REPLBatch sessionyour jar or .py driverCluster managerYARN or Kubernetes1 HTTP JSONrecovery2 launch3a3b4 RPC backexecutors
Livy accepts REST calls, launches a Spark driver per session through spark-submit, and for interactive sessions keeps an RPC channel to a remote driver that runs code snippets. Batch sessions are fire-and-track: Livy only follows the application state.

An interactive session (POST /sessions) launches a Spark application whose driver runs Livy's remote Spark context (RSC) together with REPL interpreters for Scala, Python, R and SQL. When the driver starts it opens an RPC connection back to the Livy server. Each statement you post travels over that channel, executes in the driver's interpreter, and the output is stored on the server until you fetch it. Since 0.5 one session can run all four languages; you choose the language per statement with kind (spark, pyspark, sparkr or sql), and the interpreters share one SparkContext, so a table registered from Python is visible to SQL.

A batch session (POST /batches) is a remote spark-submit: you name a jar or Python file already readable by the cluster, Livy launches it, and from then on it only tracks the application's state and log. There is no RPC channel and no statement API. If the driver needs results returned, it writes them somewhere, typically object storage or a table.

Both kinds launch through spark-submit on the Livy host, so that host needs a Spark distribution and the cluster client configuration. The master and deploy mode come from livy.spark.master and livy.spark.deploy-mode. Cluster deploy mode is the norm in production: drivers run in YARN containers or Kubernetes pods, not on the Livy host, so one Livy server can front many sessions without becoming a memory bottleneck.

Advertisement

The REST API by example

The API is small. The following sequence creates a session, runs PySpark, polls, and cleans up. Field names and paths are those in the Livy REST documentation.

# 1. create an interactive session (kind is optional; it sets the default statement language)
curl -s -X POST http://livy:8998/sessions -H 'Content-Type: application/json' -d '{
  "kind": "pyspark", "name": "analyst-jo", "proxyUser": "jo",
  "driverMemory": "4g", "executorMemory": "8g", "executorCores": 4,
  "conf": {"spark.dynamicAllocation.enabled": "true",
           "spark.dynamicAllocation.maxExecutors": "20"},
  "heartbeatTimeoutInSecond": 600
}'
# -> {"id": 17, "state": "starting", "appId": null, ...}

# 2. wait until state == "idle"
curl -s http://livy:8998/sessions/17/state        # {"id":17,"state":"idle"}

# 3. run a statement
curl -s -X POST http://livy:8998/sessions/17/statements -H 'Content-Type: application/json' \
  -d '{"kind": "pyspark", "code": "spark.range(10**8).selectExpr(\"sum(id)\").collect()"}'
# -> {"id": 0, "state": "waiting", "output": null}

# 4. poll the statement; output appears when state == "available"
curl -s http://livy:8998/sessions/17/statements/0
# -> {"id":0,"state":"available","output":{"status":"ok","execution_count":0,
#     "data":{"text/plain":"[Row(sum(id)=4999999950000000)]"}}}

# 5. cancel a runaway statement, read the driver log, end the session
curl -s -X POST http://livy:8998/sessions/17/statements/0/cancel
curl -s "http://livy:8998/sessions/17/log?from=0&size=100"
curl -s -X DELETE http://livy:8998/sessions/17

# batch: run a packaged job and track it
curl -s -X POST http://livy:8998/batches -H 'Content-Type: application/json' -d '{
  "file": "s3a://jobs/etl/daily_rollup.py", "args": ["--date", "2026-10-02"],
  "proxyUser": "etl", "conf": {"spark.sql.shuffle.partitions": "400"}
}'

Two details matter. First, output is whatever the REPL prints, returned as MIME-typed data, so a collect() of a million rows becomes a multi-megabyte JSON string held on the driver and the Livy server. Aggregate on the cluster, return small results, and write large ones to storage. Second, the session lives until you delete it or it times out; every abandoned session holds a driver and possibly executors.

States and a well-behaved client

Clients should be written against the state machine rather than against timings. Session states are not_started, starting, idle, busy, shutting_down, error, dead, killed and success. Statement states are waiting, running, available, error, cancelling and cancelled. Note that a statement whose Python code raised an exception is usually available with output.status equal to error; the statement state error means Livy itself failed to run it. A small client that respects both:

import time, requests

TERMINAL_SESSION = {"error", "dead", "killed", "success"}

class LivyClient:
    def __init__(self, url, user=None, auth=None):
        self.url, self.s = url.rstrip("/"), requests.Session()
        self.s.auth = auth                      # e.g. HTTPKerberosAuth()
        self.s.headers["X-Requested-By"] = user or "client"   # needed when CSRF protection is on

    def _wait(self, path, ok, bad, timeout):
        delay, deadline = 0.5, time.monotonic() + timeout
        while time.monotonic() < deadline:
            body = self.s.get(self.url + path, timeout=30).json()
            if body["state"] in ok:
                return body
            if body["state"] in bad:
                raise RuntimeError(f"{path} ended in {body['state']}")
            time.sleep(delay)
            delay = min(delay * 1.5, 5.0)       # back off; do not hammer the server
        raise TimeoutError(path)

    def open(self, timeout=600, **spec):
        sid = self.s.post(self.url + "/sessions", json=spec, timeout=30).json()["id"]
        self._wait(f"/sessions/{sid}", {"idle"}, TERMINAL_SESSION, timeout)
        return sid

    def run(self, sid, code, kind="pyspark", timeout=3600):
        st = self.s.post(f"{self.url}/sessions/{sid}/statements",
                         json={"code": code, "kind": kind}, timeout=30).json()
        body = self._wait(f"/sessions/{sid}/statements/{st['id']}",
                          {"available"}, {"error", "cancelled"}, timeout)
        out = body["output"]
        if out["status"] != "ok":
            raise RuntimeError(f"{out.get('ename')}: {out.get('evalue')}")
        return out["data"].get("text/plain")

    def close(self, sid):
        self.s.delete(f"{self.url}/sessions/{sid}", timeout=30)

Wrap open and close in a context manager or a try/finally so a crashing client still deletes its session.

Configuration that decides safety and cost

Most of Livy's behaviour is set in livy.conf. The settings below are the ones that decide whether a deployment is safe; defaults are from the shipped template.

livy.server.port = 8998
livy.spark.master = yarn                    # template default is local; k8s://https://... since 0.9
livy.spark.deploy-mode = cluster            # drivers in containers, not on the Livy host

livy.server.session.timeout-check = true
livy.server.session.timeout = 1h            # idle sessions are reclaimed after this
livy.server.session.state-retain.sec = 600s # how long finished sessions stay visible

livy.impersonation.enabled = true           # run sessions as the caller, not as 'livy'
livy.server.access-control.enabled = true   # restrict who may view or modify which sessions
livy.server.csrf-protection.enabled = true  # mutating calls must carry X-Requested-By

livy.server.recovery.mode = recovery
livy.server.recovery.state-store = zookeeper
livy.server.recovery.state-store.url = zk1:2181,zk2:2181,zk3:2181

livy.file.local-dir-whitelist = /opt/livy/approved-jars   # local paths sessions may reference
livy.server.yarn.app-lookup-timeout = 120s

The idle timeout is the single biggest cost lever. One hour is generous for a notebook user who walked away holding twenty executors; with dynamic allocation enabled the executors are released sooner, but the driver container stays until Livy reclaims the session. A per-session ttl or idleTimeout field can tighten this for individual clients.

Security: Livy is remote code execution

Livy executes arbitrary code by design, so every reachable endpoint is a remote code execution service. Treat it accordingly:

  • Authenticate every request. On Hadoop clusters this is SPNEGO with Kerberos (livy.server.auth.type = kerberos with a principal and keytab); elsewhere, put Livy behind a gateway such as Apache Knox or an authenticating reverse proxy and never expose port 8998 directly.
  • Impersonate. With livy.impersonation.enabled, the Spark application runs as the authenticated user, so HDFS permissions, Ranger policies and YARN queue ACLs apply to the human, not to a shared service account. proxyUser in the request, and the doAs query parameter, are honoured only for configured superusers, which is what a trusted gateway such as a notebook hub uses.
  • Enable access control. Without it, any authenticated user can list sessions and post statements into another user's session, which runs with that user's identity.
  • Limit file paths. Leave livy.file.local-dir-whitelist empty or narrow, so a request cannot ask the driver to load an arbitrary local file from the Livy host.
  • Keep secrets out of conf. Session conf values appear in the Spark UI and in Livy's responses. Use the mechanisms in managing secrets in Spark instead.

High availability and recovery

A Livy server is a single process holding the session table and, for interactive sessions, the RPC endpoints drivers connect to. If it dies with recovery off, running drivers are orphaned: they keep holding containers until they time out, and clients lose their session ids.

With livy.server.recovery.mode = recovery, Livy persists session metadata to the state store and, on restart, reloads it and reconnects to drivers that are still alive. Use ZooKeeper (or a highly available filesystem) as the store, run Livy under a supervisor that restarts it quickly, and keep the restart well inside the driver's heartbeat timeout. Recovery is not active-active load balancing: the session table belongs to one server, so scaling out means running several independent Livy instances and routing each user or tenant to one of them consistently.

Worked example: a notebook gateway for 40 analysts

A team runs Livy for 40 analysts on a YARN cluster through a JupyterHub with sparkmagic kernels. At peak, 30 notebooks are open but only about 8 are executing at any moment.

Before. Sessions use static allocation of 10 executors at 8 GB each, plus a 4 GB driver, and the default one-hour idle timeout. Thirty sessions hold 30 × (4 + 80) = 2,520 GB of container memory, although only about 670 GB is doing work. New sessions wait in the YARN queue, and users respond by never closing their notebooks.

After. Sessions enable dynamic allocation with a minimum of 0 and a maximum of 20 executors and a 60-second executor idle timeout, and Livy's session timeout drops to 30 minutes. Idle sessions now cost only their 4 GB driver: 22 idle × 4 GB plus 8 active sessions averaging 12 executors gives 88 + 8 × (4 + 96) = 888 GB. The busiest analysts get more executors than before, and the queue clears. The dynamic allocation guide covers the shuffle-tracking settings this depends on.

Failure modes

  • Sessions stuck in starting. The YARN queue is full or the application was rejected. Check the session log endpoint and the YARN application; raise livy.server.yarn.app-lookup-timeout only if submission is genuinely slow.
  • Driver cannot reach the server. Interactive sessions need the driver to connect back to Livy; a firewall or wrong advertised address leaves the session dead after start. Open the RSC port range and set the server's reachable address.
  • Huge statement outputs. A large collect() or toPandas() can exhaust driver or server memory. Cap result sizes in the client and with spark.driver.maxResultSize.
  • Session leaks. Clients that crash before cleanup leave drivers holding containers until the timeout. Monitor the count of idle sessions per user.
  • Lost sessions after a restart. Recovery was off, or the state store was local to a replaced host.
  • Dependency drift. Python packages differ between the Livy host, driver and executors. Ship environments as archives or container images rather than relying on node installs.

Livy, Spark Connect or spark-submit

Livy is not the only way to run Spark remotely. Spark Connect, built into Spark since 3.4, gives clients a thin DataFrame API over gRPC with the plan executed on a server-side driver; it returns structured Arrow results instead of REPL text, and suits applications. Livy remains the fit for notebook kernels that send code strings, for batch submission over HTTP, and for Hadoop estates that rely on its Kerberos impersonation.

NeedLivySpark ConnectPlain spark-submit
Notebook kernel sending code textNativePossible, different modelNo
Typed DataFrame client in an appAwkwardNativeApp is the driver
Batch submission over HTTPYesNoNeeds client config
Per-user impersonation on HadoopBuilt inDepends on deploymentVia keytab or proxy user
KubernetesSince 0.9.0-incubatingYesYes, see Spark on Kubernetes

What to do next

  1. Put Livy behind authentication, turn on impersonation, access control and CSRF protection, and confirm port 8998 is not reachable from outside the gateway.
  2. Set livy.spark.deploy-mode = cluster so drivers do not run on the Livy host.
  3. Enable recovery with a ZooKeeper or HA filesystem state store, then kill the Livy process in staging and confirm sessions survive.
  4. Turn on dynamic allocation for interactive sessions and shorten the idle timeout to match how your users actually work.
  5. Write clients against session and statement states, with back-off polling and guaranteed session deletion.
  6. Alert on idle session count per user, sessions stuck in starting, and statement output size.
  7. For new application integrations, evaluate Spark Connect before adding more Livy clients.
Key takeaway: Livy moves the Spark driver onto the cluster and puts an HTTP API in front of it. Interactive sessions keep a remote driver with REPL interpreters reachable over RPC; batch sessions are a remote spark-submit that Livy only tracks. Because it executes arbitrary code, run it authenticated, impersonating, access-controlled and in cluster deploy mode; enable recovery; and keep sessions cheap with dynamic allocation and short idle timeouts.