Repository navigation
[improve][broker] Resume Shared dispatch when consumer channels become writable - #26630
Merged
Merged
Conversation
…e writable Assisted-by: Codex
1 of 10 tasks
lhotari
added this pull request to stack #26636
September 17, 2026 14:55
merlimat
approved these changes
Sep 17, 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.
Fixes the Shared/Key_Shared channel-backpressure and missing writable-wakeup behavior described in #24926. The Exclusive/Failover behavior also discussed in that issue is outside this PR's scope.
Motivation
A Shared or Key_Shared consumer can still have receive permits when its connection's outbound buffer is full. The dispatcher should stop assigning new work to that connection and resume promptly when the socket can accept more data. Polling through small reads or retry backoffs can instead accumulate queued writes or delay recovery after the connection becomes writable.
A writable channel event provides the wakeup directly. It also allows healthy consumers to remain eligible without continuing to feed a blocked socket.
Observed impact under socket backpressure
A recorded diagnostic compared the same test fixture with and without these changes: one producer and one consumer on a persistent topic, 2,048 unbatched 8 KiB messages, and 8,192 initial consumer receive permits. The consumer's socket reads were paused with small TCP send/receive buffers. After the broker channel became unwritable, the test held the pause for one second and sampled Netty's pending outbound bytes before resuming client reads.
The Shared dispatcher previously consumed permits for all 2,048 messages despite the unwritable socket. With the change, only 22 permits had been consumed at the sample: the dispatcher finished its selected batch and stopped assigning further entries to the blocked connection. This substantially reduces outbound buffering while the consumer cannot read.
Key_Shared's existing send-completion coordination already limited buffering to less than one read batch in this fixture. Its smaller byte difference is timing-dependent and does not establish a comparable improvement. Exact queue sizes depend on batch/read timing; this is not a hard byte limit.
Both subscription types delivered all messages in order after socket reads resumed, without additional publishing, acknowledgments or FLOW being needed; messages were acknowledged afterward. Both baseline and changed runs passed without failures, skips or retries. These previously collected diagnostic results demonstrate reduced buffering for Shared, not a recovery-time or throughput improvement. Drain times were similar in this deliberately constrained TCP setup, and these are not measurements of the final revised regression-test fixture.
Modifications
This changes the Shared and Key_Shared dispatchers, without changing channel watermark defaults or the single-active and classic dispatcher implementations.
Verifying this change
Added tests cover writable recovery without new FLOW or publish traffic, actual paused client sockets, writability changing during notification, recursive write/flush callbacks, fixed and exponential retry timers, duplicate-read prevention, subscription-monitor independence and exhausted dispatch quotas. Socket tests verify complete delivery and ordering after reads resume.
59 scoped test cases passed without failures, skips or retries. Required Spotless and Checkstyle checks passed.
Does this pull request potentially affect one of the following parts:
Netty writable events schedule the existing dispatcher read loop. Notifications do not perform inline dispatch or wait for dispatcher/subscription monitors.