Skip to content

How to recover from a Flink JobManager failure with HA

Flink
Chad Harris·October 3, 2026·12 min read

The JobManager pod was evicted, the node under it was lost, or a network partition cut it off from ZooKeeper or the Kubernetes API server. Afterwards the Flink web UI shows the jobs back and running from a checkpoint, shows them restarting every time leadership moves, or shows no running jobs at all because the JobManager came back with an empty job list and every running job and its state pointer is gone.

JobManager high availability (HA) is the Flink feature that keeps the second JobManager able to pick up where the first one stopped. The JobManager High Availability documentation describes it as hardening a Flink cluster against JobManager failures, so that the cluster re-executes the applications that were running at the time of the failure. Without it, the same page says there is a single JobManager per cluster, and if it crashes no new programs can be submitted and running programs fail.

This page covers the JobManager. If a TaskManager was lost and the job keeps failing and restarting, that is the restart strategy at work, and the guide to a Flink job stuck in a restart loop is the better starting point.

What JobManager HA protects

The JobManager coordinates the cluster. The Flink architecture documentation says there is always at least one JobManager, and that a high availability setup might have several, one of which is always the leader while the others are standby. The TaskManagers are separate processes that execute the tasks. A Flink deployment therefore runs both kinds of process, and HA is the part that covers the JobManager side.

The HA overview lists three things that Flink’s HA services provide.

  • Leader election, which selects a single leader out of a pool of candidates.
  • Service discovery, which lets the rest of the cluster find the address of the current leader.
  • State persistence, which keeps what the successor needs to resume the jobs: JobGraphs, user code jars and completed checkpoints.

The state is split in two places. The Kubernetes HA documentation and the ZooKeeper HA documentation both say JobManager metadata is persisted in a file system, and only a pointer to it is stored in Kubernetes or ZooKeeper. Both the pointer and the files have to survive for a recovery to work. The file system location is set by this key.

high-availability.storageDir

HA does not replace checkpointing. The HA overview lists completed checkpoints among the state the successor needs to resume the jobs, and the checkpoint data itself is written by the job to its configured checkpoint storage, which the checkpoints documentation describes. A streaming job resumes from the latest successful checkpoint after a failover, so the job needs checkpointing switched on for there to be state to resume. If jobs recovered but from an old or missing checkpoint, see Flink checkpoints that fail or time out. For how the JobManager fits into running Flink in production more widely, see the complete Flink guide.

What happens to running jobs during a failover

The Flink page on recovering batch job progress from JobMaster failures sets out the two outcomes. If HA is disabled, the job fails. If HA is enabled, a failover happens and the job is restarted. Streaming jobs resume from their latest successful checkpoint. Batch jobs have no checkpoints and start over from the beginning, unless batch job recovery is switched on, which itself requires HA.

execution.batch.job-recovery.enabled: true

A streaming job therefore redoes the work done since its last completed checkpoint. The production readiness checklist notes that exactly once sinks, such as Kafka, only make results visible when a checkpoint completes, so the checkpoint interval also sets how much output waits behind a failover.

In the current Flink documentation, jobs before and after a failure are matched by name, and jobs with identical names are then matched by submission order. The HA overview recommends giving each job a unique name, so that a non-deterministic submission order cannot attach a job to the wrong recovered state. The name is the argument to the execute call.

execute(jobName)

Confirm HA is configured

Check the configuration the running JobManager actually loaded, not the file in the repository. Flink’s REST API is the monitoring API Flink’s own dashboard uses, and it is designed to be used by custom monitoring tools as well, so any management tool reads the same values from the same endpoint. One call returns the cluster configuration.

GET /jobmanager/config

Look for the HA type first. The Flink configuration reference gives its default as NONE, and HA is only on when it is set to ZOOKEEPER, KUBERNETES or the class name of a custom factory.

high-availability.type

