Kafka flush timeout returns 202 with delivery unconfirmed
I maintainer di solito rispondono entro 2 giorni
Nessuno ha ancora preso questa issue.
Valutazione
- Difficoltà
- 3/5
- Tempo stimato
- 1-2 giorni
- Idoneità per principianti
- 70/100
Direzione di ricerca
Inizia in src/writers/writer_kafka.py, in KafkaWriter.write() e nel ramo terminale di timeout del flush che emette un avviso e poi raggiunge il flusso di successo; quindi segui come viene registrato il successo in _write_to_all() per POST /topics/{topic_name}. Esamina adr/002-observability/002-observability.md per il contratto richiesto sulla semantica di timeout rispetto a quella degli errori. Esegui o estendi gli unit test del writer Kafka per il caso di timeout remaining > 0 senza eccezione e verifica che il caso non venga più trattato come riuscito (Kafka non viene segnalato in writers_ok).
Scritto dal modello di indicizzazione a partire dal testo della issue.
Descrizione
Describe the bug
KafkaWriter.write() (src/writers/writer_kafka.py) retries the flush, and when messages are still pending after the last attempt it logs a WARNING and falls through:
if isinstance(remaining, int) and remaining > 0:
logger.warning("Kafka flush timed out with messages still pending.", extra={...})
A pure timeout appends nothing to errors and raises nothing, so the if errors: branch is never reached and write() returns normally. _write_to_all() counts Kafka in writers_ok and the request returns 202, while the message may never have been delivered.
Only a WARNING records it, so an alarm on level = "ERROR" is blind to this by construction — and ADR-002's count(ERROR) == count(5xx) invariant holds here only because the request is never considered failed at all.
Steps to Reproduce
- Configure a topic with a Kafka writer and a reachable broker, so the producer initializes and
produce()succeeds. - Make the broker unable to acknowledge the write while keeping the connection alive — e.g. take the partition leader offline after produce, or set
min.insync.replicasabove the number of live in-sync replicas. POST /topics/{topic_name}with a valid message.flush()returnsremaining > 0on all 3 attempts (KAFKA_FLUSH_RETRIES, 7s timeout each perKAFKA_FLUSH_TIMEOUT) and raises noKafkaException.- Observe the response:
202, withkafkalisted inwriters_ok. - Observe the logs: one
WARNING("Kafka flush timed out with messages still pending."), zeroERROR.
Expected state
A message whose delivery was never confirmed must not be reported to the caller as written.
202 is meant to mean "accepted by every configured sink". Either the response reflects that the Kafka write did not complete, or writers_ok stops listing a sink whose delivery is unconfirmed — but a caller must not be told the message landed when the service does not know that it did.
Impact / Severity
High
Attachments / Evidence
src/writers/writer_kafka.py — flush loop, terminal timeout branch, and the if errors: gate that the timeout path never reaches:
# Warn if messages still pending after retries
if isinstance(remaining, int) and remaining > 0:
logger.warning(
"Kafka flush timed out with messages still pending.",
extra={"pending_messages": remaining, "flush_timeout_sec": _KAFKA_FLUSH_TIMEOUT_SEC},
)
duration_ms = round((time.perf_counter() - started_at) * 1000, 2)
if errors: # <- empty on a pure timeout
...
raise WriteError(failure_text)
logger.debug("Kafka accepted the message.", ...) # <- reached instead
errors is appended to only by delivery_report (on a delivery error) and by the two except KafkaException blocks. A flush that simply does not drain in time hits none of them.
Related / References
Two options, and the choice is a product decision rather than a cleanup:
- Treat a terminal flush timeout as a
WriteError. Correct on the contract, but it turns a current202into a500, so it is a caller-visible behaviour change. - Keep the
202and alarm on this specificWARNING, treating unconfirmed delivery as an operational signal rather than a request failure.
Option 1 is the honest one if 202 is meant to mean "accepted by every configured sink". Whichever is chosen, the decision belongs in ADR-002 and writers_ok must stop reporting an unconfirmed sink as OK.
Acceptance:
- Decision recorded in ADR-002 §Logging strategy.
writers_okno longer reports a sink whose delivery was never confirmed.- Unit test covers flush timeout with
remaining > 0and no raised exception.
Related: ADR-002 (adr/002-observability/002-observability.md), #193, PR #204, #219 — the other invariant gap found in the same review.
- Lingua principale
- Python
- Stelle
- 4
- Fork
- 0
- Merge medio
- 1g 20h
- PR unite (30g)
- 8
Preparare l'ambiente
- Include un Dockerfile o un file Docker Compose
- Ha un modello di pull request
- Nessuna guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di AbsaOSS/EventGate
-
refactoring type:tech-debt
Difficoltà 2/5 1-3 ore Idoneità per principianti 84/100
I maintainer di solito rispondono entro 2 giorni
-
enhancement
Difficoltà 2/5 1-3 ore Idoneità per principianti 70/100
I maintainer di solito rispondono entro 2 giorni
-
Make Writes IdempotentApertaenhancement
Difficoltà 5/5 Più di una settimana Idoneità per principianti 32/100
I maintainer di solito rispondono entro 2 giorni
-
infrastructure type:tech-debt
Difficoltà 3/5 1-2 giorni Idoneità per principianti 70/100
I maintainer di solito rispondono entro 2 giorni
-
refactoring type:tech-debt
Difficoltà 3/5 1-2 giorni Idoneità per principianti 71/100
I maintainer di solito rispondono entro 2 giorni
Tutte le issue di AbsaOSS/EventGate
Issue simili
-
comp/tools duplicate P2 sweeper:risk-compatibility tool/mcp type/bug
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 88/100
NousResearch/hermes-agent#132042 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 78/100
deepset-ai/haystack#13092 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 86/100
EverMind-AI/Raven#845 ·
I maintainer di solito rispondono entro 1 giorno
-
[Feature] 移除「切换到旧版知识库」入口Aperta
Difficoltà 2/5 1-3 ore Idoneità per principianti 72/100
AstrBotDevs/AstrBot#10340 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 78/100
BasedHardware/omi#20401 · 1 commento ·
I maintainer di solito rispondono entro 1 giorno