How do WebSocketSession.receive() and WebSocketSession.send() work with Flux<WebSocketMessage>, and how do you wire them together in a handler?
answer
- receive() = Flux<WebSocketMessage> in; send(Publisher) = Mono<Void> out
- session.textMessage(...) / binaryMessage(...) factories
- getPayloadAsText() vs getPayload() DataBuffer
- Combine in+out with Mono.zip / Flux.merge / then()
- Pooled DataBuffers released after consume — retain if cached
basics
~20 ssession.receive() returns a Flux<WebSocketMessage> of inbound frames; session.send(Publisher<WebSocketMessage>) takes an outbound stream and returns Mono<Void>. You transform the inbound Flux (or a separate source) into outbound WebSocketMessages and pass it to send(), returning that Mono.
solid answer
~40 sA WebSocketSession exposes the connection as two reactive streams. `receive()` gives a `Flux<WebSocketMessage>` — one element per inbound frame; you read text via `WebSocketMessage.getPayloadAsText()` or the raw `DataBuffer` via `getPayload()`. `send(Publisher<WebSocketMessage>)` consumes an outbound `Flux`/`Mono` and returns a `Mono<Void>` completing when that source completes. You create messages with factory methods on the session: `session.textMessage(String)`, `session.binaryMessage(...)`. Typical patterns: echo — `send(receive().map(...))`; outbound-only push — `send(someFlux.map(session::textMessage))` and ignore `receive()`; bidirectional — merge an inbound-processing stream (that emits nothing but has side effects) with an outbound source using `Mono.zip`/`Flux.merge`/`then`. Backpressure is honored end to end. Crucially the inbound `DataBuffer` payloads are pooled and released once consumed, so don't cache them without retaining.
code
java · 24 linesimport org.springframework.web.reactive.socket.*;
import reactor.core.publisher.*;
import java.time.Duration;
public class ChatHandler implements WebSocketHandler {
@Override
public Mono<Void> handle(WebSocketSession session) {
// Inbound: process each frame, emit nothing, then complete.
Mono<Void> input = session.receive()
.doOnNext(msg -> handleIncoming(msg.getPayloadAsText()))
.then();
// Outbound: an independent server-push stream mapped to text frames.
Flux<WebSocketMessage> outbound = Flux.interval(Duration.ofSeconds(1))
.map(i -> session.textMessage("server-tick-" + i));
Mono<Void> output = session.send(outbound);
// Return ONE Mono that runs both; session closes when the combined stream ends.
return Mono.zip(input, output).then();
}
private void handleIncoming(String text) { /* ... */ }
}go deeper
Know receive() returns inbound frames and send() takes an outbound stream returning Mono<Void>.
Wire echo/push/bidirectional patterns and use session.textMessage; combine streams into one Mono.
Reason about backpressure on slow clients and the single-writer constraint per session.
Discuss pooled DataBuffer retain/release semantics and backpressure strategy for hot broadcast sinks.
## The two streams of a session `WebSocketSession` models the live connection as a duplex pair of reactive streams: - **Inbound:** `Flux<WebSocketMessage> receive()` — emits one `WebSocketMessage` per frame the client sends, completes when the client closes the connection, errors on transport failure. - **Outbound:** `Mono<Void> send(Publisher<WebSocketMessage> messages)` — you hand it a publisher; the framework writes each emitted message to the wire, applying backpressure, and the returned `Mono<Void>` completes when your source `Publisher` completes. ## WebSocketMessage `WebSocketMessage` wraps a payload plus a `Type` (`TEXT`, `BINARY`, `PING`, `PONG`). Reading: - `getPayloadAsText()` → decodes the payload `DataBuffer` as UTF-8 text. - `getPayload()` → the raw `org.springframework.core.io.buffer.DataBuffer`. - `getType()` → the frame type. Creating messages — **use the session factory methods** so the correct `DataBufferFactory` is used: - `session.textMessage(String)` - `session.binaryMessage(Function<DataBufferFactory, DataBuffer>)` - `session.pingMessage(...)` / `session.pongMessage(...)` ## Wiring patterns **1. Echo (inbound drives outbound):** ```java return session.send(session.receive().map(m -> session.textMessage(m.getPayloadAsText()))); ``` **2. Server push / outbound-only** (e.g., a ticker) — you ignore `receive()` but note the caveat below: ```java Flux<WebSocketMessage> ticks = Flux.interval(Duration.ofSeconds(1)) .map(i -> session.textMessage("tick-" + i)); return session.send(ticks); ``` **3. Bidirectional / independent in and out.** `receive()` and `send()` are separate streams, but you must return **one** `Mono<Void>`. Combine them so both run and the session closes when appropriate: ```java Mono<Void> input = session.receive() .doOnNext(m -> process(m.getPayloadAsText())) .then(); // completes, emits nothing Mono<Void> output = session.send(outboundFlux); return Mono.zip(input, output).then(); // or Flux.merge(input, output).then() ``` Use `Mono.zip`/`Flux.merge` when you want the session to end when *either*/*both* streams end depending on your semantics. ## Backpressure and threading - Backpressure is honored: `send()` only pulls from your outbound publisher as fast as the socket can write; if the client is slow, demand slows upstream. Unbounded hot sources (e.g., a broadcast sink) should use `onBackpressureBuffer`/`onBackpressureDrop` to avoid overwhelming a slow client. - **You generally get one `send()` subscription per session.** Don't call `send()` multiple times concurrently and don't write to the session from multiple threads — sessions are not designed for concurrent writes. Merge everything into a single outbound `Flux`. ## DataBuffer lifecycle — the classic leak Inbound `WebSocketMessage` payloads are backed by **pooled** `DataBuffer`s (especially on Reactor Netty). The framework **releases** the buffer once the message has been consumed within your `receive()` subscription. Two consequences: - If you **cache/delay/defer** a received message beyond the immediate stream (e.g., buffer, replay, hand to another async boundary), you must **retain** it (`DataBufferUtils.retain(...)`) and release later, or you get use-after-free / leak warnings. - If you call `getPayloadAsText()` eagerly in the stream (the common case), you're safe because you extract the String before the buffer is released. - Also: if you **never subscribe to `receive()`** in an inbound-capable handler, incoming frames may not be drained. For a pure push handler this is usually fine, but if the client sends data, prefer `session.receive().then()` merged into your pipeline to drain and release buffers. ## When to use raw receive()/send() Use it for low-level, per-frame control (custom binary protocols, simple echo/push, proxies). For structured pub/sub messaging use STOMP or RSocket (different leaves).
- Why can caching a received WebSocketMessage cause a memory/leak problem?Inbound payloads are backed by pooled DataBuffers that the framework releases once the message is consumed in the receive() stream. If you hold the message past that point without DataBufferUtils.retain(...), you risk use-after-free or a leak; extract the text/bytes eagerly instead.
- You have an inbound stream and an independent outbound stream. How do you return them from one handle() call?Convert inbound side-effects to Mono<Void> via .then(), wrap outbound in session.send(...), and combine with Mono.zip(input, output).then() or Flux.merge(input, output).then() so both subscribe and the session ends per your chosen semantics.
saying these in an interview costs you the question
- Calling session.send() multiple times concurrently or writing from multiple threads
- Assuming getPayload() returns a String rather than a DataBuffer
- Caching received messages without retaining the DataBuffer
- Building outbound frames with new WebSocketMessage(...) instead of the session factory methods
- Thinking receive() and send() cannot both be active on the same session