skip to content

Reactive Streams Specification

Publisher, Subscriber, Subscription and the onSubscribe/onNext/onError/onComplete signal contract, with request(n) as the demand mechanism. Interviewers ask because it is the spec that makes reactive libraries interoperable rather than a Reactor invention.

part ofSpring Frameworkoverview, primer and where to startread it →
on this pageshow

questions

5

What are the four interfaces defined by the Reactive Streams specification, and what is each responsible for?

level: juniorimportance: must knowfreq 70%

answer

  1. Publisher / Subscriber / Subscription / Processor
  2. subscribe returns void; Subscription arrives via onSubscribe
  3. Processor = Subscriber + Publisher
  4. org.reactivestreams package, vendor-neutral
  5. Flux/Mono ARE Publishers

basics

~10 s

Publisher produces data, Subscriber consumes it, Subscription links the two and lets the Subscriber request items or cancel, and Processor is both a Subscriber and a Publisher (a middle stage).

solid answer

~40 s

Reactive Streams defines four interfaces in the org.reactivestreams package. Publisher<T> has a single method subscribe(Subscriber). Subscriber<T> has onSubscribe(Subscription), onNext(T), onError(Throwable), onComplete(). Subscription has request(long n) and cancel(); it is the private channel between one Publisher and one Subscriber that carries demand upstream and cancellation upstream. Processor<T,R> extends both Subscriber<T> and Publisher<R>, representing an intermediate processing stage that consumes one stream and produces another. Reactor's Flux and Mono are Publisher implementations; Spring WebFlux uses these interfaces as the common contract so any compliant library (Reactor, RxJava, Akka Streams) interoperates. The whole point is a standard, vendor-neutral handshake for asynchronous stream processing with flow control.

code

java · 21 lines
java
// The four Reactive Streams interfaces (org.reactivestreams), paraphrased:
public interface Publisher<T> {
    void subscribe(Subscriber<? super T> s);
}

public interface Subscriber<T> {
    void onSubscribe(Subscription s); // exactly once, first
    void onNext(T t);                 // 0..N times
    void onError(Throwable t);        // terminal
    void onComplete();                // terminal
}

public interface Subscription {
    void request(long n); // signal demand upstream
    void cancel();        // stop and release
}

public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {}

// Reactor's Flux IS a Publisher, so this is legal:
org.reactivestreams.Publisher<String> pub = reactor.core.publisher.Flux.just("a", "b");

go deeper

for a junior

Must name all four interfaces and roughly what each does.

for a middle

Should know the exact method signatures and that subscribe returns void.

for a senior

Explains the interop rationale and how Reactor maps onto these interfaces.

for a principal

Discusses the interface set as an integration contract across libraries and the JDK Flow mirror.

