skip to content

Reactive and Async Kafka Clients

Reactive Kafka clients and how backpressure propagates through a non-blocking pipeline with manual acknowledgement. Interviewers ask when the stack is Reactor or Vert.x and a blocking poll loop does not fit.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

6

What is Reactor Kafka, and how do KafkaSender and KafkaReceiver differ from the plain KafkaProducer/KafkaConsumer?

level: juniorimportance: must knowfreq 55%

answer

  1. Reactor types over real Kafka clients
  2. KafkaSender.send(Flux) -> Flux<SenderResult>
  3. KafkaReceiver.receive() -> Flux<ReceiverRecord>
  4. no while-true poll, no .get()
  5. same org.apache.kafka configs underneath

basics

~20 s

Reactor Kafka is a library that wraps Kafka's producer/consumer in Project Reactor types. KafkaSender publishes records as a Flux and returns results as a Flux; KafkaReceiver exposes incoming records as a Flux you subscribe to, instead of a blocking poll loop.

solid answer

~40 s

Reactor Kafka adapts the Apache Kafka clients to Project Reactor's reactive types (Mono/Flux). KafkaSender wraps a KafkaProducer: you call sender.send(Flux<SenderRecord>) and get back a Flux<SenderResult> with metadata or errors per record, so sends are non-blocking and composable. KafkaReceiver wraps a KafkaConsumer: receiver.receive() returns a Flux<ReceiverRecord> that drives the consumer's poll loop for you under the hood on a dedicated scheduler. Versus the raw clients, you no longer write an imperative while-true poll loop or block on producer.send().get(); instead you compose operators (map, flatMap, concatMap) and rely on the reactive runtime for scheduling and backpressure. The underlying clients are still the standard org.apache.kafka clients, so configs (bootstrap.servers, acks, group.id) are identical — Reactor Kafka only changes the programming model, not the protocol.

go deeper

for a junior

Know it wraps Kafka clients in Mono/Flux; sender sends a Flux of records, receiver gives a Flux of records.

for a middle

Explain SenderRecord/SenderResult correlation and the ReceiverRecord + ReceiverOffset handle for acks.

for a senior

Discuss the single-threaded consumer event loop, scheduling, and where blocking work breaks polling.

for a principal

Reason about when reactive Kafka is justified architecturally versus the operational simplicity of the plain client or Spring Kafka.

