Skip to content

fix: reset backoff counter per-record on success to prevent compounding delay after chaos recovery - #274

Merged
dlg99 merged 2 commits into
feat/kafka-supportfrom
feat/kafka-support-heesung
Sep 24, 2026
Merged

dlg99 merged 2 commits into
feat/kafka-supportfrom
feat/kafka-support-heesung

Conversation

@heesung-sohn

Copy link
Copy Markdown

After a chaos/outage window, consecutiveUnavailableException can grow large (e.g. 5+ retries during a 30s partition). Previously the counter only reset after the entire batch completed, so subsequent records in the same batch inherited the high counter and slept for seconds before their first retry — even though CQL was already recovering. In testConnectionFailure this caused the 240s consumer receive timeout to expire before all 3 messages were delivered.

Fix: reset consecutiveUnavailableException to 0 each time a single record's CQL query succeeds in waitForCqlWithRetry (both KafkaCassandraSourceTask and CassandraSource). The end-of-batch reset at poll() is kept as a redundant safety net.

…ng delay after chaos recovery

After a chaos/outage window, consecutiveUnavailableException can grow
large (e.g. 5+ retries during a 30s partition). Previously the counter
only reset after the entire batch completed, so subsequent records in
the same batch inherited the high counter and slept for seconds before
their first retry — even though CQL was already recovering. In
testConnectionFailure this caused the 240s consumer receive timeout to
expire before all 3 messages were delivered.

Fix: reset consecutiveUnavailableException to 0 each time a single
record's CQL query succeeds in waitForCqlWithRetry (both
KafkaCassandraSourceTask and CassandraSource). The end-of-batch reset
at poll() is kept as a redundant safety net.

@dlg99 dlg99 left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM.
Is it covered by tests?

Also, I'd rename consecutiveUnavailableException to consecutiveUnavailableExceptionCount or similar

@heesung-sohn
heesung-sohn force-pushed the feat/kafka-support-heesung branch from 2455316 to e489764 Compare September 24, 2026 18:49
@heesung-sohn

Copy link
Copy Markdown
Author

LGTM. Is it covered by tests?

Also, I'd rename consecutiveUnavailableException to consecutiveUnavailableExceptionCount or similar

added test cases and renamed the var.

@dlg99
dlg99 merged commit b594e71 into feat/kafka-support Sep 24, 2026
30 checks 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