DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Author: Vaquar Khan
Status: Draft
Discussion Thread:
Target Release: Apache Kafka 4.5 (follow-on to KIP-1191 / AK 4.4)
Depends On: KIP-1191 (DLQ for Share Groups), KIP-1249 (Pause Delivery for Share Groups)
Motivation
KIP-1191 introduces a Dead Letter Queue (DLQ) mechanism for Kafka Share Groups. While this is a significant operational improvement, it introduces two distinct systemic risks at scale that this KIP addresses.
Risk 1: DLQ Flood (Downstream Service Failure)
If a downstream consumer service becomes unavailable or consistently rejects messages, every message in the share group will exhaust its max.delivery.attempts and be routed to the DLQ. In a high-throughput pipeline, this can drain the main topic into the DLQ within minutes, causing:
Loss of message ordering: DLQ records are no longer in the original partition order relative to unprocessed records.
Re-processing complexity: Operators must triage and replay a large DLQ backlog, often under incident pressure.
Masking of root cause: A flood of DLQ writes obscures whether the failure is a poison pill (single bad record) or a systemic outage (downstream service down).
Broker resource pressure: Sustained DLQ writes consume broker I/O, network bandwidth, and coordinator CPU.
Risk 2: Head-of-Line Blocking via DLQ Infrastructure Failure
While group.share.partition.max.record.locks correctly prevents broker memory pressure by bounding in-flight records , hitting that lock ceiling has a critical liveness consequence: fetching operations yield no further records from the primary partition. If the DLQ topic suffers a prolonged metadata or storage outage, the coordinator holds record locks until the limit is exhausted — stalling the primary pipeline entirely.
This re-introduces head-of-line blocking via a secondary infrastructure failure — precisely the problem Share Groups were originally designed to eliminate. A DLQ (a secondary concern) taking down the primary stream is an architectural single point of failure (SPOF). The safer operational posture is to prioritize system liveness over strict DLQ durability during infrastructure degradation.
This KIP proposes:
A configurable rolling-window circuit breaker that automatically pauses a share group when the DLQ write rate exceeds a defined threshold (addressing Risk 1).
A configurable DLQ write timeout with a fall-through to the
ARCHIVEDstate and logged error (addressing Risk 2).
Public Interfaces
New Share Group Configuration Properties
| Config Key | Type | Default | Description |
| share.group.dlq.circuit.breaker.enable | boolean | false | Enables the circuit breaker for this share group. |
| share.group.dlq.circuit.breaker.threshold.percent | int | 20 | Percentage of messages routed to DLQ within the rolling window that triggers the circuit breaker. Range: 1–100. |
| share.group.dlq.circuit.breaker.window.ms | long | 60000 | Rolling window duration (ms) over which the DLQ rate is evaluated. Default: 60 seconds. |
| share.group.dlq.circuit.breaker.min.messages | int | 100 | Minimum number of messages processed within the window before the circuit breaker can trigger. Prevents false positives on low-volume groups. |
| share.group.dlq.circuit.breaker.auto.resume.enable | boolean | false | If true, the group automatically resumes after auto.resume.ms. |
| share.group.dlq.circuit.breaker.auto.resume.ms | long | 300000 | Time (ms) to wait before auto-resuming a paused group. Default: 5 minutes. |
| errors.deadletterqueue.write.timeout.ms | long | 30000 | Maximum time (ms) the coordinator will wait for a DLQ write to complete before falling through to ARCHIVED state with a logged error. Default: 30 seconds. |
| share.group.dlq.strict.isolation | boolean | false | If true, coordinator rejects group startup if the configured DLQ topic is already registered by another share group. |
New Broker Metrics
| Metric | Description |
| kafka.share.group:type=ShareGroupMetrics,name=DlqCircuitBreakerTripped,group={group-id} | Counter incremented each time the circuit breaker trips for a given share group. |
| kafka.share.group:type=ShareGroupMetrics,name=DlqWriteTimeouts,group={group-id} | Counter incremented each time a DLQ write times out and falls through to ARCHIVED. |
| kafka.share.group:type=ShareGroupMetrics,name=DlqTopicSharedByMultipleGroups,topic={dlq-topic} | Gauge emitted at coordinator startup and on group registration when a DLQ topic is shared by more than one share group. No enforcement — observability only. |
Share Group State Extension
The share group state machine gains a new observable reason for the PAUSED state:
PAUSED (reason=DLQ_CIRCUIT_BREAKER_TRIPPED)
This reason is surfaced in:
DescribeGroupsAPI response (newpause-reasonfield)Coordinator logs at
WARNlevelThe
DlqCircuitBreakerTrippedbroker metric above
Proposed Changes
Trigger Condition 1: Rolling-Window Rate Threshold (DLQ Flood)
The Share Group Coordinator tracks a sliding window counter per share group:
For each record disposition event, the coordinator increments either a
deliveredcounter or adlq_writtencounter within the current window bucket.At the end of each window bucket, the coordinator evaluates:
dlq_rate = dlq_written / (delivered + dlq_written)
if dlq_rate >= threshold AND (delivered + dlq_written) >= min.messages:
TRIP circuit breaker → PAUSE group (reason=DLQ_CIRCUIT_BREAKER_TRIPPED)
The coordinator uses the KIP-1249 pause mechanism to halt delivery. No new pause infrastructure is introduced.
The paused state persists until:
An operator explicitly resumes the group via admin tooling, OR
auto.resume.enable=trueandauto.resume.mshas elapsed.
Trigger Condition 2: DLQ Write Timeout (Liveness / HOL-Blocking Prevention)
When the coordinator attempts a DLQ write:
A per-write timer is started, bounded by
errors.deadletterqueue.write.timeout.ms.If the write completes within the timeout: normal flow continues.
If the timeout expires before the write completes:
The record state transitions directly to
ARCHIVED(notARCHIVING).A
WARN-level log entry is emitted with the record offset, partition, and failure reason.The
DlqWriteTimeoutsmetric is incremented.The record lock is released, allowing the primary partition to continue delivery.
DLQ Isolation (Opt-In)
When share.group.dlq.strict.isolation=true:
At group startup, the coordinator checks whether the configured DLQ topic is already registered by another active share group.
If a conflict is detected, group startup is rejected with a descriptive error.
Default is
falseto preserve the intentional flexibility of the current design.
Regardless of the strict.isolation setting, the DlqTopicSharedByMultipleGroups metric is always emitted when a shared DLQ condition is detected.
Interaction with Existing Configs
max.delivery.attempts(KIP-1191): Unchanged. The circuit breaker fires after DLQ writes occur.KIP-1289 (Transactional acknowledgement): A failed transaction that does not result in a committed DLQ write is not counted toward the rate threshold. A timed-out write is counted as a timeout event, not a DLQ write event.
KIP-1249 (Pause delivery): The circuit breaker uses KIP-1249's existing pause mechanism directly.
Admin Tooling
The kafka-share-groups.sh tool is extended with:
Bash
# Describe circuit breaker state for a group
kafka-share-groups.sh --bootstrap-server <host> \
--describe --group <group-id> --circuit-breaker
# Manually resume a circuit-breaker-paused group
kafka-share-groups.sh --bootstrap-server <host> \
--resume --group <group-id> --reason "incident resolved"
# Show DLQ topology across all share groups
kafka-share-groups.sh --bootstrap-server <host> \
--describe --dlq-topology
Compatibility, Deprecation, and Migration Plan
Backward compatibility: All new configs default to off/safe values. Existing share groups are unaffected unless explicitly configured.
Protocol changes: The
DescribeGroupsresponse is extended with an optionalpause-reasonfield. Older clients ignore it gracefully.No deprecations.
Test Plan
Unit tests: Circuit breaker threshold calculation, window bucket rotation, counter reset on resume, timeout fall-through logic.
Integration tests: End-to-end trip and auto-resume under simulated DLQ write failures; timeout fall-through with primary pipeline continuity verification.
Soak tests: 72-hour automated soak with a consumer rejecting 25% of messages; verify group pauses, metrics fire, and no broker resource leak.
HOL-blocking regression tests: Simulate sustained DLQ topic outage; verify primary partition delivery continues after timeout fall-through and lock release.
Chaos tests: Coordinator failover mid-window; verify counter state integrity and no false-trip on recovery.
Isolation tests: Multi-group DLQ sharing scenarios with
strict.isolation=trueandfalse.Upgrade tests: Rolling upgrade from AK 4.4 to AK 4.5; verify no disruption.
Rejected Alternatives
Client-side circuit breaker only: Rejected because client-side logic cannot observe the full DLQ write rate across all consumers in a share group. The coordinator has the authoritative view.
Hard-coded 20% threshold: Rejected in favor of a configurable threshold. Different workloads have different tolerances.
Hard coordinator enforcement of DLQ isolation: Rejected in favor of opt-in
strict.isolationconfig + observability metric, preserving intentional flexibility while giving regulated environments a safety valve.Automatic DLQ topic compaction on circuit breaker trip: Out of scope. DLQ remediation tooling is a separate concern.
Open Questions
Should the circuit breaker window counter be persisted in the share group state topic (for coordinator failover resilience), or is an in-memory counter acceptable given that a failover naturally resets the window?
Should
auto.resume.enabledefault totruefor development environments? (Proposed: no — safer to default to manual resume in all environments.)Should the
pause-reasonfield be added to the existingDescribeGroupsRPC, or introduced as a newDescribeShareGroupStateRPC?Should timed-out DLQ writes (
ARCHIVEDfall-through) count toward the rate threshold circuit breaker, or be tracked as a separate trigger condition only?