What is Reactor Kafka, and how do KafkaSender and KafkaReceiver differ from the plain KafkaProducer/KafkaConsumer?
answer
- Reactor types over real Kafka clients
- KafkaSender.send(Flux) -> Flux<SenderResult>
- KafkaReceiver.receive() -> Flux<ReceiverRecord>
- no while-true poll, no .get()
- same org.apache.kafka configs underneath
basics
~20 sReactor 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 sReactor 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
Know it wraps Kafka clients in Mono/Flux; sender sends a Flux of records, receiver gives a Flux of records.
Explain SenderRecord/SenderResult correlation and the ReceiverRecord + ReceiverOffset handle for acks.
Discuss the single-threaded consumer event loop, scheduling, and where blocking work breaks polling.
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.