## Background Apache Kafka ships two core clients: **KafkaProducer** (sends records to topics) and **KafkaConsumer** (reads them). The standard usage is *imperative and blocking*: a producer call like `producer.send(record).get()` blocks the calling thread until the broker acknowledges, and a consumer runs a `while (true) { records = consumer.poll(timeout); process(records); }` loop on a single thread. **Project Reactor** is a library for *reactive programming* on the JVM. Its two main types are: - **Mono<T>** — a stream that emits 0 or 1 item then completes (or errors). - **Flux<T>** — a stream that emits 0..N items then completes (or errors). Reactive code is *declarative* (you describe a pipeline of operators) and *non-blocking* (threads aren't parked waiting for I/O). ## What Reactor Kafka is **Reactor Kafka** (`io.projectreactor.kafka:reactor-kafka`) is a thin adapter that exposes the standard Kafka clients through Reactor types. It does not reimplement the Kafka protocol — under the hood it still uses `org.apache.kafka.clients.producer.KafkaProducer` and `KafkaConsumer`, so all the familiar configs apply unchanged. ### KafkaSender - Created from `SenderOptions` (which holds the producer config map). - Core method: `send(Publisher<SenderRecord<K,V,T>>)` returns `Flux<SenderResult<T>>`. - A **SenderRecord** is a ProducerRecord plus an arbitrary *correlation* object `T` you attach; the matching **SenderResult** carries either the `RecordMetadata` (partition/offset) on success or an exception on failure, plus your correlation value so you can match results to inputs. - Sends are non-blocking: you never call `.get()`. Errors surface as `SenderResult.exception()` per record (not by throwing). ### KafkaReceiver - Created from `ReceiverOptions` (consumer config + subscription). - Core method: `receive()` returns `Flux<ReceiverRecord<K,V>>`. - Subscribing to that Flux starts a managed poll loop running on a dedicated single-threaded scheduler (the KafkaConsumer is not thread-safe, so all access is funnelled through one event thread). - Each **ReceiverRecord** is a ConsumerRecord plus a `ReceiverOffset` handle used for manual acknowledgement/commit. ## Why use it - Composability: chain `map`/`flatMap`/`concatMap` instead of hand-writing loops. - Non-blocking integration with reactive stacks (WebFlux, R2DBC, reactive HTTP clients). - Built-in **backpressure**: a slow downstream subscriber automatically slows the consumer's polling (the receiver pauses partitions) — something you'd have to implement manually with the raw client. ## Edge cases / gotchas - The KafkaConsumer remains single-threaded; doing blocking work inline on the receive Flux stalls polling. Offload with `publishOn`/`subscribeOn` or use `flatMap` with bounded concurrency. - Ordering: `flatMap` interleaves; use `concatMap` to preserve per-partition order. - It is still at-least-once by default; exactly-once needs transactions (`sender.transactionManager()`).

  • Does Reactor Kafka replace the Apache Kafka client or wrap it?
    It wraps it. The standard KafkaProducer/KafkaConsumer do the actual protocol work; Reactor Kafka only adapts them to Mono/Flux, so all native configs (acks, group.id, etc.) still apply.
  • How do you correlate a send result back to the record you sent?
    Attach a correlation value T when building the SenderRecord; the matching SenderResult exposes that same value via correlationMetadata(), letting you pair inputs and outputs even though sends complete out of order.

saying these in an interview costs you the question

  • Claiming Reactor Kafka reimplements the Kafka wire protocol (it doesn't; it wraps the standard clients).
  • Saying KafkaReceiver makes the underlying KafkaConsumer thread-safe (it stays single-threaded; access is serialized on one event thread).
  • Thinking send errors are thrown synchronously; they arrive as SenderResult.exception() per record.

context

open as a page

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%

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.

open as a page

How does manual acknowledgement work in a Reactor Kafka stream, and what is the role of ReceiverOffset.acknowledge() vs commit()?

level: seniorimportance: must knowfreq 45%

basics

~20 s

Each ReceiverRecord exposes a ReceiverOffset. Call acknowledge() after you finish processing a record to mark its offset eligible for committing; Reactor Kafka periodically commits acknowledged offsets in batches. commit() forces an immediate, explicit commit and returns a Mono.

open as a page

What is the Vert.x Kafka client, and how does its flow-control model (pause/resume/fetch) compare to Reactor Kafka?

level: middleimportance: should knowfreq 30%

basics

~20 s

The Vert.x Kafka client is the Eclipse Vert.x toolkit's non-blocking wrapper around the Kafka clients, exposing KafkaConsumer/KafkaProducer as Vert.x streams on the event loop. Flow control is explicit: you call pause(), resume(), and fetch(n) to control how many records you receive.

open as a page

How does backpressure propagate in Reactor Kafka, and how is it implemented in terms of consumer pause/resume?

level: seniorimportance: should knowfreq 40%

basics

~20 s

When a downstream subscriber requests fewer records than are available, Reactor Kafka stops the KafkaConsumer from fetching more by calling consumer.pause() on its partitions. When demand returns, it calls consumer.resume(). This keeps fast brokers from overwhelming slow processing.

open as a page

When is a reactive Kafka client (Reactor/Vert.x) actually justified over the plain client or Spring for Apache Kafka, and what are the pitfalls?

level: principalimportance: should knowfreq 28%

basics

~20 s

Reactive Kafka is justified when your service is already non-blocking end-to-end (WebFlux/Vert.x/R2DBC) and you want backpressure to flow across stages. If your processing is blocking or you just need a listener, Spring @KafkaListener or the plain client is simpler and just as performant.

open as a page