skip to content

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%

answer

  1. onSubscribe first & once; onNext <= demand
  2. signals serialized (no concurrency) -> no locks in operators
  3. one terminal, nothing after terminal or cancel
  4. no nulls; request(n<=0) -> onError
  5. TCK proves compliance; Sinks serialize for you

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.

solid answer

~40 s

The spec's rules turn the four interfaces into a reliable contract. A Publisher must invoke onSubscribe exactly once before any other signal; must never emit more onNext than the total request(n) demand; must deliver all signals to a given Subscriber serially with a happens-before relationship, never concurrently; must terminate with at most one signal (onComplete xor onError) after which nothing more is sent; must stop (best-effort) after cancel(); must never pass null in onNext/onError; and must respond to request(n<=0) with onError(IllegalArgumentException). These guarantees are what make operators composable: because each operator can assume serialized, demand-bounded, single-terminal input, it can implement map/filter/flatMap without defensive locking or overflow checks, and chains of operators remain correct. The TCK exists precisely to verify a Publisher honors all of this. Reactor's operators rely on these invariants internally.

code

java · 17 lines
java
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;

// Sinks give a COMPLIANT, serialized Publisher instead of a hand-rolled one
// that could violate Rule 1.3 (serialization) or 1.1 (demand) or 1.7 (single terminal).
Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();
Flux<String> flux = sink.asFlux(); // a spec-compliant Publisher

// tryEmitNext honors the contract; emitComplete is the single terminal.
sink.tryEmitNext("a");
sink.tryEmitNext("b");
sink.tryEmitComplete();          // terminal — nothing may follow
// sink.tryEmitNext("c");        // would be rejected: no signals after terminal

// null is illegal (Rule 2.13): sink.tryEmitNext(null) -> IllegalArgumentException

flux.subscribe(System.out::println); // prints a, b

go deeper

for a junior

Can state onSubscribe-first and single-terminal at a high level.

for a middle

Lists the main guarantees: demand bound, serialization, one terminal, no null.

for a senior

Explains why serialization removes locking and how demand-bounding underpins flow control.

for a principal

Connects the rule set to operator composability and cites the TCK as the compliance mechanism; reasons about how custom emitters break the contract.

The Reactive Streams **rule set** (numbered rules across Publisher §1, Subscriber §2, Subscription §3, and Processor §4) is what makes the four bare interfaces trustworthy. The load-bearing guarantees a **compliant Publisher** must provide: **1. onSubscribe first, exactly once (Rule 1.9, 2.12).** Before any `onNext`/`onError`/`onComplete`, the Publisher calls `onSubscribe(Subscription)` exactly once. Even if subscription fails immediately, it must still call `onSubscribe` (with a Subscription) and then `onError`. A Subscriber can therefore always rely on having its control channel before data arrives. **2. Demand is never exceeded (Rule 1.1).** The total number of `onNext` calls must be ≤ the sum of `request(n)` demand issued so far. This is the flow-control backbone — a Publisher may emit fewer than requested (e.g., stream ends), but never more. **3. Signals are serialized (Rule 1.3).** All signals to a single Subscriber are issued sequentially, with a happens-before relationship between them; they are **never concurrent**. This is arguably the most important rule for implementers: it means operator state doesn't need locks on the signal path. **4. At most one terminal, and it's final (Rule 1.7, 1.2).** A stream ends with `onComplete` **or** `onError`, never both, never repeated. After a terminal signal, the Publisher emits nothing further, and the Subscription is considered cancelled. **5. Emissions stop after cancel (Rule 1.8, best-effort).** After the Subscriber calls `cancel()`, the Publisher must **eventually** stop signaling. It's best-effort/asynchronous, so a few trailing `onNext` may still arrive, but no new work should be scheduled. **6. No null values (Rule 2.13 / 1.x).** `onNext` and `onError` must never be passed `null`; doing so is a violation and must throw `NullPointerException` at the call site. This is why **Reactor forbids null elements** — you model absence with `Mono.empty()` or by filtering, never a null onNext. **7. Invalid demand is an error (Rule 3.9).** `request(n)` with `n <= 0` must result in `onError(IllegalArgumentException)`. **8. Subscription is single-use / one-to-one (Rule 2.12, 3.x).** A Subscription belongs to exactly one Publisher–Subscriber pair; re-subscribing produces a new Subscription. Calling `request`/`cancel` after a terminal is a no-op. **Why these matter for building operators.** An **operator** (Reactor's `map`, `filter`, `flatMap`, `take`, etc.) is itself a `Processor`-like element: it subscribes upstream and publishes downstream. Because the contract guarantees (a) serialized signals, (b) demand-bounded onNext, (c) a single terminal, and (d) no nulls, an operator can be written as a small state machine **without** defensive synchronization on the signal path or overflow buffers of its own — it simply transforms/forwards signals and relays demand. If any of these guarantees could be violated, every operator would need to re-validate its input, and composition of long operator chains would be unsound. This is exactly why the **TCK (Technology Compatibility Kit)** exists: a library like Reactor runs the TCK against its Publisher implementations to prove every rule holds, so downstream operators and application code can trust the invariants. When you write a custom `Publisher`/`Subscriber` for interop, running the TCK (`org.reactivestreams:reactive-streams-tck`) is the recommended validation step. **Practical implications / gotchas:** - A hand-rolled Publisher that emits from multiple threads without serialization is the classic compliance bug — it violates Rule 1.3 and breaks downstream operators nondeterministically. Use Reactor `Sinks` (which serialize) rather than rolling your own. - Emitting an element after `onComplete` (e.g., a late async callback) silently corrupts downstream state — guard terminal state. - Because demand must never be exceeded, a Publisher that ignores `request(n)` and pushes eagerly is non-compliant even if it 'works' under `Long.MAX_VALUE` demand. - These are the *contract* guarantees; the concrete overflow/backpressure *strategies* for when a Subscriber can't keep up are covered by the dedicated backpressure topic.

  • How would you validate that a custom Publisher you wrote is spec-compliant?
    Run it against the Reactive Streams TCK (org.reactivestreams:reactive-streams-tck), which is a suite of tests exercising every rule — onSubscribe ordering, demand bounding, serialization, terminal-signal finality, null handling, and invalid-demand behavior. Passing the TCK is the standard proof of compliance.
  • Why is the serialization guarantee (Rule 1.3) so important for operator authors?
    Because it means signals to a downstream Subscriber never overlap, an operator can hold mutable state (counters, buffers, terminal flags) and mutate it on the signal path without locks. Without it, every operator would need synchronization, killing performance and complicating correctness.
  • What is a common way to accidentally build a non-compliant Publisher?
    Emitting onNext from multiple threads concurrently (violating serialization), or letting a late async callback fire onNext after onComplete/onError. Using Reactor Sinks or create/generate with proper synchronization avoids this; hand-rolled multi-threaded emitters are the usual culprit.

saying these in an interview costs you the question

  • Thinking a Publisher may emit more onNext than requested if it 'has data ready'
  • Believing concurrent onNext from multiple threads is allowed
  • Assuming you can emit after onComplete for a 'final cleanup' value
  • Saying null onNext is acceptable and handled downstream
  • Not knowing the TCK exists to verify compliance

context