skip to content

Explain why subscribeOn is position-independent while publishOn is not, and how you'd reason about which thread each operator runs on in a chain that mixes both.

level: principalimportance: should knowfreq 40%

answer

  1. subscription flows UP, data flows DOWN
  2. subscribeOn hooks upward → position-independent, closest-to-source wins
  3. publishOn hooks downward → position matters, stacks
  4. trace: source→first publishOn on subscribeOn thread; each publishOn re-segments
  5. subscribeOn no-op on already-async sources; interval/delay default to parallel

basics

~20 s

subscribeOn hooks the subscription, which flows upward to the source, so its position doesn't matter — it sets where the source runs. publishOn hooks the downward data signals, switching threads only for operators below it, so position matters.

solid answer

~50 s

Signals travel two ways: **subscription upward** (subscriber → source) and **data downward** (source → subscriber). `subscribeOn` intercepts the *subscription* pass: wherever it sits, when the subscribe signal reaches it, it schedules the upstream subscription (and thus the source's emission) on its Scheduler — so **position is irrelevant** and only the one **closest to the source** wins. `publishOn` intercepts the *data* pass: as onNext flows down through it, it re-dispatches everything **below** onto its Scheduler until the next publishOn — so **position matters** and multiples stack. To trace a mixed chain: the source and operators above the first publishOn run on the (closest-to-source) subscribeOn Scheduler, or the subscribing thread if none; then each publishOn switches downstream to its Scheduler for the operators beneath it. In WebFlux, 'subscribing thread' means the Netty event loop, which is why offloading uses subscribeOn(boundedElastic) at the source.

code

java · 18 lines
java
import reactor.core.publisher.Flux;
import reactor.core.scheduler.Schedulers;

Flux.range(1, 2)                                   // source
    .map(i -> tag("A", i))                          // segment 1
    .publishOn(Schedulers.parallel())               // -> switch downstream
    .map(i -> tag("B", i))                          // segment 2 (parallel)
    .publishOn(Schedulers.single())                 // -> switch again
    .map(i -> tag("C", i))                          // segment 3 (single)
    .subscribeOn(Schedulers.boundedElastic())       // position-independent: sets SOURCE
    .subscribe(i -> tag("sub", i));                 // runs on 'single' (last publishOn)

// Threads:
//   A   -> boundedElastic-1   (source + up to first publishOn; from subscribeOn)
//   B   -> parallel-1         (after publishOn(parallel))
//   C   -> single-1           (after publishOn(single))
//   sub -> single-1           (last active Scheduler)
// Note: the subscribeOn sits at the BOTTOM yet still governs the source thread.

go deeper

for a junior

Recall the two one-liners; not expected to explain the signal-direction mechanism.

for a middle

Trace a mixed chain's thread segments correctly.

for a senior

Explain via upward-subscription vs downward-data and publishOn's async boundary/backpressure.

for a principal

Discuss subscribeOn no-ops on async sources, Context propagating upward, default-parallel time operators, and ordering implications across publishOn boundaries.

## Two signal directions Every Reactor pipeline is assembled top-to-bottom but *executes* via two opposite signal flows once you subscribe: 1. **Subscription (control) — upward.** Calling `subscribe()` sends a subscribe signal that propagates *up* the chain, operator by operator, until it reaches the source, which then starts producing. 2. **Data (onNext/onComplete/onError) — downward.** The source emits, and each element flows *down* through the operators to the subscriber. `subscribeOn` and `publishOn` each latch onto one of these directions — that's the whole reason for their different positional behavior. ### `subscribeOn` = intercepts the upward subscription When the subscribe signal, travelling upward, hits a `subscribeOn(S)`, that operator arranges for the **rest of the upward subscription (and therefore the source's work/emission) to happen on Scheduler S**. Because it acts on the upward pass — which always ends at the source no matter where the operator physically sits — **its position in the chain is irrelevant** to *which thread the source uses*. And if there are several `subscribeOn`s, the subscribe signal hits them from bottom to top; the **last one it processes before reaching the source (i.e., the one closest to the source) determines** the source thread — the earlier (lower) ones are overridden. Hence: *closest-to-source wins, position otherwise doesn't matter*. ### `publishOn` = intercepts the downward data `publishOn(S)` sits in the downward data path. As each onNext passes through it, it **hands the element to Scheduler S** and everything **below** it runs on S — until another `publishOn` changes it again. Because it acts on the downward pass, **only operators physically below it are affected**, so **position matters** and multiple `publishOn`s create successive thread segments. ## Tracing algorithm for a mixed chain 1. Find the **closest-to-source `subscribeOn`** (if any). The source and every operator **down to the first `publishOn`** run on that Scheduler. If there's no `subscribeOn`, they run on the **subscribing thread** (the Netty event loop in WebFlux). 2. At each **`publishOn(S)`**, switch: operators from there **down to the next `publishOn`** run on S. 3. The **subscriber's** `onNext`/`onComplete` callbacks run on whatever the last active Scheduler is (the final `publishOn`, or the subscribeOn/subscribing thread if no publishOn). ## Subtleties a principal should raise - **`subscribeOn` may be a no-op for emission thread on already-async sources.** If the source emits on its *own* threads (a non-blocking driver, `interval`, an external callback), `subscribeOn` controls the *subscription* thread but the source may still deliver onNext on its native thread. `subscribeOn` truly matters for **synchronous/blocking** sources (`fromCallable`, `range`, iterable), where the emitting thread *is* the subscribing thread. - **`publishOn` introduces an async boundary** with an internal bounded queue (Queues.SMALL_BUFFER_SIZE by default) and participates in backpressure — it can affect latency, batching and ordering guarantees across the boundary. - **`flatMap`/`concatMap` inner publishers** may run on their own Schedulers, so a chain's effective threading isn't only about the outer subscribeOn/publishOn. - **Context** (`Context`/`contextWrite`) also propagates *upward* like subscription — a useful parallel that reinforces the mental model. - Operators such as `interval`, `timeout`, `delayElements` default to `Schedulers.parallel()`, quietly introducing a thread switch even without an explicit publishOn. ## Why WebFlux cares Because 'subscribing thread' is the shared Netty event loop, the practical upshot is: put a single `subscribeOn(Schedulers.boundedElastic())` near a blocking source to lift the whole upstream off the event loop, and use `publishOn` to steer specific downstream stages — knowing exactly which operators each affects.

  • You put subscribeOn at the very bottom of the chain — does the source still run on it?
    Yes. The subscribe signal travels upward through it to the source, so the source runs on that Scheduler regardless of the operator's physical position — assuming it's the closest-to-source subscribeOn.
  • When does subscribeOn NOT change the thread an operator observes onNext on?
    When the source is already asynchronous (emits on its own threads) or when a downstream publishOn has already switched threads. subscribeOn only sets the source/subscription thread, not downstream-of-publishOn segments.
  • An operator you didn't add is running on parallel-N — why?
    Time-based operators like interval, timeout, delayElements, and delaySubscription default to Schedulers.parallel(), introducing a thread switch even without an explicit publishOn.

saying these in an interview costs you the question

  • 'subscribeOn only affects operators after it in the chain'
  • 'the last subscribeOn wins' (it's the closest to the source)
  • 'publishOn changes the source's emission thread'
  • Assuming subscribeOn changes emission thread on an already-async source
  • Forgetting interval/delay operators default to parallel

context