Skip to content

How to fix a Flink job stuck in a restart loop

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

A Flink job in a restart loop shows the same few lines over and over in its event history, a few seconds or minutes apart.

Flink job entered state RESTARTING
Flink job entered state RUNNING
Flink job entered state RESTARTING

The job never stays up long enough to make progress, and it never stops either. On the default exponential delay strategy the number of attempts is unlimited, so a job with a permanent fault never reaches FAILED. The section on the restart strategy in force explains why.

A restart loop is a job that keeps failing and keeps being restarted by its restart strategy, because the cause of the failure is still there after each restart.

The counter that confirms a loop is the job-scope metric on the JobManager, which the Flink metrics documentation defines as the total number of restarts since the job was submitted, including full restarts, fine-grained restarts and restarts triggered by rescaling.

numRestarts

This page is about a Flink job that restarts itself. If the thing restarting is a Kafka Connect connector or task, the guide to diagnosing a failed Kafka Connect connector is the better starting point. For how restarts fit into Flink’s wider fault tolerance model, see the complete Flink guide.

The sections below follow this order:

  1. Rule out a checkpoint failure first.
  2. Find the root exception.
  3. Tell a transient failure from a permanent one.
  4. Check which restart strategy is in force.
  5. Make fixes, one change at a time, and watch numRestarts.

Rule out a checkpoint failure first

A failed checkpoint can cause a restart loop that does not look like a checkpoint problem. Flink tolerates no checkpoint failures by default, so the first checkpoint that fails fails the job over, and the job restarts. To the operator this looks like a restart loop, not like a checkpoint timeout. The Flink configuration reference says the setting applies to an IOException on the JobManager, a failure in the asynchronous phase on the TaskManagers and a checkpoint that expires after a timeout.

execution.checkpointing.tolerable-failed-checkpoints: 0   # default

The check is the checkpoint tab of the job in the Flink web UI, which the checkpoint monitoring documentation describes, including its Failed counter. If the Failed counter rises in step with the restarts, or the root exception in the next section names a failed or expired checkpoint, the fix is in the checkpoint path. The guide to Flink checkpoints that fail or time out walks that diagnosis stage by stage.

If checkpoints complete normally between restarts, the failure is elsewhere and the exception history is the next place to look.

Find the root exception

Flink keeps a history of the failures it has handled for each job. The Flink REST API documentation describes the endpoint that returns it, and notes that the same monitoring API is used by Flink’s own dashboard. The API listens on the JobManager, on port 8081 by default.

GET /jobs/:jobid/exceptions

The response holds an exceptionHistory object with a list of entries. Each entry is one root failure and carries these fields:

  • exceptionName, the class of the exception.
  • stacktrace, the full stack trace.
  • taskName, the task where it happened.
  • taskManagerId and endpoint, where that task was running.
  • timestamp, when it happened.
  • failureLabels, key and value labels attached to the failure.
  • concurrentExceptions, other exceptions with the same fields, collected alongside the root failure.

The cause is usually at the bottom of the root entry’s stack trace, under the last “Caused by” line, not in the wrapper exception at the top.

The history is short. The configuration reference sets the maximum number of failures collected per job.

web.exception-history-size: 16   # default

On a fast loop, 16 entries can all be recent repeats, and the first failure that started the loop may already be gone. The truncated flag in the response shows whether entries were left out because of the maxExceptions query parameter. Older scripts that read the top-level root-exception or all-exceptions members still work, but the REST API documentation marks them as deprecated in favor of exceptionHistory.

Label failures across many jobs

Teams that see many loops across many jobs can label failures automatically. The Flink failure enrichers documentation describes a plugin interface that runs on every failure the JobManager reports and attaches labels, such as a label marking the failure as a system error, which the REST API then returns in failureLabels. Enrichers start only when they are listed in the configuration.

jobmanager.failure-enrichers

Tell a transient failure from a permanent one

Flink’s restart strategies do not distinguish transient from permanent exceptions. They count failures. The exception history is enough to tell the two apart in practice, and failure enrichers can add the distinction as a label.

A permanent failure repeats. The same exceptionName and the same stack trace appear on the same taskName at every timestamp, at a steady interval, because each restart replays the same failure. A record that fails deserialization on every read is the Flink form of a Kafka poison pill, and the guide to Kafka deserialization errors covers how such records arise.

A transient failure clears. The history shows different exceptions, or a single failure followed by a long stretch in RUNNING. The task failure recovery documentation notes that delaying a retry helps when the job talks to external systems where connections or pending transactions should time out before the job runs again.

