Skip to content

[fix][broker] Prevent slow Key_Shared sockets from stalling other consumers - #26635

Merged
lhotari merged 1 commit into
masterfrom
lh-fix-key-shared-write-progress
Sep 18, 2026
Merged

lhotari merged 1 commit into
masterfrom
lh-fix-key-shared-write-progress

Conversation

@lhotari

@lhotari lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member

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.

Socket-pressure diagnostic Wait for all consumer writes Continue after any consumer write completes
Healthy-consumer messages delivered while the other socket was paused 16 of 1,024, then stalled All 1,024
Healthy-consumer progress without new FLOW or acknowledgments Blocked by the unrelated pending write Continues independently

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:

Metric Wait for all writes Continue after any completed write Change
Steady dispatch 100,320 msg/s 100,157 msg/s −0.16%
Estimated allocation per dispatched message 1,961 B 2,080 B +6.1%
CPU samples per million dispatched messages 5,557 5,623 +1.2%

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 quickCheck passed.

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:

  • 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

Any consumer batch's write completion can schedule the existing conflated dispatcher loop, instead of waiting for all consumer writes in the read batch.

@lhotari
lhotari added this pull request to stack #26636 September 17, 2026 14:55
Base automatically changed from lh-improve-dispatch-backpressure to master September 17, 2026 16:13
@merlimat
merlimat force-pushed the lh-fix-key-shared-write-progress branch from afd0f52 to 11c22af Compare September 17, 2026 16:13
@lhotari
lhotari merged commit 95d170e into master Sep 18, 2026
43 checks passed
@lhotari
lhotari deleted the lh-fix-key-shared-write-progress 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.

3 participants