Skip to content

fix: hold a commit that hits a transient KafkaError and back off its retries - #98

Merged
lesnik512 merged 1 commit into
mainfrom
fix/commit-retry-backoff
Oct 4, 2026
Merged

lesnik512 merged 1 commit into
mainfrom
fix/commit-retry-backoff

Conversation

@lesnik512

Copy link
Copy Markdown
Member

Closes #62.

Problem

On a transient KafkaError from consumer.commit(), the committer re-put the batch's tasks onto its intake queue. Depending on what else was pending, that caused one of three failures:

  • Lost urgency. If the re-queue emptied pending work, the flush ended and the re-absorbed tasks waited a full commit_batch_timeout_sec. With the defaults (10s/10s) the retry landed just after commit_all gave up, so the offsets were redelivered after the rebalance.
  • Hot loop. If another partition still had a user task in flight, the flush stayed open and every re-put woke the loop to commit again: about 100 commit() calls, each with an ERROR traceback, in a 0.3s flush. It kept going after commit_all returned.
  • Dropped on shutdown. During close() the queue is no longer read, so re-put offsets were silently left uncommitted.

Change

  • A failed commit stays pending (ADR 0004). PendingCommits.hold() keeps the failed ReadyCommit. The next take_ready() merges it with newly ready work per consumer, taking the max offset per partition. Held tasks still count in len() and are neither task_done()'d nor uncounted, so commit_all's join() and the backpressure count stay honest.
  • Backoff in CommitScheduler. Retries are spaced 0.1s, doubling, capped at 2s (COMMIT_RETRY_BACKOFF_BASE_SEC / COMMIT_RETRY_BACKOFF_MAX_SEC, not public options). There is one streak per committer; a round without a transient error resets it. The scheduler stays clock-free: the driver feeds transient_error to note_committed. wait_timeout() wakes the loop for the retry only during a flush or shutdown.
  • Flush release. commit_all counts its waiters. When the last one times out (or is cancelled), the committer calls CommitScheduler.release_flush() on its next pass. That happens before evaluate(), so a flush opened in the meantime still wins. Overlapping commit_all calls from several consumers keep the flush open until all of them give up.
  • Logging. The first failure of a streak logs ERROR with a traceback. Later failures log WARNING without one. The first round after a streak with no transient error logs INFO.
  • Docs. A Flush entry in CONTEXT.md (the glossary goes from six terms to seven, and AGENTS.md follows), plus ADR 0004.

One deviation from the brief

The brief said backoff applies "while a flush is active". In this PR the retry-not-before gate holds back every commit trigger, batch-size included, not just flush urgency. Without that, a held batch large enough to reach commit_batch_size re-commits on every loop wake-up outside a flush, which is the same hot loop through a different trigger. Outside a flush the gate is shorter than the batch timeout and never wakes the loop itself, so the common path's cadence is unchanged. test_transient_error_outside_a_flush_retries_after_the_batch_timeout pins that.

Tests

Regression tests written first. These fail on main:

  • test_commit_all_retries_a_transient_error_within_the_flush_timeout
  • test_commit_retries_back_off_during_a_flush
  • test_flush_urgency_ends_when_commit_all_times_out
  • test_close_retries_a_transient_error_before_exiting
  • test_a_streak_of_transient_errors_logs_one_error

These pass on main and guard behaviour this change must keep: test_one_commit_all_timing_out_keeps_the_flush_for_another_waiter and test_transient_error_outside_a_flush_retries_after_the_batch_timeout. There are also clock-free scheduler tests (gating, doubling to the cap, reset, shutdown gating, release ordering) and PendingCommits hold/merge tests. Four existing tests that pinned the re-queue contract were updated.

just lint passes. just test --cov=. --cov-branch gives 227 passed, 100% branch coverage. The unit suite was repeated 6 times under xdist with no flakes.

Note for #96

#96 re-keys PendingCommits by consumer. Held commits here are already per consumer, so whichever lands second needs only a mechanical rebase.

…retries

A transient commit error re-queued the batch. That either ended the flush
(commit_all timed out at the default 10s/10s and offsets were redelivered),
hot-looped commit() with a traceback per attempt while other partitions had
work in flight, or, during close(), left the offsets uncommitted.

The failed commit now stays in PendingCommits and merges into the next
take_ready(); CommitScheduler spaces retries 0.1s doubling to 2s; a flush is
released once every commit_all waiting on it has given up; a streak logs one
ERROR, then WARNINGs, then INFO on recovery.

Closes #62
@lesnik512
lesnik512 merged commit 092d31b into main Oct 4, 2026
11 checks passed
@lesnik512
lesnik512 deleted the fix/commit-retry-backoff branch October 4, 2026 08:52
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.

commit_all loses flush urgency when a commit hits a transient KafkaError

1 participant