The distinction decides the fix. A transient failure needs a restart strategy that retries with enough delay. A permanent failure needs a code, dependency, data or configuration fix, and a restart strategy that stops retrying so the job fails visibly instead of looping.

F1 Transient and permanent failures in a restart loop, side by side
Transient failure Permanent failure
What you see first A burst of restarts, then a long stretch in RUNNING. Restarts at a steady rhythm that never ends.
Evidence in the exception history Different exceptions, or one entry then a quiet period. The same exceptionName and stack trace at every timestamp.
Usual cause An external system or a lost TaskManager. A code bug, a missing dependency, bad configuration or an unprocessable record.
Fix Keep restarting with a backoff. Fix the cause, redeploy, and bound the restart strategy.
Risk of the fix A long delay lowers availability. A tight bound fails a job a longer wait would have saved.
The transient and permanent split is Factor House's own classification. Restart strategy behavior is from the Apache Flink task failure recovery documentation, and the exception history is from the REST API documentation for GET /jobs/:jobid/exceptions. Task failure recovery documentation

Check which restart strategy is in force

The Flink job scheduling documentation describes the states behind a loop. On a failure the job first switches to FAILING and cancels its running tasks. If the job can be restarted it enters RESTARTING, and once it has been completely restarted it reaches CREATED and runs again. The job reaches FAILED only when it is not restartable, which in practice means the restart strategy has given up. The restart strategy therefore decides whether a loop ever ends.

The task failure recovery documentation sets the default. If checkpointing is disabled, the job does not restart at all. From Flink 1.19, if checkpointing is enabled and no restart strategy is configured, Flink uses the exponential delay strategy. The Flink 1.18 documentation describes a different default for that version, the fixed delay strategy with Integer.MAX_VALUE restart attempts, so a job on an older engine also has no practical limit. The exponential delay defaults that matter for a loop are these.

restart-strategy.type: exponential-delay
restart-strategy.exponential-delay.max-backoff: 1 min
restart-strategy.exponential-delay.reset-backoff-threshold: 1 h
restart-strategy.exponential-delay.attempts-before-reset-backoff: infinite

The documentation also lists an initial backoff of 1 second, a multiplier of 1.5 and a jitter factor of 0.1. With an infinite number of attempts, a job on the default strategy keeps retrying a permanent failure indefinitely. The delay between attempts grows by 1.5 times after each failure until it reaches one minute, and then stays at one minute. The backoff and the attempt counter reset only after the job has run for an hour without a failure. A job that fails every few minutes therefore settles into one restart roughly every minute, and never reaches FAILED.

A restart strategy set on the job overrides the cluster default from the Flink configuration file, so the strategy in force may not be the one in the cluster configuration. The REST API returns the configuration of a job, which is the place to check before changing anything.

GET /jobs/:jobid/config

Failover strategy is a separate setting. It decides which tasks restart, not whether the job restarts. The default restarts the region containing the failed task, which is a group of tasks connected by pipelined data exchanges. As the region failover section describes, it also restarts any region that produces a result the restarted region needs and that is no longer available, and every region that consumes from a restarted region.

jobmanager.execution.failover-strategy: region   # default

Fixes, one change at a time

Pick the metric first. A fix has worked when numRestarts stops climbing and the job stays in RUNNING. Change one setting, then watch that counter and the exception history before changing the next.

Fix the cause of a permanent failure

When the exception history shows the same failure on every attempt, no restart strategy will fix it. Correct the code, dependency, data or configuration named in the root exception, then redeploy. If the job’s state must be kept, redeploy from a savepoint. The Flink savepoint documentation explains the restore options, including allowing the job to skip state that cannot be mapped to the new program.

Bound the restart strategy so a permanent failure stops

A bounded strategy turns an endless loop into a FAILED job that alerting can catch. Unbounded automatic restarts are not specific to Flink. The Strimzi documentation says the automatic restart of failed Kafka Connect connectors and tasks is attempted indefinitely by default, and offers a maxRestarts property to set a limit. The task failure recovery documentation offers three ways to set the bound in Flink.

Fixed delay retries a set number of times with the same delay between attempts, then fails the job.

restart-strategy.type: fixed-delay
restart-strategy.fixed-delay.attempts: 3
restart-strategy.fixed-delay.delay: 10 s

Failure rate keeps restarting until a set number of failures falls within one time interval, then fails the job.

