What concurrency, backpressure, and buffer-management pitfalls must you handle when broadcasting to many reactive WebSocket sessions, and how do you address them?
answer
- One session = one ordered send(); never concurrent/multi-thread writes
- Fan-out via Sinks.many().multicast(); each session own send()
- Slow consumer: onBackpressureBuffer(bounded)/drop/latest + publishOn
- Map messages per-session; don't share one WebSocketMessage/DataBuffer
- No blocking on event loop; boundedElastic for blocking; broker for multi-instance
basics
~20 sA session supports a single ordered write stream, so never call send() concurrently or write from multiple threads. For broadcast, use a shared Sinks.many() and let each session subscribe with its own send(), applying backpressure strategy (buffer/drop) per slow client, and mind pooled DataBuffer retain/release when reusing payloads.
solid answer
~50 sEach WebSocketSession is effectively single-writer: you get one send() subscription and must merge all outbound sources into one ordered Flux — concurrent send() calls or cross-thread writes corrupt framing. For fan-out, a common pattern is a hot source such as Reactor's Sinks.many().multicast() (or a Flux shared via .share()/replay), where each connected session does session.send(sharedFlux.map(session::textMessage)). The hard problem is slow consumers: a single slow client must not stall the whole broadcast, so each subscriber needs an isolated backpressure strategy — onBackpressureBuffer with a bounded queue plus an overflow policy, or onBackpressureDrop / onBackpressureLatest, and possibly publishOn to decouple. Also mind pooled DataBuffers: don't build one shared WebSocketMessage and send it to many sessions, because each session's DataBufferFactory and buffer lifecycle differ — map per session, or retain/release explicitly. Finally, size Reactor Netty worker threads and monitor for slow-client memory growth.
code
java · 29 linesimport org.springframework.web.reactive.socket.*;
import reactor.core.publisher.*;
import reactor.core.scheduler.Schedulers;
import reactor.util.concurrent.Queues;
public class BroadcastHandler implements WebSocketHandler {
// Shared hot source; publishers call sink.tryEmitNext(...) elsewhere.
private final Sinks.Many<String> sink =
Sinks.many().multicast().onBackpressureBuffer(Queues.SMALL_BUFFER_SIZE, false);
private final Flux<String> broadcast = sink.asFlux();
public void publish(String payload) { sink.tryEmitNext(payload); }
@Override
public Mono<Void> handle(WebSocketSession session) {
Flux<WebSocketMessage> outbound = broadcast
// Per-subscriber isolation so one slow client can't stall others or OOM:
.onBackpressureBuffer(256,
dropped -> { /* metric */ },
reactor.core.publisher.BufferOverflowStrategy.DROP_OLDEST)
.publishOn(Schedulers.parallel())
// Fresh message per session (correct DataBufferFactory, no shared buffer):
.map(session::textMessage);
return session.send(outbound) // single ordered writer
.and(session.receive().then()); // drain inbound + detect close
}
}go deeper
Understand a session should have a single send() and not be written from many threads.
Use a shared Sinks source and map messages per session for broadcast.
Apply per-subscriber backpressure strategies and avoid blocking the event loop.
Reason end-to-end about slow-consumer isolation, pooled-buffer correctness, and the single-instance-Sinks vs external-broker scaling boundary.
## Constraint 1 — a session is single-writer The reactive `WebSocketSession` is designed around **one** outbound subscription. WebSocket framing is ordered and stateful, so: - **Do not call `session.send(...)` more than once concurrently**, and do not write to the same session from multiple threads. Interleaved writes corrupt frame boundaries. - Combine every outbound source (business events, heartbeats, control messages) into **one** `Flux<WebSocketMessage>` via `Flux.merge`, `Flux.mergeWith`, or a per-session `Sinks.many()` you push into, then a single `send()`. ## Constraint 2 — fan-out / broadcast topology To push one event to N sessions you need a **shared hot source**. Idiomatic Reactor: ```java Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer(); Flux<String> broadcast = sink.asFlux(); // publish an event: sink.tryEmitNext(payload); ``` Each session subscribes independently in its handler: ```java return session.send(broadcast.map(session::textMessage)) .and(session.receive().then()); ``` Alternatives: a shared `Flux....share()` / `.publish().autoConnect()`, or an external broker. (Structured pub/sub via STOMP/RSocket is a different leaf; here we mean raw fan-out.) ## Constraint 3 — the slow consumer problem (the crux) Multicasting to many sessions means each subscriber consumes at its own pace. One slow/stalled client (slow network, paused tab) creates **backpressure** that, if unmanaged, either (a) blocks the shared source for everyone, or (b) buffers unboundedly and OOMs the server. Per-subscriber isolation is essential: - `onBackpressureBuffer(maxSize, overflowStrategy)` — bounded per-client queue; on overflow choose `DROP_LATEST`, `DROP_OLDEST`, or `ERROR` (which closes just that session). Bounding is what protects heap. - `onBackpressureDrop(...)` / `onBackpressureLatest()` — for state-snapshot streams where stale frames are worthless (e.g., live price ticks) — keep only the newest. - `publishOn(scheduler)` — decouple the slow write from the shared upstream so it doesn't stall siblings; combined with a bounded buffer this isolates a laggard. - With `Sinks.many().multicast().onBackpressureBuffer()`, understand that the sink's buffer is shared until subscribers attach; per-subscriber protection still comes from operators each session adds. - Consider closing chronically slow clients (`session.close(CloseStatus.POLICY_VIOLATION)` after a buffer-overflow error) rather than degrading the whole system. ## Constraint 4 — pooled DataBuffer / message reuse `WebSocketMessage` payloads are backed by pooled `DataBuffer`s tied to a `DataBufferFactory`. Two traps: - **Don't build one `WebSocketMessage` and send the same instance to many sessions.** Each session may have a different factory, and a message's buffer is released after being written; sending a released/foreign buffer risks corruption or leaks. Instead map **per session**: `broadcast.map(session::textMessage)` creates a fresh message with that session's factory. - If you must share a binary payload, encode to a byte[]/String and let each session wrap it, or use `DataBufferUtils.retain(...)` and manage release counts explicitly. - On the **inbound** side, remember received buffers are released after consumption; caching without retaining leaks. ## Constraint 5 — threading / event loop Reactor Netty runs handlers on a small pool of event-loop threads. Blocking work inside `doOnNext` (DB calls, `block()`) starves the loop and stalls all sessions on that thread. Offload blocking work with `subscribeOn/publishOn(Schedulers.boundedElastic())` and never call `block()` inside a handler. Size worker threads via Reactor Netty config if needed, and monitor buffer depth / connection count. ## Constraint 6 — ordering and heartbeats Since all outbound merges into one stream, interleave heartbeats (`session.pingMessage(...)`) into that single Flux rather than a second `send()`. For idle detection use `timeout(...)` on the merged pipeline (mapping the timeout error to a graceful close). ## When to reach for a broker instead Hand-rolled fan-out with `Sinks` is fine for single-instance, moderate N. For horizontal scaling (multiple app instances), a per-instance sink can't see events emitted on other instances — you need an external broker/relay (Redis pub/sub, a message broker, or STOMP with a broker relay — separate leaf). Recognizing that boundary is a principal-level judgment.
- Why not create one WebSocketMessage and send the same object to every session?Each session may use a different DataBufferFactory, and a message's pooled DataBuffer is released after it's written. Reusing one instance risks writing a released/foreign buffer — corruption or leaks. Map per session (broadcast.map(session::textMessage)) so each gets a fresh message.
- One client on a slow network stalls your broadcast to everyone. What's the fix?Give each subscriber isolated backpressure: a bounded onBackpressureBuffer with an overflow strategy (or onBackpressureLatest for stale-tolerant streams), decouple with publishOn, and optionally close chronically slow sessions on buffer overflow so they can't degrade the system.
- Your single-instance Sinks-based broadcast works, but events emitted on other app instances aren't delivered. Why, and what do you do?A Sinks.many() is in-JVM; it can't see events from other instances. For horizontal scale you need an external fan-out — Redis pub/sub, a message broker, or a STOMP broker relay — so all instances share the event stream.
saying these in an interview costs you the question
- Calling session.send() multiple times or writing from several threads to one session
- Unbounded buffering per subscriber (OOM on slow clients)
- Sharing a single WebSocketMessage/DataBuffer across many sessions
- Blocking (DB calls, .block()) inside doOnNext on the event loop
- Assuming an in-JVM Sinks broadcast scales across multiple instances