skip to content

How do you perform a non-blocking Kafka send with Reactor Kafka using Mono/Flux, and why must you never call .block() inside a reactive pipeline?

level: middleimportance: must knowfreq 50%

answer

  1. Mono/Flux<SenderRecord> -> send -> Flux<SenderResult>
  2. cold/lazy: no subscribe = no send
  3. block() parks event-loop thread -> deadlock
  4. offload blocking via Schedulers.boundedElastic()
  5. errors as SenderResult.exception(), not thrown

basics

~20 s

Build a Flux of SenderRecords and pass it to sender.send(...), then subscribe to the returned Flux of results. Calling .block() parks the carrier thread, defeating non-blocking I/O and risking deadlock on the small event-loop thread pool.

solid answer

~40 s

With Reactor Kafka you create SenderRecords (a ProducerRecord plus a correlation key), wrap them in a Mono or Flux, and call kafkaSender.send(records), which returns a Flux<SenderResult>. Nothing happens until something subscribes — Reactor is lazy. You compose downstream with operators (doOnNext to log offsets, onErrorResume to handle failures) and let the framework (e.g. WebFlux) subscribe at the edge. You must never call .block() inside the pipeline because reactive runtimes use a tiny pool of event-loop threads; block() parks that carrier thread waiting for the result, which kills throughput and can deadlock when all event-loop threads are blocked waiting on work that needs those same threads to progress. For genuinely blocking dependencies, isolate them on a bounded elastic scheduler via subscribeOn/publishOn(Schedulers.boundedElastic()) rather than calling block().

code

java · 11 lines
java
SenderRecord<String, String, String> rec =
    SenderRecord.create(new ProducerRecord<>("orders", key, value), key);

Flux<SenderResult<String>> results =
    kafkaSender.send(Mono.just(rec))
        .doOnNext(r -> {
            if (r.exception() != null) log.error("send failed {}", r.correlationMetadata(), r.exception());
            else log.info("sent {} -> offset {}", r.correlationMetadata(), r.recordMetadata().offset());
        });

// return results to WebFlux (it subscribes). NEVER results.blockLast().

go deeper

for a junior

Know that you build a Flux of records, call send, and subscribe; don't block.

for a middle

Explain laziness, why block() is harmful on event-loop threads, and offloading via boundedElastic.

for a senior

Reason about deadlock conditions, maxInFlight backpressure, and ordering with flatMap vs concatMap.

for a principal

Define team guardrails (e.g. BlockHound) to catch accidental blocking and set scheduler sizing policy.

## The non-blocking send pattern In the **plain Kafka producer**, a synchronous send looks like `producer.send(record).get()` — `.get()` blocks the calling thread until the broker acks. That's fine on a thread-per-request server but ruinous on a reactive event-loop server. With **Reactor Kafka** you instead describe the send as a stream: 1. Build one or more `SenderRecord<K,V,T>` (a `ProducerRecord` + a correlation object `T`). 2. Wrap them: `Mono.just(record)` for one, or a `Flux` for many. 3. Call `kafkaSender.send(recordPublisher)` → returns `Flux<SenderResult<T>>`. 4. Compose downstream and return the Flux; the web framework subscribes. ### Laziness Reactor is **lazy/cold**: building the pipeline does nothing. The send only executes when a subscriber arrives (`.subscribe()`, or the framework subscribing to a controller's returned `Mono`/`Flux`). Forgetting to return/subscribe means the record is silently never sent. ## Why never call .block() Reactive frameworks (Netty/WebFlux) run on a **small fixed pool of event-loop threads** — typically a handful (often number-of-cores). These threads must never park; they continuously pick up ready work. `Mono.block()` *parks the current thread* until the value arrives. Consequences: - **Throughput collapse:** a parked event-loop thread can't service other requests. - **Deadlock:** if every event-loop thread is blocked waiting on results whose completion itself requires an event-loop thread (e.g. the Kafka send callback or downstream operator), nothing can make progress. Reactor even throws an explicit error if you call `block()` on a thread marked non-blocking (a `BlockingOperationError`). ## Correct handling of blocking dependencies Sometimes you genuinely must call blocking code (a JDBC call, a legacy SDK). Don't `block()`; instead **offload** it: - `Mono.fromCallable(() -> blockingCall()).subscribeOn(Schedulers.boundedElastic())` runs it on a separate elastic thread pool sized for blocking work, keeping the event loop free. - `publishOn(Schedulers.boundedElastic())` shifts *subsequent* operators onto that pool. ## Error handling Send failures don't throw — they arrive as `SenderResult.exception()`. Use `.doOnNext` to inspect each result, or rely on the Flux's error channel for terminal failures, and `onErrorResume`/`retryWhen` to recover. ## Ordering and concurrency For a Flux of many records, `send` preserves the order of the input Publisher into the producer. If you fan out per record with `flatMap`, you lose order; use `concatMap` to preserve it. ## Edge cases - **maxInFlight** in SenderOptions caps how many sends are outstanding — a backpressure lever. - Pipelines that never subscribe = silent no-op (a classic bug). - Mixing `block()` in a `@Scheduled` plain thread is technically allowed but still an anti-pattern that defeats the reactive model.

  • A developer's send 'does nothing' but throws no error. What's the most likely cause?
    Nobody subscribed to the returned Flux. Reactor pipelines are cold/lazy; without a subscriber (or returning it to a framework that subscribes), the send never executes.
  • You must enrich each record with a synchronous JDBC lookup. How do you do it without blocking the event loop?
    Wrap the JDBC call in Mono.fromCallable(...) and subscribeOn(Schedulers.boundedElastic()), so the blocking work runs on the elastic pool while event-loop threads stay free.

An event-loop thread is like a single waiter serving many tables. block() is the waiter standing frozen at one table waiting for the kitchen; offloading to boundedElastic is sending a runner to the kitchen so the waiter keeps serving everyone else.

saying these in an interview costs you the question

  • Calling .block() to 'get the result' inside a WebFlux handler.
  • Believing the send fires immediately when the Flux is built, without a subscriber.
  • Running blocking JDBC/legacy calls inline on the receive/send Flux instead of offloading to boundedElastic.
  • Expecting send errors to be thrown synchronously.

context