skip to content

Schedulers & subscribeOn/publishOn

Schedulers decide which threads run your work, with subscribeOn moving the whole chain and publishOn moving everything downstream. The practical use is offloading a blocking call to boundedElastic, and interviewers ask exactly that.

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

questions

5

What is a Reactor Scheduler, and what are the four built-in Schedulers (boundedElastic, parallel, single, immediate) each meant for?

level: juniorimportance: must knowfreq 70%

answer

  1. boundedElastic = blocking I/O, capped 10×cores
  2. parallel = CPU work, fixed = #cores
  3. single = one thread; immediate = no switch
  4. never block the event loop
  5. Mono.fromCallable + subscribeOn(boundedElastic)

basics

~10 s

A Scheduler decides which thread(s) run reactive work. boundedElastic is for blocking I/O, parallel for CPU work, single for one shared thread, and immediate runs on the current thread (no switch).

solid answer

~40 s

A `Scheduler` is Reactor's abstraction over thread pools; operators like `subscribeOn`/`publishOn` take one to move work off the calling thread. `Schedulers.boundedElastic()` — a capped, elastic pool for blocking/long I/O (JDBC, blocking HTTP); it caps threads (default ~10× CPU cores) and queues extra tasks, evicting idle threads. `Schedulers.parallel()` — a fixed pool sized to CPU cores for short CPU-bound work; never block on it. `Schedulers.single()` — one reusable thread for low-latency serialized tasks. `Schedulers.immediate()` — a no-op that runs on the current thread, useful as a null-object when an API requires a Scheduler but you want no switch. In WebFlux the golden rule is: keep the event-loop threads non-blocking and push any blocking call to boundedElastic.

code

java · 13 lines
java
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;

// Blocking JDBC call wrapped and pushed off the event loop
Mono<User> user = Mono.fromCallable(() -> userJdbcDao.findById(id)) // blocking
    .subscribeOn(Schedulers.boundedElastic());                     // run on I/O pool

// CPU-bound, non-blocking transform on the parallel pool
Mono<BigInteger> hash = Mono.fromSupplier(() -> expensiveHash(data))
    .subscribeOn(Schedulers.parallel());

// immediate(): required by an API but we want NO thread switch
flux.publishOn(Schedulers.immediate()); // stays on current thread

go deeper

for a junior

Name the four Schedulers and the one-line purpose of each; know 'blocking → boundedElastic'.

for a middle

Explain the bound (10×cores + queue) and why the event loop must stay non-blocking.

for a senior

Discuss isolating heavy subsystems with newBoundedElastic, daemon-thread naming, BlockHound detection.

for a principal

Reason about pool sizing under SLA/backpressure, thread-starvation deadlocks, and when to wrap a custom ExecutorService.