Then check the keys each HA service requires. Both need the storage directory. Kubernetes HA also needs a cluster id, and ZooKeeper HA needs the quorum address, with the root path and the cluster id recommended.

Kubernetes HA looks like this.

high-availability.type: kubernetes
high-availability.storageDir: s3://flink/recovery
kubernetes.cluster-id: <cluster-id>

ZooKeeper HA looks like this.

high-availability.type: zookeeper
high-availability.storageDir: hdfs:///flink/recovery
high-availability.zookeeper.quorum: address1:2181,address2:2181
high-availability.zookeeper.path.root: /flink
high-availability.cluster-id: /cluster_one

The ZooKeeper HA documentation says that cluster id should not be set by hand on YARN, native Kubernetes or another cluster manager, where it is generated, and that several HA clusters on bare metal each need their own.

ZooKeeper HA needs a running ZooKeeper quorum, according to the HA overview, and Kubernetes HA only works on Kubernetes. If Flink HA shares the ZooKeeper ensemble that ran a Kafka cluster, note that the Apache Kafka 4.0 release is the first major release to operate entirely without ZooKeeper, with metadata handled by KRaft instead. Before that ensemble is decommissioned after a Kafka upgrade, confirm that Flink HA still has a ZooKeeper quorum of its own, or move to Kubernetes HA on Kubernetes.

Confirm HA works before you need it

A configuration that reads correctly can still fail on the first failover, so check three things the Flink documentation names as requirements.

First, the storage directory must be a file system every JobManager can read and write. The native Kubernetes documentation says that HA setups with several JobManagers sharing a persistent volume need a claim that supports ReadWriteMany, and that ReadWriteOnce in that case can cause mount failures or I/O errors.

Second, on Kubernetes, the pods need a service account that can create, edit and delete ConfigMaps, as the Kubernetes HA documentation lists in its prerequisites. The native Kubernetes documentation says the JobManager uses it to create the leader ConfigMaps, and the TaskManagers watch those ConfigMaps to find the JobManager and ResourceManager.

-Dkubernetes.service-account=flink-service-account

Third, on standalone Kubernetes, the standalone Kubernetes documentation says that with HA enabled Flink uses its own HA services for service discovery, so the JobManager pods should be started with their IP address, not a Kubernetes service, as the RPC address.

jobmanager.rpc.address

Then test it. Delete the leading JobManager pod in a staging cluster with a job running and checkpoints completing. After the new leader is up, the job should be RUNNING again, and its checkpoint statistics should show a restore from the latest checkpoint.

GET /jobs/:jobid/checkpoints
latest.restored

An empty restored block after a test failover means the job started without state, which points to the branches in the next section.

Diagnose a failover that did not recover

Rule out the TaskManager first. A failure that names a task, a TaskManager that left the cluster or a user code exception is handled by the restart strategy, not by HA, and the task failure recovery documentation covers it. If checkpointing is on and no restart strategy is configured, that documentation says Flink uses the exponential delay strategy, so a job that keeps restarting with backoff and no JobManager event in the log belongs on the restart loop page.

restart-strategy.type

When the JobManager itself is the problem, the JobManager log and a thread dump from the JobManager JVM are the two sources to read. The REST API reference lists the JobManager log files and returns the thread dump.

GET /jobmanager/logs
GET /jobmanager/thread-dump

HA was never enabled

The JobManager came back and the job list is empty, or every job is FAILED. That is the documented behavior without HA. The fix has two parts. Enable HA as above so the next failure recovers, and restart the lost jobs from the newest state that still exists. The checkpoints documentation says checkpoints are not retained by default and are deleted when a program is cancelled, so the last retained checkpoint or the last savepoint is the restore point. The guide to upgrading from a savepoint without losing state walks through that restore and how to confirm it.

execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION

The storage directory is not reachable

The new leader is elected but cannot load the job metadata, or a standby JobManager fails to start because its volume will not mount. The pointer in ZooKeeper or the ConfigMap is there, but the files it points to are not readable from this JobManager. Check that the storage directory is a shared file system rather than a local path, that the credentials work from every JobManager pod and that a shared volume is ReadWriteMany.

