Repository navigation
[fix][broker] Prevent slow Key_Shared sockets from stalling other consumers - #26635
Merged
Merged
Conversation
…sumers Assisted-by: Codex
merlimat
force-pushed
the
lh-fix-key-shared-write-progress
branch
from
September 17, 2026 16:13
afd0f52 to
11c22af
Compare
nodece
approved these changes
Sep 18, 2026
dao-jun
approved these changes
Sep 18, 2026
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.
Builds on #26630, which prevents dispatch to unwritable consumer channels and resumes dispatch when they become writable.
Motivation
Key_Shared dispatch groups a read batch by consumer and waits for every consumer's socket write to complete before requesting another read. If one consumer stops reading from its socket, that shared completion barrier can also stall healthy consumers assigned to other key ranges, even when they have receive permits and writable connections.
Progress should not depend on unrelated FLOW, acknowledgment or channel events to wake the dispatcher. A completed write for one consumer is enough to consider more work for eligible consumers.
Observed benefit and performance tradeoffs
A recorded two-consumer socket diagnostic paused one consumer's socket while the other remained writable. Healthy-consumer messages were 128 bytes and slow-consumer messages were 8 KiB, keeping healthy writes below the channel watermark so an unrelated writable notification could not mask the completion barrier. The healthy consumer had receive permits and sent no new FLOW or acknowledgments during this phase.
Restoring the old completion barrier reproduced the stall with the same diagnostic fixture. This is a progress/correctness benefit, not a 64-fold throughput claim. The final regression test uses a larger slow-consumer stream to ensure the socket stays backpressured, as described below. Configured look-ahead limits still bound how far the dispatcher can advance past blocked keys.
To check whether this behavior change regresses ordinary workloads without artificially delayed or paused consumer sockets, a separate recorded comparison used 500 isolated producers and 20 consumers on one Key_Shared subscription, random keys, unbatched 128-byte messages, and a 100,000 msg/s aggregate producer rate cap. This is a regression check for normal consumer progress, separate from the slow-socket test above:
Both runs received all 12 million messages with no failed acknowledgments. Steady intervals exclude the first 20 seconds after all producers connect and end before producers disconnect; allocation and CPU estimates use matched 30-second steady windows. The producer rate cap does not imply low utilization. These sampled results from short runs with varying thermal conditions establish no throughput gain and flag a possible allocation cost; they do not isolate its cause. The reason for this change is to prevent a slow consumer's pending socket write from stalling healthy consumers.
Modifications
Remove the per-read counter that waits for all consumer writes. Each consumer batch's asynchronous write completion requests the existing conflated read loop, which rechecks channel writability, receive permits, pending reads and dispatch limits.
The dispatcher still serializes read and dispatch decisions, and each channel preserves its write order. Key ownership, replay and draining-hash rules are unchanged. Configured look-ahead limits still bound how far dispatch can advance past a blocked key range. This removes a cross-consumer completion dependency without introducing a new executor handoff or an eager read loop.
Verifying this change
Added a deterministic test with one incomplete consumer write and one completed write. The completed write must request the next read without waiting for the other consumer.
Added a test using two separate client sockets on one Key_Shared subscription. While one socket is paused, the healthy consumer receives all 1,024 messages for its key range without sending new FLOW or acknowledgments. The healthy messages are interleaved throughout a larger slow-consumer stream so that they cannot all complete before the slow socket fills. After that socket resumes, its 4,096 messages are delivered in order.
Restoring the previous completion barrier makes the healthy-consumer test fail while the slow write is pending.
62 scoped test cases passed without failures, skips or retries. Spotless, Checkstyle and
quickCheckpassed.The benefit is continued healthy-consumer progress under socket backpressure; maximum-throughput improvements are not claimed.
Does this pull request potentially affect one of the following parts:
Any consumer batch's write completion can schedule the existing conflated dispatcher loop, instead of waiting for all consumer writes in the read batch.