Skip to content

[improve][broker] Resume Shared dispatch when consumer channels become writable - #26630

Merged
merlimat merged 3 commits into
masterfrom
lh-improve-dispatch-backpressure
Sep 17, 2026
Merged

merlimat merged 3 commits into
masterfrom
lh-improve-dispatch-backpressure

Conversation

@lhotari

@lhotari lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member

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.

Subscription Pending outbound bytes before Pending outbound bytes after Reduction Drain time before Drain time after
Shared 17,490,153 135,553 99.22% 9,697 ms 9,859 ms
Key_Shared 431,837 313,285 27.45% 10,211 ms 10,143 ms

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

  • Notify connected consumers when their channel becomes writable, and forward the notification to the subscription dispatcher.
  • Coalesce nested writable notifications onto a later event-loop task. Recheck writability for each consumer because a notification can itself cause writes and change the channel state.
  • Resume the existing conflated Shared/Key_Shared read loop without taking subscription or dispatcher monitors on the Netty notification stack.
  • Exclude unwritable consumers from selection and stop polling when no connection can accept a batch.
  • Once a consumer is selected, enqueue its bounded batch even if writability changes during the write. Preserve remaining entry positions for replay rather than rewinding entries already sent to other consumers.
  • Recheck receive permits, pending reads and dispatch-rate quotas when processing a notification. Existing retry timers cannot create a duplicate pending read.

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:

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Netty writable events schedule the existing dispatcher read loop. Notifications do not perform inline dispatch or wait for dispatcher/subscription monitors.

@lhotari
lhotari added this pull request to stack #26636 September 17, 2026 14:55
@merlimat
merlimat merged commit e0a166e into master Sep 17, 2026
43 checks passed
@lhotari
lhotari deleted the lh-improve-dispatch-backpressure branch October 1, 2026 00:02
@lhotari lhotari added this to the 5.0.0 milestone Oct 1, 2026
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