What is the java.util.concurrent.Flow API, and how do Flow.Publisher, Flow.Subscriber, and Flow.Subscription cooperate to deliver backpressure?
answer
- Four interfaces: Publisher, Subscriber, Subscription, Processor
- Subscriber methods: onSubscribe/onNext/onError/onComplete
- request(n) = pull-based demand = backpressure
- Exactly one terminal: onComplete or onError; then no onNext
- Signals to one subscriber are serialized; SubmissionPublisher is built in
basics
~20 sFlow (added in Java 9) is the JDK's reactive-streams API: a Publisher produces items, a Subscriber consumes them. When you subscribe, you get a Subscription and must call request(n) to ask for n items — the publisher only sends up to what you requested. That request mechanism is backpressure, so a fast producer can't overwhelm a slow consumer.
solid answer
~40 sjava.util.concurrent.Flow is the JDK's standard reactive-streams contract (Java 9), an asynchronous, backpressure-aware Observer. It defines four nested interfaces: Publisher (has subscribe(Subscriber)), Subscriber (onSubscribe, onNext, onError, onComplete), Subscription (request(long n), cancel), and Processor (both). The handshake: a Subscriber calls publisher.subscribe; the publisher calls subscriber.onSubscribe(subscription); the subscriber stores it and calls subscription.request(n) to pull at most n items; the publisher then delivers up to n onNext calls, and the subscriber requests more as it processes. This demand-driven pull is backpressure — the consumer dictates the rate so a fast producer can't overrun a slow one. The stream ends with exactly one terminal signal, onComplete or onError, after which no more onNext. The JDK ships SubmissionPublisher as a ready Publisher. Flow is just the interfaces (the reactive-streams SPI); libraries like Reactor or RxJava interoperate with it.
code
java · 19 linesclass PrintSubscriber implements Flow.Subscriber<String> {
private Flow.Subscription subscription;
@Override public void onSubscribe(Flow.Subscription s) {
this.subscription = s;
s.request(1); // pull one item -> backpressure
}
@Override public void onNext(String item) {
System.out.println(item);
subscription.request(1); // ask for the next only after handling this
}
@Override public void onError(Throwable t) { t.printStackTrace(); }
@Override public void onComplete() { System.out.println("done"); }
}
try (SubmissionPublisher<String> publisher = new SubmissionPublisher<>()) {
publisher.subscribe(new PrintSubscriber());
publisher.submit("a");
publisher.submit("b");
} // close() triggers onComplete after buffered items draingo deeper
Knows Flow is a publisher/subscriber API added in Java 9 for streams of data.
Can name the four interfaces and describe that the subscriber requests items so the producer doesn't overwhelm it.
Walks the full subscribe/onSubscribe/request/onNext/terminal handshake, explains request(n) as backpressure, and the single-terminal-signal rule.
Adds the serialization guarantee, Flow's role as the reactive-streams SPI/interop point, choice of Flow vs synchronous listeners vs full reactive libraries, and integration with SubmissionPublisher / the HTTP client at scale.
## Why Flow exists The deprecated `java.util.Observable` and synchronous listeners deliver events by *pushing*: the subject calls the observer immediately, on the subject's thread, with no way for the observer to say "slow down." For **asynchronous streams of many items** — events, rows, messages — a fast producer can overwhelm a slow consumer, causing unbounded buffering or dropped data. **`java.util.concurrent.Flow`** (Java 9) is the JDK's answer: the standard **reactive-streams** contract, an Observer pattern with **backpressure**. Flow is deliberately *just interfaces* — a service-provider contract — mirroring the external Reactive Streams specification so libraries (Reactor, RxJava, Akka Streams) interoperate through it. ## The four interfaces All are nested inside `Flow`: - **`Flow.Publisher<T>`** — the source. One method: `subscribe(Flow.Subscriber<? super T> s)`. - **`Flow.Subscriber<T>`** — the consumer. Four methods: `onSubscribe(Subscription)`, `onNext(T item)`, `onError(Throwable)`, `onComplete()`. - **`Flow.Subscription`** — the per-subscription control link. Two methods: `request(long n)` (demand for up to n more items) and `cancel()`. - **`Flow.Processor<T,R>`** — both a Subscriber and a Publisher (a stage that transforms a stream). ## The handshake, step by step 1. A subscriber calls `publisher.subscribe(subscriber)`. 2. The publisher creates a `Subscription` and calls `subscriber.onSubscribe(subscription)`. 3. Inside `onSubscribe`, the subscriber saves the subscription and calls `subscription.request(n)` to signal it can handle **n** items. **Until it requests, it receives nothing** — demand is opt-in. 4. The publisher delivers **at most n** items via repeated `onNext(item)` calls. 5. As the subscriber processes items it calls `request(...)` again to ask for more, keeping a bounded outstanding demand. 6. The stream terminates with **exactly one** terminal signal: `onComplete()` (success) or `onError(throwable)` (failure). After a terminal signal, `onNext` must not be called. The subscriber may also `cancel()` at any time to stop. ## Backpressure — the whole point **Backpressure** is the mechanism by which a slow consumer limits a fast producer. Because items flow only in response to `request(n)`, the *consumer* controls the rate. A producer that respects the contract never sends more than the outstanding requested amount, so there is no unbounded buffering and no overrun. This is the key difference from push-only listeners and from the old Observable. ## Concurrency contract Signals to a given subscriber (`onSubscribe`/`onNext`/`onError`/`onComplete`) must be **serialized** — never invoked concurrently for the same subscriber, even if the publisher uses many threads. This lets the subscriber be written as if single-threaded. Delivery is asynchronous and typically happens on publisher-managed threads. ## What the JDK gives you out of the box The JDK ships **`SubmissionPublisher<T>`**, a concrete `Flow.Publisher` you can `submit(item)` to; it buffers and delivers to subscribers with backpressure on an executor. The HTTP client (`java.net.http`) also exposes bodies as Flow publishers/subscribers. For richer operators (map/filter/flatMap) you use a library that builds on the same Flow contract. ## Relationship to the rest of the Observer story Flow is the *asynchronous, backpressure-aware* member of the JDK's Observer family. Use `PropertyChangeListener` or Swing listeners for synchronous in-VM notifications; use Flow for async, possibly-infinite, rate-controlled streams. It is the modern replacement for the deprecated `java.util.Observable` when you need streaming. ## Deriving your own answer A senior should name the four interfaces, walk the subscribe/onSubscribe/request/onNext/terminal handshake, and articulate that `request(n)` *is* backpressure. A principal adds the serialization guarantee, the SPI/interop role, and when to choose Flow over synchronous listeners.
- If a Subscriber never calls subscription.request(n), what does it receive?Nothing. Demand is opt-in: the publisher sends onNext only up to the total requested. With zero outstanding demand, no items are delivered (though it may still get onComplete/onError if the stream ends). This is exactly how backpressure throttles a fast producer to zero.
- Does Flow provide map/filter/flatMap operators?No. java.util.concurrent.Flow is only the four SPI interfaces plus SubmissionPublisher. Operator-rich pipelines come from libraries (Reactor, RxJava) that implement the same Flow/Reactive-Streams contract and interoperate through it.
saying these in an interview costs you the question
- Saying the publisher pushes freely regardless of request(n) (that breaks backpressure)
- Thinking the subscriber receives onNext before calling request
- Allowing onNext after onComplete/onError, or multiple terminal signals
- Confusing Flow (async, backpressure) with synchronous listeners/Observable
- Believing Flow ships rich operators like map/filter (it's just the SPI interfaces)