Flink

Mid-level

How do I configure Flink restart strategies so a failure doesn't crash-loop forever?

Flink's out-of-the-box restart behavior is built for a job that occasionally hiccups, not one that's actually broken. Left unconfigured, a genuinely bad deploy just restarts forever.

Read in:
failure-rate exponential-delay fixed-delay (default)
Restart strategies roughly ordered by how well they suit a production job with real failure modes - not a measured benchmark.

When a Flink job fails, its restart strategy decides what happens next: restart immediately, restart after a delay, give up, or something in between. The cluster-wide default is fixed-delay with a small number of attempts - fine for a flaky external call, actively bad for a bug that fails on every single restart.

The three strategies, and what each one is actually for

fixed-delay retries a fixed number of times with a constant delay between attempts - simple, and fine for transient failures, but a job that's broken for a real reason (bad code in the latest deploy, a schema it can no longer parse) just burns through its attempts and then stops, or restarts forever if the attempt count is set too high. exponential-delay backs off further after each failure, which helps against a downstream dependency that's struggling rather than a job that's actually wrong. failure-rate is the one built for production: it tracks failures per time window and only gives up once that rate is exceeded, so a single blip doesn't kill the job but a job that's failing continuously does eventually stop instead of looping.

Configuring failure-rate

# flink-conf.yaml 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: 30 s

That configuration allows up to 3 failures in any rolling 5-minute window, waiting 30 seconds between restart attempts - past 3 failures in 5 minutes, the job stops entirely rather than continuing to restart. The right numbers depend on how expensive a restart actually is for this job (state size, checkpoint recovery time) - there's no universal default worth copying without adjusting.

What a crash-loop actually costs

A job stuck restarting isn't free even while it's "handling" the failure automatically: every restart attempt reloads state from the last checkpoint, which for a job with meaningful state size is real I/O and real time, repeated on every cycle. A crash-loop with a short fixed delay and a high attempt count can generate more checkpoint-storage read traffic than the job's actual steady-state workload - worth checking before assuming a restarting-but-not-alerting job is harmless.

What this doesn't replace: a restart strategy controls what Flink does after a failure - it doesn't diagnose why the job failed in the first place. opti-pipe's rule engine reads a job's exported metrics for exactly that: parallelism/backpressure and checkpoint-duration signals that point at a root cause, restart strategy aside.

See what this looks like on your own pipeline.

Upload a real Spark event log, dbt run_results.json, or Flink metrics export and get concrete, approve-before-apply recommendations back - not another rule of thumb.