Skip to content

Add event coalescing for slow-client send buffers - #33

Merged
ralscha merged 2 commits into
ralscha:mainfrom
ABin-Huang:feat/coalesce-overflow-policy
Sep 19, 2026
Merged

ralscha merged 2 commits into
ralscha:mainfrom
ABin-Huang:feat/coalesce-overflow-policy

Conversation

@ABin-Huang

Copy link
Copy Markdown
Contributor

Motivation

The per-client send buffer (PR #32) protects slow clients from head-of-line blocking, but a slow client that receives a high-frequency stream of small events — for example LLM token streaming — can still fill its buffer quickly. Each event is written to the connection individually, so the write count and network frame count grow with the token rate.

What this PR adds

An opt-in event coalescer that merges consecutive buffered events of the same client into a single SSE frame before they are written:

  • EventCoalescer — functional interface; return null for pairs that must be sent separately.
  • DefaultEventCoalescer — merges plain data events (same event name, no id / retry / comment / JSON view) by joining their data with a line break. The SSE protocol encodes that as multiple data: lines of one event, so a client that concatenates data lines receives the original data in order.
  • ClientSendBuffer — when a coalescer is configured, the dispatcher drains a bounded batch (up to 32 events) from the queue and merges adjacent events. Fewer, larger writes reduce system call and network overhead and keep the per-client queue from filling up. Heartbeats are never merged.
  • SseEventBusConfigurer.eventCoalescer() — default null (coalescing disabled, behaviour unchanged); only used when the send buffer is enabled.
  • Metrics — new sse.eventbus.client.buffer.coalesced.events counter.

Example

@Configuration
@EnableSseEventBus
public class MyConfig implements SseEventBusConfigurer {

    @Override
    public int clientSendBufferCapacity() {
        return 128;
    }

    @Override
    public EventCoalescer eventCoalescer() {
        return new DefaultEventCoalescer();  // or a custom coalescer
    }
}

Tests

  • Unit tests: adjacent events merge into one write; non-mergeable pairs are sent separately; the metric is recorded.
  • Integration test (SseEventBusBackpressureCoalesceTest): 50 published events result in fewer than 50 writes (afterEventSent), the coalesced counter is > 0 and the client stays connected.

Local verification with the CI command:

./mvnw -B -ntp -Perror-prone verify

→ BUILD SUCCESS, 162 tests, 0 failures (NullAway, spring-javaformat, license checks included).

Introduce an opt-in EventCoalescer that merges consecutive buffered
events of the same client into a single SSE frame before they are
written to the connection.

- EventCoalescer: functional interface, returns null for pairs that
  must be sent separately.
- DefaultEventCoalescer: merges plain data events (same event name, no
  id/retry/comment/jsonView) by joining their data with a line break;
  the SSE protocol encodes that as multiple data: lines of one event.
- ClientSendBuffer: when a coalescer is configured, the dispatcher
  drains a bounded batch (up to 32 events) from the queue and merges
  adjacent events, so a high-frequency stream of small events (for
  example LLM token streaming) produces fewer, larger writes and keeps
  the per-client queue from filling up. Heartbeats are never merged.
- SseEventBusConfigurer.eventCoalescer(): default null (disabled,
  behaviour unchanged); only used when the send buffer is enabled.
- Metrics: new 'coalesced.events' counter.
- Tests: unit tests for merging, non-mergeable pairs and the metric;
  integration test showing fewer writes than published events.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

There is a correctness risk in DefaultEventCoalescer’s data handling (String.valueOf fallback) and the new integration test can pass without actually verifying full delivery completion.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 Medium severity · 1 Low severity

Open (3)
What changed in this PR

This PR introduces an opt-in event coalescing mechanism for the per-client send buffer to reduce write/system-call overhead for slow clients receiving high-frequency small SSE events.

Changes:

  • Adds EventCoalescer SPI and DefaultEventCoalescer for merging consecutive compatible events into a single SSE frame.
  • Updates ClientSendBuffer dispatch loop to drain bounded batches and coalesce adjacent events (with a new Micrometer counter).
  • Extends configuration and tests to cover coalescing behavior and metrics.
File Description
src/​main/​java/​ch/​rasc/​sse/​eventbus/​EventCoalescer.java New functional interface for merging consecutive buffered events.
src/​main/​java/​ch/​rasc/​sse/​eventbus/​DefaultEventCoalescer.java Default coalescing implementation for “plain data” events.
src/​main/​java/​ch/​rasc/​sse/​eventbus/​ClientSendBuffer.java Implements bounded batch draining + adjacent-event coalescing in dispatcher thread.
src/​main/​java/​ch/​rasc/​sse/​eventbus/​SseEventBus.java Wires configurer-provided coalescer into per-client send buffer creation.
src/​main/​java/​ch/​rasc/​sse/​eventbus/​SseBackpressureMetrics.java Adds coalesced.events counter and recording API.
src/​main/​java/​ch/​rasc/​sse/​eventbus/​config/​SseEventBusConfigurer.java Adds eventCoalescer() hook (default disabled).
src/​main/​java/​ch/​rasc/​sse/​eventbus/​ClientEvent.java Exposes getConvertedValue() to support safe coalescing decisions.
src/​test/​java/​ch/​rasc/​sse/​eventbus/​ClientSendBufferTest.java Updates constructor signature usage and adds unit tests for coalescing + metrics.
src/​test/​java/​ch/​rasc/​sse/​eventbus/​SseEventBusBackpressureCoalesceTest.java Adds integration coverage for coalescing and metric emission.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment on lines +58 to +65
private static @Nullable String eventData(ClientEvent event) {
String converted = event.getConvertedValue();
if (converted != null) {
return converted;
}
Object data = event.getSseEvent().data();
return data != null ? String.valueOf(data) : null;
}
Comment on lines +115 to +122
// the dispatcher merged adjacent events into single SSE frames
await().atMost(Duration.ofSeconds(5))
.untilAsserted(() -> assertThat(
REGISTRY.get("sse.eventbus.client.buffer.coalesced.events").counter().count())
.isGreaterThan(0));
// far fewer writes than published events, and at least one write happened
await().atMost(Duration.ofSeconds(5)).untilAsserted(() -> assertThat(SENT_NOTIFICATIONS.get()).isGreaterThan(0));
assertThat(SENT_NOTIFICATIONS.get()).isLessThan(50);
Comment on lines +282 to +283
List<BufferedEvent> rest = new ArrayList<>();
this.queue.drainTo(rest, MAX_COALESCE_BATCH - 1);
Address the first review round:

- DefaultEventCoalescer no longer stringifies unconverted payloads: the
  merged data is taken from the converted value when present, or directly
  from a String payload (which the send path writes as-is); unconverted
  non-String payloads are never merged.
- SseEventBusBackpressureCoalesceTest waits until the per-client queue
  is drained (queue.size gauge == 0) before comparing the number of
  writes against the number of published events, so the assertion cannot
  pass on a partially delivered batch.
- ClientSendBuffer pre-sizes the batch list in the hot path.
@ABin-Huang

Copy link
Copy Markdown
Contributor Author

All three findings are addressed in the new commit:

1. No stringifying of unconverted payloads — DefaultEventCoalescer now takes the merged data from the converted value when present, or directly from a String payload (which the send path writes as-is, keeping the converted value null). Unconverted non-String payloads are never merged, so no toString() fallback can diverge from the DataObjectConverter output.

2. Await full delivery in the integration test — the test now waits until the per-client queue is drained (sse.eventbus.client.buffer.queue.size gauge == 0) before asserting that writes < 50, so it cannot pass on a partially delivered batch.

3. Pre-sized batch list — new ArrayList<>(MAX_COALESCE_BATCH - 1) in the dispatch hot path.

Local verification with the CI command:

./mvnw -B -ntp -Perror-prone verify

→ BUILD SUCCESS, 162 tests, 0 failures.

@copilot-pull-request-reviewer could you re-review?

@ralscha
ralscha merged commit 87e2bd8 into ralscha:main Sep 19, 2026
5 checks passed
@ralscha

ralscha commented Sep 19, 2026

Copy link
Copy Markdown
Owner

Released with version 3.3.1

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