Leadership is revoked over and over

Failovers repeat with no job exception behind them, and the JobManager log shows leadership being granted and revoked, which separates a failover loop from a job restart loop.

With ZooKeeper, the ZooKeeper HA documentation says Flink treats a suspended connection as an error by default, invalidates the leadership of its components and triggers a failover. On an unstable network, that means every short connection blip is a failover. The documented trade is to tolerate suspended connections and act only when a connection is lost, which the same page says makes Flink more resilient to temporary connection problems but increases the risk of ZooKeeper timing problems.

high-availability.zookeeper.client.tolerate-suspended-connections: true

The ZooKeeper client retries with exponential backoff before it gives up. The documentation gives the defaults as a 5 second initial wait, a 60 second maximum wait and 3 attempts.

high-availability.zookeeper.client.retry-wait: 5 s
high-availability.zookeeper.client.max-retry-wait: 1 min
high-availability.zookeeper.client.max-retry-attempts: 3

With Kubernetes, the leader keeps its leadership by renewing a lease. The configuration reference says the leader gives up leadership if it cannot renew the lease within the renew deadline, and gives the defaults as a 15 second lease duration, a 15 second renew deadline and a 5 second retry period. A JobManager that loses leadership regularly is failing to renew its lease in time, so check the JobManager log for failed renewals and the health of the Kubernetes API server before changing these values. If the lease does need to be longer, change one value at a time.

high-availability.kubernetes.leader-election.lease-duration: 15 s
high-availability.kubernetes.leader-election.renew-deadline: 15 s
high-availability.kubernetes.leader-election.retry-period: 5 s

The wrong jobs came back

A cluster recovers jobs that belong to another cluster, or two jobs swap state. The cluster id is the node under which all coordination data for a cluster is placed, according to the ZooKeeper HA documentation, so the documentation requires a separate cluster id for every HA cluster on bare metal. On Kubernetes, the Kubernetes HA documentation lists the cluster id as required to identify the cluster.

kubernetes.cluster-id

Inside one cluster, duplicate job names let the name and submission order matching attach state to the wrong job.

A restart deleted what it should have kept

The Kubernetes HA documentation says that to keep HA data while restarting a cluster, delete the deployment and nothing else.

kubectl delete deployment <cluster-id>

The HA ConfigMaps are kept because they do not set an owner reference, and when the cluster restarts, all previously running jobs are recovered from the latest successful checkpoint. Deleting those ConfigMaps by hand, or tearing down the namespace, removes the pointers and turns the next start into a fresh cluster.

The HA overview also says HA data is kept only until a job reaches a terminal state, finished, cancelled or failed, and is then deleted. A job that reached FAILED before the JobManager restarted is not recovered by HA, because its HA data is already gone.

F1 Why a JobManager failure lost jobs, or a failover did not settle
What shows first What proves it, and the fix
HA was never enabled The JobManager comes back, but the job list is empty or the jobs are FAILED Proof: high-availability.type is missing or NONE in GET /jobmanager/config. Fix: set the type, the storageDir and the cluster id, and restore from the last retained checkpoint or savepoint.
HA storage not reachable The new leader starts but cannot read job metadata, or a standby JobManager fails to mount its volume Proof: high-availability.storageDir and the JobManager log. Fix: a file system every JobManager can reach, and a ReadWriteMany volume when several JobManagers share a PVC.
Missing ConfigMap permissions Leader election never completes on Kubernetes Proof: the service account cannot create, edit and delete ConfigMaps. Fix: grant those permissions to the account the JobManager and TaskManager pods run as.
Leadership revoked over and over Failovers repeat with no job exception behind them Proof: ZooKeeper suspended connections, or Kubernetes lease renewals that miss their deadline, in the JobManager log. Fix: the network or API server first, then the lease or suspended-connection settings.
Wrong cluster id or job name A cluster recovers another cluster's jobs, or recovers the wrong job Proof: high-availability.cluster-id or kubernetes.cluster-id shared between clusters, or duplicate job names. Fix: one cluster id per cluster and a unique name for every job.
Drawn from the Apache Flink JobManager high availability, Kubernetes HA, ZooKeeper HA, configuration and Kubernetes deployment documentation, all linked on this page. JobManager High Availability documentation

