Skip to content

Bound Flink sink retries by time, not attempt count - #206

Merged
karthikk-01 merged 2 commits into
mainfrom
karthik/flink-retry-timeout
Aug 26, 2026
Merged

karthikk-01 merged 2 commits into
mainfrom
karthik/flink-retry-timeout

Conversation

@karthikk-01

Copy link
Copy Markdown
Contributor

Problem

Production Flink jobs restart on transient UNAVAILABLE during query-server rollouts:

io.grpc.StatusRuntimeException: UNAVAILABLE: upstream connect error or disconnect/reset
  before headers ... delayed connect error: Connection refused
  at core.ChalkFeatureUploader.uploadWithRetry(ChalkFeatureUploader.java:105)
Wrapped by: ChalkUploadException: upload_features failed after 4 attempt(s)

The retries worked — after 4 attempt(s) is maxRetries=3 + 1, and UNAVAILABLE is retryable. The budget was just too short. With the default 500ms base backoff:

attempt 0 fails -> sleep 0.5s
attempt 1 fails -> sleep 1.0s
attempt 2 fails -> sleep 2.0s
attempt 3 fails -> give up      total coverage ~3.5s

Connection refused fails fast rather than hanging to the 30s upload timeout, so the real coverage is ~3.5s — shorter than a routine deploy. Channel-level gRPC retry doesn't help: maxAttempts=3 with backoff capped at 0.1s adds ~0.16s.

Change

An attempt count can't express "survive a rollout" — the real budget silently shifts whenever the backoff is tuned. Replaced it with an explicit wall-clock budget:

  • retryTimeout (new, default 120s) — total budget for retrying one batch's transient failures. Sized to outlast a normal query-server rollout.
  • maxRetries — now an opt-in hard cap, unlimited by default, so retryTimeout is the effective bound. Setting it explicitly preserves the old exact-attempt semantics.
  • Timing uses nanoTime so an NTP step can't move the deadline, and the loop breaks when the next sleep would overrun the budget instead of sleeping past it.

Retries still run inline on the task thread, so the budget is also the worst-case subtask stall — documented next to the checkpoint-timeout interaction.

Also included

1cd1803 — the sources/javadoc fix for the chalk-java-shaded publication, currently only on release/1.3.3 (#205). Without it on main, cutting 1.3.4 fails Central validation exactly as 1.3.3 did: the incomplete shaded component sinks the whole deployment bundle, taking chalk-java with it, while CI still reports green. Harmless duplicate if #205 merges first.

Tests

New: stopsRetryingWhenTimeBudgetSpent, timeBudgetAllowsMoreAttemptsThanTheOldFixedCap, explicitMaxRetriesStillCapsBeforeBudget, retryBudgetDefaultsToTimeNotAttempts. All pre-existing retry tests pass unchanged.

The 4 initializationError failures in TestChalkClient* / TestGrpcClient / TestAllClients are pre-existing credential-dependent integration tests — they fail identically on a clean main locally.

Behaviour note

An outage longer than the budget still surfaces as a task failure and a checkpoint restart. That is at-least-once working as intended, not data loss — this change makes the common case (a rollout) stop causing restarts at all.

Central requires sources and javadoc on every component, and both coordinates
upload as a single deployment bundle -- so the incomplete chalk-java-shaded
component failed validation and took chalk-java 1.3.3 down with it. Both publish
runs (8/10, 8/14) uploaded successfully and were rejected server-side afterward,
which is why CI stayed green with nothing on Central.

Reuse the unshaded jars: relocation is a bytecode rewrite, so the sources are
identical and the public ai.chalk.* API is unrelocated. Deferred to afterEvaluate
because the vanniktech plugin registers both tasks after this file evaluates.
Production Flink jobs restarted on transient UNAVAILABLE during query-server
rollouts: "upload_features failed after 4 attempt(s)". The retries themselves
worked, but a fixed maxRetries=3 with a 500ms base backoff only covers
0.5+1+2 = 3.5s -- shorter than a routine deploy, so the budget was spent while
the backend was still coming back and the task failed.

An attempt count can't express "survive a rollout": the real budget shifts
whenever the backoff is tuned. Replace it with an explicit wall-clock budget,
retryTimeout, defaulting to 120s. maxRetries stays as an opt-in hard cap and is
now unlimited by default, so the time budget is the effective bound. Timing uses
nanoTime so an NTP step can't move the deadline, and the loop stops when the
next sleep would overrun it rather than sleeping past it.

Retries still run inline on the task thread, so the budget is also the worst-case
subtask stall; documented alongside the checkpoint-timeout interaction.
@karthikk-01
karthikk-01 requested review from bqin01 and sjmignot August 26, 2026 18:57
@karthikk-01
karthikk-01 merged commit 3025b05 into main Aug 26, 2026
1 check passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants