fix: hold a commit that hits a transient KafkaError and back off its retries - #98
Merged
Merged
Conversation
…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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Closes #62.
Problem
On a transient
KafkaErrorfromconsumer.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:commit_batch_timeout_sec. With the defaults (10s/10s) the retry landed just aftercommit_allgave up, so the offsets were redelivered after the rebalance.commit()calls, each with an ERROR traceback, in a 0.3s flush. It kept going aftercommit_allreturned.close()the queue is no longer read, so re-put offsets were silently left uncommitted.Change
PendingCommits.hold()keeps the failedReadyCommit. The nexttake_ready()merges it with newly ready work per consumer, taking the max offset per partition. Held tasks still count inlen()and are neithertask_done()'d nor uncounted, socommit_all'sjoin()and the backpressure count stay honest.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 feedstransient_errortonote_committed.wait_timeout()wakes the loop for the retry only during a flush or shutdown.commit_allcounts its waiters. When the last one times out (or is cancelled), the committer callsCommitScheduler.release_flush()on its next pass. That happens beforeevaluate(), so a flush opened in the meantime still wins. Overlappingcommit_allcalls from several consumers keep the flush open until all of them give up.CONTEXT.md(the glossary goes from six terms to seven, andAGENTS.mdfollows), 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_sizere-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_timeoutpins that.Tests
Regression tests written first. These fail on
main:test_commit_all_retries_a_transient_error_within_the_flush_timeouttest_commit_retries_back_off_during_a_flushtest_flush_urgency_ends_when_commit_all_times_outtest_close_retries_a_transient_error_before_exitingtest_a_streak_of_transient_errors_logs_one_errorThese pass on
mainand guard behaviour this change must keep:test_one_commit_all_timing_out_keeps_the_flush_for_another_waiterandtest_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) andPendingCommitshold/merge tests. Four existing tests that pinned the re-queue contract were updated.just lintpasses.just test --cov=. --cov-branchgives 227 passed, 100% branch coverage. The unit suite was repeated 6 times under xdist with no flakes.Note for #96
#96 re-keys
PendingCommitsby consumer. Held commits here are already per consumer, so whichever lands second needs only a mechanical rebase.