**Reactive Streams** is a small specification (just four Java interfaces plus a written rule set) that standardizes asynchronous stream processing with non-blocking flow control. It lives in the `org.reactivestreams` package (artifact `org.reactivestreams:reactive-streams`). Spring WebFlux, Project Reactor (Flux/Mono), and RxJava all implement it so they can interoperate. The four interfaces: 1. **`Publisher<T>`** — the source of a potentially unbounded sequence of elements. Its only method is `void subscribe(Subscriber<? super T> s)`. A Publisher is passive: it does nothing until subscribed to (this is why Reactor pipelines are 'cold' and only run when subscribed). A Publisher may serve multiple Subscribers; each subscription is independent. 2. **`Subscriber<T>`** — the consumer. It has exactly four callback methods: - `onSubscribe(Subscription s)` — called first, exactly once, handing the Subscriber its private control channel. - `onNext(T t)` — delivers one element; called zero-to-many times. - `onError(Throwable t)` — terminal; signals failure. - `onComplete()` — terminal; signals successful end. 3. **`Subscription`** — the one-to-one contract token between a single Publisher and a single Subscriber. It has: - `request(long n)` — the Subscriber tells the Publisher it is ready to receive up to `n` more elements (this is the demand / flow-control mechanism). - `cancel()` — the Subscriber asks the Publisher to stop sending and release resources. The Subscription is what makes the stream *pull-and-push*: elements are pushed by the Publisher, but only up to the amount the Subscriber has pulled via `request`. 4. **`Processor<T,R>`** — `extends Subscriber<T>, Publisher<R>`. It is a stage that is simultaneously a Subscriber (of an upstream) and a Publisher (to a downstream) — e.g., a transformation or a bridge. In practice application developers rarely implement Processor directly; Reactor's operators (`map`, `filter`, etc.) fill this role internally, and `Sinks` largely replaced the older `Processor` implementations (`DirectProcessor`, `UnicastProcessor` are deprecated). **Why it matters in Spring:** WebFlux controllers return `Flux`/`Mono` (Reactor Publishers), but the framework only depends on the `Publisher` interface, so a controller could return an RxJava `Flowable` and it still works. This interface set is the lingua franca. **Gotcha:** `Publisher.subscribe` returns `void`, not the Subscription — the Subscription is delivered *to the Subscriber* via `onSubscribe`. Nothing flows until the Subscriber calls `request(n)`. **Note:** The JDK later mirrored these exact interfaces as `java.util.concurrent.Flow.Publisher/Subscriber/Subscription/Processor` (Java 9+); they are semantically identical, and adapters bridge the two.

  • Why does Publisher.subscribe return void instead of the Subscription?
    Because subscription is asynchronous: the Publisher may need to do setup before it can hand back a control channel. It delivers the Subscription to the Subscriber via onSubscribe(Subscription), which is guaranteed to be called before any onNext. This keeps the API uniform for sync and async sources.
  • Do you often implement Subscriber or Processor by hand in Spring apps?
    Rarely. You compose Reactor operators (map/filter/flatMap) and return Flux/Mono; Reactor implements the interfaces internally. You implement a raw Subscriber only for low-level custom sinks or interop, and Processor almost never — Sinks replaced the old Processor types.

saying these in an interview costs you the question

  • Saying Publisher.subscribe returns the Subscription
  • Confusing Subscription (control channel) with Subscriber (consumer)
  • Claiming a Publisher starts emitting immediately without a subscriber
  • Thinking Processor is a separate thing unrelated to Subscriber/Publisher

context

open as a page

What do Subscription.request(n) and Subscription.cancel() do, and what are the rules around calling them?

level: middleimportance: must knowfreq 60%

basics

~20 s

request(n) tells the Publisher the Subscriber is ready to receive up to n more elements — it is the demand signal. cancel() tells the Publisher to stop emitting and release resources. Both are called by the Subscriber on its Subscription.

open as a page

Describe the signal protocol a Subscriber receives: onSubscribe, onNext, onError, onComplete — their ordering and cardinality.

level: middleimportance: must knowfreq 65%

basics

~10 s

onSubscribe is called first, exactly once. Then onNext happens zero or more times. The stream ends with exactly one terminal signal: either onError or onComplete — never both, and nothing comes after it.

open as a page

Why does the Reactive Streams specification exist, how does Project Reactor relate to it, and what is its relationship to java.util.concurrent.Flow?

level: seniorimportance: should knowfreq 45%

basics

~10 s

The spec gives asynchronous stream libraries a common, vendor-neutral contract so they interoperate with built-in flow control. Reactor implements it (Flux/Mono are Publishers). Java 9's java.util.concurrent.Flow contains the identical interfaces, mirrored into the JDK.

open as a page

What contractual guarantees must a compliant Publisher provide to its Subscriber, and why do these rules matter for building operators?

level: principalimportance: should knowfreq 30%

basics

~20 s

A compliant Publisher must: call onSubscribe first and exactly once, never emit more onNext than requested, deliver signals serially (never concurrently), end with at most one terminal signal, never emit after a terminal or after cancel, and never pass null. These guarantees let operators compose safely.

open as a page