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.
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.
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.
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 = 120sThe 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 = kerberoswith 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.proxyUserin the request, and thedoAsquery 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-whitelistempty 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
confvalues 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-timeoutonly 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()ortoPandas()can exhaust driver or server memory. Cap result sizes in the client and withspark.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.
| Need | Livy | Spark Connect | Plain spark-submit |
|---|---|---|---|
| Notebook kernel sending code text | Native | Possible, different model | No |
| Typed DataFrame client in an app | Awkward | Native | App is the driver |
| Batch submission over HTTP | Yes | No | Needs client config |
| Per-user impersonation on Hadoop | Built in | Depends on deployment | Via keytab or proxy user |
| Kubernetes | Since 0.9.0-incubating | Yes | Yes, see Spark on Kubernetes |
What to do next
- Put Livy behind authentication, turn on impersonation, access control and CSRF protection, and confirm port 8998 is not reachable from outside the gateway.
- Set
livy.spark.deploy-mode = clusterso drivers do not run on the Livy host. - Enable recovery with a ZooKeeper or HA filesystem state store, then kill the Livy process in staging and confirm sessions survive.
- Turn on dynamic allocation for interactive sessions and shorten the idle timeout to match how your users actually work.
- Write clients against session and statement states, with back-off polling and guaranteed session deletion.
- Alert on idle session count per user, sessions stuck in starting, and statement output size.
- For new application integrations, evaluate Spark Connect before adding more Livy clients.