## What a Scheduler is Reactor is thread-agnostic: by default a reactive pipeline runs on whatever thread called `subscribe()` (or emitted the signal). A `reactor.core.scheduler.Scheduler` is an abstraction over a source of threads (usually a thread pool) that lets you *explicitly* decide where work runs. You never create threads by hand — you hand a `Scheduler` to `subscribeOn(...)` or `publishOn(...)` and Reactor dispatches the signals onto that Scheduler's worker(s). `Schedulers` (the factory class) exposes shared, lazily-initialised singletons. ## The four built-ins **`Schedulers.boundedElastic()`** — the workhorse for **blocking** or long-running I/O (legacy JDBC, `RestTemplate`, blocking file/network calls). It's *elastic*: it creates worker threads on demand and reuses them, but it's *bounded* — by default the thread cap is **10 × the number of CPU cores**, and once all threads are busy, extra tasks are **queued** (default queue cap 100 000) rather than spawning unbounded threads. Idle threads are evicted after 60s. The bound is deliberate: it prevents the classic "thread explosion" that unbounded elastic pools caused, while still letting blocking work park a thread without starving the rest of the app. **`Schedulers.parallel()`** — a **fixed-size** pool with one thread per CPU core (`Runtime.availableProcessors()`), tuned for **short, non-blocking, CPU-bound** work (computation, fan-out with `flatMap`, parallel streams). Because it's small and fixed, **blocking on it starves the whole pool** and can deadlock the app. **`Schedulers.single()`** — a **single** reusable thread for tasks that must be serialized or need a dedicated low-latency thread (e.g., a periodic timer). `Schedulers.newSingle(...)` gives a fresh private one. **`Schedulers.immediate()`** — not a pool at all: it executes the task **synchronously on the current thread**. It's a *null-object* — used when an operator requires a `Scheduler` argument but you explicitly want no thread switch. There's also `Schedulers.fromExecutor(...)`/`fromExecutorService(...)` to wrap your own pool. ## Why this matters in WebFlux Spring WebFlux serves requests on a **small number of non-blocking event-loop threads** (Reactor Netty). If you block one of those threads (a synchronous DB/HTTP call), you stall every request that thread was multiplexing — throughput collapses. The rule: **never block the event loop; offload blocking calls to `boundedElastic`** (via `subscribeOn`/`publishOn`). CPU-bound work that's non-blocking can go to `parallel`. ## Gotchas - boundedElastic threads are **daemon** threads named `boundedElastic-…`; parallel threads are `parallel-…` — handy when reading stack traces/logs. - The bound means boundedElastic is **not infinite**: if you flood it with slow blocking calls you can still exhaust it (queued tasks pile up). Size/isolate with `Schedulers.newBoundedElastic(...)` for heavy subsystems. - Schedulers are shared singletons; `Schedulers.parallel()` returns the *same* instance every call. Don't `dispose()` the global ones. - `Mono.fromCallable(blockingCall).subscribeOn(Schedulers.boundedElastic())` is the canonical wrapper for a blocking method.

  • Why is boundedElastic bounded instead of unbounded elastic?
    Unbounded elastic pools spawned a thread per blocking task and could exhaust OS threads/memory under load. boundedElastic caps threads (default 10×cores) and queues overflow tasks, so blocking work parks a thread safely without a thread explosion.
  • What happens if you run a blocking call on Schedulers.parallel()?
    You occupy one of the few fixed CPU-core threads. Under concurrency you starve the pool, stalling other non-blocking work and risking deadlock. BlockHound can detect such blocking calls on non-blocking Schedulers.

saying these in an interview costs you the question

  • Thinking boundedElastic is unbounded/infinite
  • Using parallel() for blocking JDBC/HTTP calls
  • Believing immediate() uses a thread pool
  • Assuming parallel() grows under load like a cached pool

context

open as a page

What is the difference between subscribeOn and publishOn, and how does each affect which thread the operators run on?

level: middleimportance: must knowfreq 85%

basics

~10 s

publishOn switches the thread for operators placed after it (downstream). subscribeOn sets the thread for the source and the whole subscription, no matter where you put it in the chain.

open as a page

You must call a legacy blocking library (JDBC / RestTemplate) inside a WebFlux handler. How do you keep from blocking the event loop, and what does the correct code look like?

level: seniorimportance: must knowfreq 80%

basics

~10 s

Wrap the blocking call in Mono.fromCallable(...) (or Flux) and add subscribeOn(Schedulers.boundedElastic()), so the blocking work runs on the I/O pool instead of the Netty event-loop thread.

open as a page

Why is Schedulers.parallel() the wrong place for blocking calls, and boundedElastic the wrong place for tight CPU-bound work? What can go wrong?

level: seniorimportance: should knowfreq 55%

basics

~20 s

parallel() has only one thread per CPU core, so blocking those threads starves the pool and can deadlock. boundedElastic can hold many threads, so putting CPU work there causes oversubscription and context-switching instead of speedup.

open as a page

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%

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.

open as a page