restart-strategy.type: failure-rate
restart-strategy.failure-rate.max-failures-per-interval: 3
restart-strategy.failure-rate.failure-rate-interval: 5 min
restart-strategy.failure-rate.delay: 10 s

Exponential delay keeps its growing backoff but stops after a set number of consecutive attempts. The documentation’s example sets these two values alongside an initial backoff, a multiplier, a maximum backoff and a jitter factor. With them, the job fails if it still encounters exceptions after 8 consecutive retries, and the delay and the retry counter reset when the job runs for 6 minutes without an exception.

restart-strategy.type: exponential-delay
restart-strategy.exponential-delay.attempts-before-reset-backoff: 8
restart-strategy.exponential-delay.reset-backoff-threshold: 6 min

Keep a backoff for transient failures

The task failure recovery documentation says it strongly recommends the exponential delay strategy, because jobs can be retried quickly when exceptions occur occasionally, and avalanches of external components can be avoided when exceptions occur frequently. Its example is many Flink jobs consuming from one Kafka cluster. When that cluster goes down, jobs with a short fixed delay all retry at once, which is likely to cause an avalanche.

The jitter factor adds or subtracts a random amount from each delay so that jobs with identical settings restart at different times. The documentation advises against setting it to 0 in production.

restart-strategy.exponential-delay.jitter-factor: 0.1   # default

A longer delay is the right change only when the history shows a transient cause. For a permanent cause, a longer delay only slows the loop down.

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. Flex is one of several ways to watch a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks. Teams on Ververica Platform can start from the Ververica Platform tool comparison.

Diagnosis. Flex documents these views for each job in its job inspection documentation.

  • The Events tab is a reverse-chronological lifecycle log. Its documented example includes “Triggered CHECKPOINT 91 with status FAILED”, “Flink job entered state RESTARTING” and “Flink job entered state RUNNING”, so a restart and a checkpoint failure read in one place.
  • The Configuration tab shows the settings the job was submitted with. The documented job execution config includes a restart-strategy entry, with the cluster level default restart strategy as its example value.
  • The Checkpoints tab counts Triggered, Completed, In progress, Restored and Failed checkpoints, for the rule-out step above.
  • The Details table lists each job’s State and Activity, the time of the last recorded event.

The documented job State values are RUNNING, FINISHED, CANCELLED and FAILED. The documentation lists no RESTARTING value, so a loop shows in the Events tab, not in the State column.

Fixes. Flex documents these actions.

  • The Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint.
  • Submitting a job from an uploaded JAR lets the operator set a savepoint path, a claim mode, a restore mode and the Allow Non Restored State option, which covers the redeploy after a permanent failure is fixed.

What stays in Flink. The job documentation does not describe the exception history, the root exception or its stack trace, or failure labels, so those stay in the Flink web UI or the /jobs/:jobid/exceptions REST endpoint. The documentation also describes no restart counter and no alert on restarts. Flink exports numRestarts through its metric reporters to an external system, where an alert can be set.

Flex reads the Flink REST API through the endpoint set in its Flink cluster configuration, and its documentation states compatibility with v1 of the Flink REST API. To see each job’s events, configuration and checkpoint history in Flex, point it at the cluster.

FAQ

What is the default restart strategy in Flink?

From Flink 1.19, if checkpointing is enabled and no strategy is configured, Flink uses the exponential delay strategy, as the task failure recovery documentation states. The Flink 1.18 documentation describes the fixed delay strategy with a practically unlimited number of attempts as that version’s default. If checkpointing is disabled, the default is no restart, so the first failure fails the job.

How many times will Flink restart a failing job?

On the default exponential delay strategy, without limit, because the attempts key defaults to infinite.

restart-strategy.exponential-delay.attempts-before-reset-backoff: infinite   # default

Per the task failure recovery documentation, fixed delay defaults to 1 attempt and failure rate to 1 failure per 1 minute interval. Set the attempts or the failure rate explicitly to bound the loop.

Why does a Flink job keep restarting after a checkpoint fails?

Flink tolerates no checkpoint failures by default, so the first checkpoint that expires or fails in its asynchronous phase fails the job over, and the restart strategy restarts it. The key is execution.checkpointing.tolerable-failed-checkpoints, documented in the Flink configuration reference. The guide to Flink checkpoints that fail or time out covers how to find the slow stage.

Where does Flink show why a job failed?

Call GET /jobs/:jobid/exceptions on the JobManager, or open the job in the Flink web UI, which uses the same monitoring API. The root entry’s stack trace names the cause. The history keeps 16 failures per job by default, set by web.exception-history-size.

Related reading