How does @SubscriptionMapping work and why must its handler return a streaming Publisher (Flux)?
answer
- Subscription root type = long-lived stream
- must return Publisher / Flux
- each emission = one pushed result
- WebSocket (graphql-transport-ws) or SSE, not plain POST
- no ThreadLocal on emission threads; cancel cleans up
basics
~20 s@SubscriptionMapping resolves a field under the Subscription root type. Because a subscription pushes many values over time, the method returns a Flux (or any Reactive Streams Publisher); each emitted item is delivered to the client as a separate result.
solid answer
~40 s@SubscriptionMapping marks a handler for a field under the root Subscription type. Unlike queries/mutations that produce a single result, a subscription is an ongoing stream, so the method must return a `Flux<T>` (or any Reactive Streams `Publisher`). Each element Flux emits becomes one GraphQL result pushed to the client, typically over a WebSocket (graphql-transport-ws) or SSE transport. The stream stays open until it completes, errors, or the client cancels. Returning a non-streaming type (a plain object or `Mono`) is invalid for a subscription and fails wiring. Key concerns: the reactive pipeline runs off the request thread, so no request-scoped/ThreadLocal assumptions; backpressure is honored by the Publisher; errors terminating the Flux end the subscription. Compared to @QueryMapping returning Flux — which is collected into a list — a @SubscriptionMapping Flux is streamed element by element.
code
java · 23 lines@Controller
public class CommentController {
private final Sinks.Many<Comment> sink =
Sinks.many().multicast().onBackpressureBuffer();
@MutationMapping
public Comment addComment(@Argument String bookId, @Argument String text) {
Comment c = commentService.save(bookId, text);
sink.tryEmitNext(c); // fan out to subscribers
return c;
}
// Streams each new comment for a book; must return a Publisher
@SubscriptionMapping
public Flux<Comment> commentAdded(@Argument String bookId) {
return sink.asFlux()
.filter(c -> c.bookId().equals(bookId));
}
}
// schema.graphqls
// type Subscription { commentAdded(bookId: ID!): Comment }go deeper
Know a subscription streams many values and returns a Flux.
Explain per-emission delivery, WebSocket/SSE transport, and completion/error/cancel semantics.
Contrast query-Flux (collected) vs subscription-Flux (streamed), threading/ThreadLocal caveats, and backpressure.
Design multicast/backpressure strategy (Sinks, hot vs cold), security-context propagation to emission threads, and resource cleanup on cancel.
GraphQL has three operation types: `query`, `mutation`, and `subscription`. A **subscription** is a long-lived operation: the server pushes a sequence of results to the client as events happen (e.g. new messages, price ticks). Spring for GraphQL models this with `@SubscriptionMapping`. **Signature requirement**: the handler must return a **Reactive Streams `Publisher`** — in practice a Project Reactor `Flux<T>` (multiple values) or occasionally `Mono<T>`/`Flux` where each item is a subscription result. graphql-java's subscription execution strategy expects a `Publisher` for the field's result. Returning a plain value or a single object is not valid for a subscription field. **Element semantics**: each element the `Flux` emits is turned into one GraphQL execution result and delivered to the client. The stream lives until: - it **completes** (Flux `onComplete`) → subscription ends normally, - it **errors** (`onError`) → the error is delivered and the subscription terminates, - the **client cancels** / disconnects → the underlying subscription is cancelled (Reactor `cancel` propagates), so you should clean up resources in the stream. **Transport**: subscriptions need a streaming transport. Spring supports WebSocket using the `graphql-transport-ws` protocol and also Server-Sent Events (SSE). Plain HTTP POST (used for query/mutation) can't carry a stream. You enable the WebSocket endpoint via `spring.graphql.websocket.path`. **Threading**: the returned `Flux` is subscribed and executed on reactive scheduler threads, not the servlet request thread. Therefore: - Don't rely on `ThreadLocal`/request-scoped state inside the stream; pass needed data through `GraphQLContext` / `@ContextValue` or capture it before returning the Flux. - Blocking calls inside the pipeline should be offloaded (`publishOn`/`subscribeOn(Schedulers.boundedElastic())`). **Contrast with @QueryMapping/@MutationMapping returning Flux**: for a **query/mutation** field whose GraphQL type is a `List`, returning a `Flux` is allowed but Spring **collects** it into a list — a single result. For a **subscription**, the `Flux` is **streamed** element-by-element. Same return type, different execution model driven by the operation type. **Arguments**: subscription handlers accept the same injectable parameters — `@Argument`, `DataFetchingEnvironment`, `@ContextValue`, `Principal`, etc. **Gotchas** - Forgetting to return a `Publisher` → the subscription field won't wire / errors at runtime. - Doing side effects eagerly (before returning the Flux) instead of inside the reactive pipeline — they run at assembly time, not per-subscription, and won't be cancelled. - Not handling `onError` — an unexpected error tears down the whole subscription for that client. - Assuming security context is present on emission threads; propagate it explicitly if needed. **When to use**: real-time push (chat, notifications, live dashboards). If the client only needs to poll or fetch once, use a query, not a subscription.
- If @QueryMapping and @SubscriptionMapping both return Flux<Comment>, how does execution differ?The query Flux is collected into a single list result; the subscription Flux is streamed element-by-element as separate results over the WebSocket/SSE transport.
- What transport is required for subscriptions and why can't a normal HTTP POST be used?A streaming transport — WebSocket (graphql-transport-ws) or SSE — because a subscription pushes many results over time; a single HTTP POST response can't carry an open-ended stream of events.
saying these in an interview costs you the question
- Saying a subscription handler can return a plain object or a single value
- Assuming subscriptions work over ordinary HTTP POST like queries
- Relying on request-scoped ThreadLocals inside the subscription Flux
- Believing a Flux from @QueryMapping is streamed like a subscription (it's collected to a list)