Describe the signal protocol a Subscriber receives: onSubscribe, onNext, onError, onComplete — their ordering and cardinality.
answer
- Grammar: onSubscribe onNext* (onError | onComplete)?
- onSubscribe first & exactly once
- one terminal, never both, nothing after
- onNext bounded by request(n); terminals are not
- signals serialized -> no locks needed
basics
~10 sonSubscribe 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.
solid answer
~40 sThe protocol is a strict grammar: `onSubscribe onNext* (onError | onComplete)?`. onSubscribe(Subscription) is always the first signal and fires exactly once — it hands over the control channel. After that, onNext(T) may fire zero or more times, but only up to the total demand the Subscriber has requested via Subscription.request(n). The sequence ends with at most one terminal signal: onComplete() for success or onError(Throwable) for failure — mutually exclusive and final. No onNext, onError, or onComplete may occur before onSubscribe or after a terminal signal. All signals for a given Subscriber must be serialized (no concurrent invocation), so the Subscriber never needs internal synchronization for these callbacks. In Reactor you usually observe these via doOnNext/doOnError/doOnComplete or by implementing BaseSubscriber.
code
java · 24 linesimport org.reactivestreams.Subscription;
import reactor.core.publisher.BaseSubscriber;
import reactor.core.publisher.Flux;
// Observe the exact signal protocol with a hand-written Subscriber.
Flux.just("a", "b", "c")
.subscribe(new BaseSubscriber<String>() {
@Override protected void hookOnSubscribe(Subscription s) {
// FIRST signal, exactly once. Must request demand or nothing flows.
request(2); // ask for 2 elements
}
@Override protected void hookOnNext(String value) {
System.out.println("onNext: " + value);
request(1); // ask for the next one
}
@Override protected void hookOnComplete() {
System.out.println("onComplete"); // terminal (success)
}
@Override protected void hookOnError(Throwable t) {
System.out.println("onError: " + t); // terminal (failure)
}
});
// Output: onNext a, onNext b, onNext c, onComplete
// Note: onComplete OR onError fires, never both.go deeper
Knows the four signals and that there is one terminal.
Can state the full grammar and the exactly-once / at-most-one rules.
Explains serialization guarantee, null prohibition, and the onSubscribe-before-error rule.
Reasons about how operators enforce the grammar and why terminals are demand-independent.
The Reactive Streams **signal protocol** is a formal grammar describing the legal sequence of method calls a `Subscriber<T>` receives. Written in regex form: ``` onSubscribe onNext* (onError | onComplete)? ``` **Signal-by-signal:** - **`onSubscribe(Subscription s)`** — Always the **first** signal, invoked **exactly once**. Until it fires, the Subscriber has no way to request data, so nothing else can legally happen. Inside it, the Subscriber typically stores the Subscription and calls `s.request(n)` to start demand. - **`onNext(T t)`** — Delivers a single element. Occurs **zero to many** times. Crucially, the Publisher must not emit more `onNext` calls than the cumulative `request(n)` demand the Subscriber has issued. `t` must never be `null` (Rule 2.13 — passing null is a spec violation; that's why Reactor forbids null elements and you use `Mono.empty()`/filtering instead). - **`onError(Throwable t)`** — A **terminal** signal indicating failure. After it, no further signals are allowed. The `Throwable` must be non-null. - **`onComplete()`** — A **terminal** signal indicating the stream finished successfully with no more elements. **Key rules / guarantees:** 1. **At most one terminal signal.** A stream ends with `onError` OR `onComplete`, never both, never twice. An empty successful stream is `onSubscribe` then immediately `onComplete` (zero onNext). A `Mono` emits at most one onNext then onComplete, or just onComplete if empty, or onError. 2. **Nothing after a terminal.** Once `onError`/`onComplete` fires, the Subscription is considered cancelled and further `onNext`/`request` calls are meaningless/illegal. 3. **Serialized (non-concurrent) signals.** All calls to a single Subscriber happen-before the next and are never concurrent (Rule 1.3). This means a Subscriber implementation does **not** need to guard these callbacks with locks — the contract guarantees a single logical thread of signals at a time (though which thread may vary). 4. **onSubscribe precedes everything.** Even an error that occurs during subscription setup must be reported by first calling `onSubscribe` (with a no-op Subscription) and then `onError` — a Subscriber can rely on `onSubscribe` always coming first. **Reactor mapping:** You rarely implement `Subscriber` directly. Instead you use lifecycle hooks — `doOnSubscribe`, `doOnNext`, `doOnError`, `doOnComplete`, `doOnTerminate` — or extend `reactor.core.publisher.BaseSubscriber`, which gives you `hookOnSubscribe`, `hookOnNext`, `hookOnComplete`, `hookOnError`, and safe `request(n)`/`cancel()`. **Gotchas:** - Emitting `onNext` before `onSubscribe`, or after `onComplete`, is a spec violation — well-behaved operators enforce this, but a hand-rolled Publisher can break it. - `onComplete`/`onError` take **no demand** — they are delivered regardless of `request(n)`; only `onNext` is demand-bounded. - Because signals are serialized, throwing from inside `onNext` is undefined-ish; you should route failures through `onError`, not by throwing from the callback.
- Can onNext deliver a null value?No. Reactive Streams Rule 2.13 forbids null in onNext (and null Throwable in onError). That's why Reactor's Flux/Mono cannot carry nulls — you use Mono.empty(), filtering, or Optional-wrapping instead. Passing null is a spec violation.
- Since signals are serialized, does my Subscriber need synchronization?Not for the callback methods themselves — the spec guarantees they are never invoked concurrently (Rule 1.3), with a happens-before relationship between successive signals. You may still need synchronization if the Subscriber shares state with other threads outside the signal path.
saying these in an interview costs you the question
- Saying onComplete and onError can both fire
- Thinking onSubscribe can come after the first onNext
- Believing onNext can deliver null
- Assuming callbacks can be invoked concurrently, requiring locks
- Thinking onComplete requires outstanding demand to be delivered