skip to content

How does @SubscriptionMapping work and why must its handler return a streaming Publisher (Flux)?

level: seniorimportance: should knowfreq 40%

answer

  1. Subscription root type = long-lived stream
  2. must return Publisher / Flux
  3. each emission = one pushed result
  4. WebSocket (graphql-transport-ws) or SSE, not plain POST
  5. 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
java
@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

for a junior

Know a subscription streams many values and returns a Flux.

for a middle

Explain per-emission delivery, WebSocket/SSE transport, and completion/error/cancel semantics.

for a senior

Contrast query-Flux (collected) vs subscription-Flux (streamed), threading/ThreadLocal caveats, and backpressure.

for a principal

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)

context