Standby JobManagers and recovery time

On Kubernetes, a single JobManager pod with HA is often enough. The standalone Kubernetes documentation says Kubernetes restarts a crashed JobManager pod, and that standby JobManagers are for faster recovery. On native Kubernetes, the number of JobManager pods is one setting, and the configuration reference says HA should be enabled when standby JobManagers are started.

kubernetes.jobmanager.replicas: 2

A standby does not shorten the replay. The amount of work redone after a failover is set by the checkpoint interval, and the checkpoint guidance in the production readiness checklist, which also calls JobManager HA highly recommended for production, is the place to tune it. If jobs recovered and then fell behind, see finding the bottleneck behind Flink backpressure.

Where Flex fits

This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. The comparison of tools for self-managed Apache Flink sets Flex beside the Flink web UI and the REST API.

JobManager view. Flex documents these in its job manager documentation.

  • A Details tab with cards for JobManager memory used, heap used, non-heap used and the Flink version, and graphs of memory use over the past hour.
  • A Configuration tab that lists the static configuration the JobManager was started with, which the documentation describes as essential for verifying that the cluster runs with the intended settings. That is where the HA keys above can be read.
  • A Logs tab with live access to the JobManager log files, which the documentation names as the primary tool for cluster-level events, including high-availability service problems.
  • A Thread dump tab that requests a live thread dump from the JobManager JVM, for an unresponsive JobManager.

Job states. The Flex management overview shows counts of running, finished, cancelled and failed jobs, and the number of task managers and available task slots, so a JobManager restart that left jobs FAILED or TaskManagers missing shows as a change in those counts.

What stays in Flink. The Flex documentation does not describe showing which JobManager holds leadership, how many standbys exist, the contents of the HA ConfigMaps or ZooKeeper nodes, or triggering a failover. Leader election and HA storage stay with Flink, ZooKeeper and Kubernetes.

Flink management tools use the REST API that Flink’s own dashboard uses. That endpoint accepts connections from any client by default and does not authenticate them, according to the Flink SSL setup documentation, which recommends an authenticating proxy in front of it. The REST API reference includes endpoints that cancel jobs and shut down the cluster, so access during an incident needs a proxy or a tool with its own access control in front. The Flex role-based access control documentation describes granting or denying roles actions on Flink resources.

To read JobManager configuration, logs and thread dumps beside job state, Flex connects to the cluster through the Flink REST URL, as its Flink provider documentation shows.

FAQ

Does Flink need ZooKeeper for high availability?

Not on Kubernetes. Flink ships two HA services. ZooKeeper HA works with every deployment and needs a running ZooKeeper quorum, and Kubernetes HA works only on Kubernetes and stores its pointers in ConfigMaps.

Do running Flink jobs stop during a JobManager failover?

Yes. With HA enabled the jobs are restarted once a new leader takes over, and streaming jobs resume from their latest successful checkpoint. Without HA, running jobs fail.

Is one JobManager with HA enough on Kubernetes?

Often it is. Kubernetes restarts a crashed JobManager pod, and HA lets it recover the jobs. Standby JobManagers make recovery faster but require HA to be enabled.

Why did my Flink jobs not recover after the JobManager restarted?

The usual causes are that high-availability.type was never set, the storage directory was not reachable from the new JobManager, the HA ConfigMaps or ZooKeeper nodes were deleted, or the job had already reached a terminal state. Start with GET /jobmanager/config and the JobManager log.

Related reading