skip to content

Client Development and Integration Patterns

Writing real applications on Kafka: client lifecycle, Spring Kafka, retries and dead-letter topics, outbox and idempotent consumers, concurrency, and testing. Interviewers spend time here because this is where most developers actually meet Kafka.

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

questions

62 · 11 sections

How do you construct a KafkaProducer and a KafkaConsumer in the Java client, and what are the minimum required configuration properties for each?

level: juniorimportance: must knowfreq 75%
basics
~20 s

Build a Properties (or Map) of config, then pass it to the constructor: new KafkaProducer<>(props) or new KafkaConsumer<>(props). Producers need bootstrap.servers plus a key and value serializer; consumers need bootstrap.servers plus a key and value deserializer (and usually group.id).

open as a page

Describe the consumer poll loop. Why must you call poll() regularly, and what is the role of max.poll.interval.ms?

level: middleimportance: must knowfreq 78%
basics
~20 s

After subscribing, you loop calling poll(timeout) to fetch records, process them, then poll again. poll() also drives group membership/heartbeats. If you take too long between polls (longer than max.poll.interval.ms) the broker assumes you're stuck, removes you from the group, and rebalances your partitions to others.

open as a page

Are KafkaProducer and KafkaConsumer thread-safe? How does that shape how you use each across threads?

level: middleimportance: must knowfreq 80%
basics
~10 s

KafkaProducer is thread-safe — share one instance across many threads. KafkaConsumer is NOT thread-safe — it must be used by a single thread; the only safe cross-thread call is wakeup().

open as a page

Explain the producer's send() async model and the roles of flush() and close(). What can go wrong if you skip them?

level: seniorimportance: must knowfreq 65%
basics
~20 s

send() is asynchronous: it buffers the record and returns a Future immediately; a background Sender thread actually transmits batches. flush() blocks until all buffered records have been sent and acknowledged. close() flushes then releases resources. Skip them and you can lose unsent buffered records on exit.

open as a page

What are serializers and deserializers in the Kafka Java client, and what happens when a deserializer hits a bad (poison) record?

level: middleimportance: should knowfreq 55%
basics
~20 s

Kafka moves only bytes. A Serializer<T> turns your key/value object into byte[] before producing; a Deserializer<T> turns byte[] back into an object when consuming. A bad record that the deserializer can't parse throws inside poll(), and by default the consumer keeps failing on that same offset — a 'poison pill' that blocks progress until you handle it.

open as a page

What is KafkaTemplate in Spring for Apache Kafka, and how do send() and sendDefault() differ?

level: juniorimportance: must knowfreq 72%
basics
~10 s

KafkaTemplate is Spring's helper for producing records to Kafka. send(topic, key, value) targets a named topic; sendDefault(key, value) uses a default topic preconfigured on the template, so you don't repeat the topic name.

open as a page

How does @KafkaListener work, and what does the concurrency setting actually control on a MessageListenerContainer?

level: middleimportance: must knowfreq 68%
basics
~20 s

@KafkaListener marks a method to consume records from topics. Spring creates a MessageListenerContainer that runs the poll loop and invokes your method. concurrency=N creates N consumer threads (each a separate KafkaConsumer), so up to N partitions are processed in parallel.

open as a page

Explain the container AckMode values (RECORD, BATCH, MANUAL/MANUAL_IMMEDIATE) and how each affects offset commits.

level: seniorimportance: must knowfreq 60%
basics
~20 s

AckMode controls when the container commits consumer offsets. RECORD commits after each record's listener returns. BATCH commits once after the whole poll batch is processed (the default). MANUAL/MANUAL_IMMEDIATE hand control to you via an Acknowledgment.acknowledge() call.

open as a page

What roles do ProducerFactory and ConsumerFactory play, and why configure serializers/deserializers there?

level: middleimportance: should knowfreq 50%
basics
~10 s

ProducerFactory builds and caches KafkaProducer instances from a config map (bootstrap servers, serializers); KafkaTemplate uses it. ConsumerFactory builds KafkaConsumer instances (with deserializers) that listener containers use. Serializers/deserializers live there because they're producer/consumer-level settings.

open as a page

How do you enable a batch @KafkaListener, and what changes about delivery and error handling compared to record mode?

level: seniorimportance: should knowfreq 38%
basics
~20 s

Set batchListener=true on the container factory; the listener method then takes a List of records (one poll's worth) instead of a single record. You process them together, and error handling/retry applies to the whole batch unless you use index-aware handlers.

open as a page

What is DefaultErrorHandler in Spring for Apache Kafka, and what does it do when a listener throws an exception?

level: juniorimportance: must knowfreq 70%
basics
~20 s

DefaultErrorHandler is the standard component that catches exceptions thrown by your @KafkaListener. It retries the record a configurable number of times with a backoff delay, and if all retries fail it calls a recoverer (by default just logs the failed record).

open as a page

A non-deserializable 'poison-pill' record keeps crashing your consumer in an infinite loop. Why does this happen with naive handling, and how do you fix it?

level: middleimportance: must knowfreq 55%
basics
~20 s

A poison pill is a record that can't be deserialized, so it throws before your listener even runs. The consumer keeps re-reading the same offset and looping forever. Fix it with ErrorHandlingDeserializer, which catches the failure and lets the error handler skip the record to a DLT.

open as a page

Contrast blocking retries (DefaultErrorHandler) with non-blocking retries (@RetryableTopic). When would you choose one over the other?

level: seniorimportance: must knowfreq 65%
basics
~20 s

Blocking retries re-process the failed record on the same consumer thread, holding up the partition until it succeeds. Non-blocking retries (@RetryableTopic) forward the record to separate retry topics with delays, so the main topic keeps flowing while failed records are retried independently.

open as a page

Compare FixedBackOff and ExponentialBackOff for retry spacing. How do you cap attempts with ExponentialBackOff, and why add jitter?

level: middleimportance: should knowfreq 40%
basics
~20 s

FixedBackOff retries at a constant interval for a set number of attempts. ExponentialBackOff grows the delay each time (×multiplier up to a max). To bound retries with exponential delays, use ExponentialBackOffWithMaxRetries. Jitter (randomness) spreads retries out to avoid all consumers hammering a recovering dependency at once.

open as a page

Walk through configuring @RetryableTopic: attempts, backoff topics, exception classification, and where the DLT fits.

level: middleimportance: should knowfreq 50%
basics
~20 s

@RetryableTopic on a listener tells Spring to auto-create retry topics with increasing delays and a final DLT. You set attempts, the backoff (delay/multiplier), which exceptions to include/exclude, and a @DltHandler method to process records that exhausted all retries.

open as a page

What is the Transactional Outbox pattern, and what problem does it solve when you need to update a database and publish a Kafka message in one operation?

level: juniorimportance: must knowfreq 78%
basics
~20 s

Instead of writing to the database AND sending to Kafka separately (which can half-fail), you write the business row and an 'outbox' row in the same DB transaction. A separate process later reads the outbox and publishes to Kafka.

open as a page

Outbox delivery is at-least-once, so consumers may see duplicate Kafka messages. How do you make a consumer idempotent?

level: middleimportance: must knowfreq 68%
basics
~20 s

Give each event a unique id. The consumer records processed ids (e.g. in a 'processed_messages' table) inside the same transaction that applies the effect, and skips any id it has already seen. That way reprocessing a duplicate does nothing.

open as a page

Compare a polling-publisher relay versus a CDC/Debezium relay for moving outbox rows to Kafka. What are the trade-offs?

level: middleimportance: must knowfreq 70%
basics
~20 s

A polling publisher repeatedly queries the outbox table for new rows and sends them to Kafka. CDC (e.g. Debezium) instead tails the DB's transaction log and streams outbox inserts with no polling. CDC is lower-latency and lower DB load; polling is simpler and needs no extra infrastructure.

open as a page

What columns belong in a well-designed outbox table, and what does each one enable?

level: juniorimportance: should knowfreq 48%
basics
~20 s

Typical columns: a unique id (UUID), aggregate type (which topic), aggregate id (the Kafka key/ordering), event type, payload (the message body, often JSON), a timestamp, and optionally a published flag or sequence. These let the relay route, key, order, and dedup events.

open as a page

How do you preserve per-aggregate message ordering when publishing outbox events to Kafka, and where can ordering break?

level: seniorimportance: should knowfreq 55%
basics
~20 s

Kafka only orders messages within a single partition. Use the aggregate id (e.g. order id) as the Kafka message key so all events for one entity land on the same partition in order. Publish them in commit/sequence order from the relay.

open as a page

Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka redelivers messages after rebalances or restarts (at-least-once), so the same record can arrive more than once. You record each processed message's id in a store (a DB unique constraint or Redis key) and skip any id you've already seen.

open as a page

Walk through the read-process-write ordering that makes a dedup store actually safe under consumer crashes. Where exactly do you record the key, side effect, and offset?

level: seniorimportance: must knowfreq 50%
basics
~20 s

Make the side effect and the dedup-key write part of the same atomic transaction, so either both happen or neither does. Commit the Kafka offset only after that transaction succeeds. On redelivery, the existing key tells you to skip.

open as a page

Why must the dedup check-and-record be a single atomic operation, and what bug appears if you use SELECT-then-INSERT or GET-then-SET?

level: middleimportance: should knowfreq 45%
basics
~20 s

If you check 'have I seen this key?' and then separately record it, two concurrent consumers can both pass the check before either records, so both process the message. You need one atomic step: a UNIQUE-constraint INSERT or Redis SET NX that checks and records together.

open as a page

When should you key your dedup store on a business idempotency key versus (topic, partition, offset)?

level: middleimportance: should knowfreq 55%
basics
~20 s

Use a business key (like orderId) when the same logical event can reappear at a different offset or on another topic — it survives re-publishing. Use (topic,partition,offset) only when there's no natural business id and duplicates come purely from Kafka redelivery.

open as a page

How would you implement a Redis-based dedup store with SETNX, and what TTL/cleanup and durability concerns must you handle?

level: seniorimportance: should knowfreq 40%
basics
~20 s

Use SET key value NX EX <ttl>: it sets the key only if absent (NX) and auto-expires after the TTL (EX). If the command returns 'not set', you've seen the message — skip it. TTL bounds memory; pick it longer than any plausible duplicate window.

open as a page

What is the difference between an event and a command in an event-driven system, and how does that distinction shape how you name and publish messages to Kafka topics?

level: juniorimportance: must knowfreq 78%
basics
~20 s

An event states something that already happened (past tense, e.g. OrderPlaced) and has no specific intended recipient. A command tells one service to do something (imperative, e.g. PlaceOrder) and expects an action. Events go to topics for any consumer; commands target one handler.

open as a page

Compare fat (event-carried state transfer) events versus thin (notification) events. What are the trade-offs, and when would you choose each for a Kafka topic?

level: middleimportance: must knowfreq 70%
basics
~20 s

A thin event carries only an id/reference, so consumers must call back to fetch details. A fat event carries the full state, so consumers need no callback but the payload is larger and may go stale. Fat reduces coupling/load; thin keeps payloads small and authoritative.

open as a page

Explain choreography versus orchestration for coordinating a multi-service business process (e.g. an order saga) over Kafka. What are the trade-offs and failure-handling implications of each?

level: seniorimportance: must knowfreq 68%
basics
~20 s

In choreography each service reacts to events and emits its own, with no central coordinator — logic is distributed. In orchestration a central orchestrator tells each service what to do and tracks progress. Choreography is decoupled but hard to trace; orchestration is centralized, explicit, and easier to monitor.

open as a page

How do you evolve the schema of a domain event published to a Kafka topic without breaking existing consumers? Discuss compatibility modes and what changes are safe.

level: seniorimportance: must knowfreq 72%
basics
~20 s

Use a schema registry (Avro/Protobuf/JSON Schema) and a compatibility mode. The safest common choice is BACKWARD: new consumers can read old data. Safe changes are adding optional fields with defaults and removing optional fields; renaming or removing required fields breaks compatibility.

open as a page

How should you design Kafka topics for domain events — granularity (one event type per topic vs many), keying, and retention/compaction — and what is a domain event versus an integration event?

level: middleimportance: should knowfreq 55%
basics
~20 s

A domain event captures something meaningful in the business domain (OrderPlaced). Design topics around an aggregate/entity, key events by the aggregate id for ordering, choose one-type-per-topic for clarity or grouped types for related events, and use compaction for state-style events and time retention for pure notifications.

open as a page

When testing a Kafka application, how do you decide what belongs in a unit test versus an integration test, and which tools fit each side?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Unit tests check your logic without a real broker, using MockProducer/MockConsumer or TopologyTestDriver. Integration tests run against a real (embedded or containerized) broker to verify serialization, partitioning, and offset behavior end to end.

open as a page

Compare @EmbeddedKafka and Testcontainers KafkaContainer for integration testing. When would you choose each?

level: middleimportance: must knowfreq 65%
basics
~10 s

@EmbeddedKafka runs a broker inside the test JVM — fast, no Docker, but not production-identical. Testcontainers KafkaContainer runs a real Kafka image in Docker — true fidelity and version-matched, but slower and Docker-dependent.

open as a page

What are MockProducer and MockConsumer, and how would you use them to unit-test producer and consumer logic?

level: juniorimportance: should knowfreq 45%
basics
~20 s

They are in-memory fakes of the Kafka client. MockProducer records every send() in a history() list so you assert what was sent; MockConsumer lets you preload records and control poll() returns to test consumer logic — both with no broker.

open as a page

Why are Awaitility-style polling assertions essential in Kafka integration tests, and how do you write a robust one?

level: middleimportance: should knowfreq 55%
basics
~10 s

Kafka delivery is asynchronous, so a record sent now isn't consumed instantly. Awaitility polls an assertion repeatedly until it passes or a timeout expires — replacing flaky Thread.sleep with a bounded, retrying wait.

open as a page

How does TopologyTestDriver work for testing Kafka Streams, and what does it let you avoid?

level: middleimportance: should knowfreq 50%
basics
~10 s

TopologyTestDriver runs a Kafka Streams Topology in-process, synchronously, with no broker. You push records into TestInputTopic and read results from TestOutputTopic, so you can unit-test stream logic deterministically and fast.

open as a page

Is the Kafka KafkaConsumer thread-safe, and what is the canonical threading model for consuming from Kafka?

level: juniorimportance: must knowfreq 70%
basics
~10 s

No. The Java KafkaConsumer is not thread-safe — only one thread may call it (except wakeup()). The standard model is one consumer instance per thread, each in its own poll loop.

open as a page

How does the Spring Kafka ConcurrentMessageListenerContainer 'concurrency' setting map to partitions and threads?

level: middleimportance: must knowfreq 60%
basics
~10 s

concurrency=N creates N child containers, each with its own consumer and thread sharing the group. Kafka distributes partitions among them, so useful concurrency is capped at the partition count; extra child consumers idle.

open as a page

When you decouple the poll loop from a worker thread pool, what offset-commit hazard arises and how do you avoid losing or reprocessing records?

level: seniorimportance: must knowfreq 55%
basics
~20 s

Workers finish out of offset order, so committing the latest offset can skip records still in flight on slower workers. Avoid it by committing only the highest contiguous completed offset per partition, and pause/resume the partition to bound in-flight work.

open as a page

What problem does the Confluent Parallel Consumer solve, and what ordering guarantees do its KEY, PARTITION, and UNORDERED modes provide?

level: seniorimportance: should knowfreq 40%
basics
~20 s

It lets you process one topic with far more parallelism than partitions while managing offsets safely. UNORDERED = max parallelism, no order; KEY = parallel across keys but ordered per key; PARTITION = ordered per partition, parallel across partitions.

open as a page

When should you scale Kafka consumer throughput by adding partitions versus using a decoupled worker pool or Parallel Consumer, and what are the trade-offs?

level: principalimportance: should knowfreq 35%
basics
~20 s

Add partitions when work is CPU-bound and partition count is moderate, accepting it's a one-way change that can break keyed ordering. Use a decoupled pool or Parallel Consumer for I/O-bound work where you'd otherwise need impractically many partitions.

open as a page

How do you configure a Kafka producer and consumer to use Avro serialization with Confluent Schema Registry?

level: juniorimportance: must knowfreq 70%
basics
~10 s

Set the producer's value serializer to KafkaAvroSerializer and the consumer's value deserializer to KafkaAvroDeserializer, then point both at the registry with schema.registry.url. Keys often stay a StringSerializer.

open as a page

A bad record causes a deserialization exception that crashes your @KafkaListener in an infinite loop. How do you handle this poison pill?

level: seniorimportance: must knowfreq 65%
basics
~20 s

Wrap your deserializer in Spring Kafka's ErrorHandlingDeserializer. It catches the exception during deserialization, hands a null payload plus the failure to your error handler, and lets you skip or DLT the record instead of looping forever.

open as a page

What does the specific.avro.reader config do, and when do you set it to true?

level: middleimportance: should knowfreq 55%
basics
~10 s

specific.avro.reader=true makes KafkaAvroDeserializer return your generated SpecificRecord class (e.g. an Order POJO) instead of a generic GenericRecord. Set it when you have compiled Avro classes on the classpath.

open as a page

How do you configure Avro Serdes in a Kafka Streams application, and how does it differ from a plain @KafkaListener?

level: seniorimportance: should knowfreq 45%
basics
~10 s

Streams uses Serde objects (a serializer+deserializer pair), not separate serializer/deserializer classes. Set default.value.serde to SpecificAvroSerde (or GenericAvroSerde), give it schema.registry.url, and pass it explicitly in Consumed/Produced for repartition/state-store topics.

open as a page

What does auto.register.schemas do, and why might you disable it on producers in production?

level: principalimportance: should knowfreq 40%
basics
~20 s

auto.register.schemas (default true) lets the producer's serializer register a new schema in the registry on first use. In production you often set it false and pre-register schemas through a governed pipeline, plus set use.latest.version=true, so apps can't silently introduce schemas.

open as a page

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

level: juniorimportance: must knowfreq 55%
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.

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

What is the request-reply (RPC-style) messaging pattern in Kafka, and what extra pieces does it need on top of plain pub/sub?

level: juniorimportance: must knowfreq 55%
basics
~20 s

Request-reply makes one service send a message and wait for a matching response. On top of pub/sub you add a reply topic to send the answer back, and a correlation ID so the requester knows which reply belongs to which request.

open as a page

When should you choose Kafka request-reply (RPC-style) over plain pub/sub event choreography, and what are the architectural trade-offs?

level: seniorimportance: must knowfreq 50%
basics
~20 s

Use request-reply only when the caller truly needs the answer before it can proceed and you specifically want it to flow over Kafka. It re-adds temporal coupling, blocking, reply topics, correlation, and timeouts. Prefer pub/sub events when the caller can react later — it's looser, more resilient, and replayable. Often a synchronous HTTP/gRPC call is a better RPC than Kafka.

open as a page

Walk through how ReplyingKafkaTemplate.sendAndReceive() works end to end, including how the reply container, correlation, and the returned future fit together.

level: middleimportance: should knowfreq 45%
basics
~20 s

You build ReplyingKafkaTemplate from a producer factory plus a listener container on the reply topic. sendAndReceive() stamps a correlation ID, sets the reply-topic header, stores a future in a map keyed by that ID, sends the request, and returns the future. When the reply container receives a record with that ID, it completes the future.

open as a page

How do you implement the responder side of a Kafka request-reply exchange so replies are correctly correlated and routed, and what must you preserve from the request?

level: middleimportance: should knowfreq 38%
basics
~20 s

The responder consumes the request topic, does its work, and sends the result to the reply topic named in the request's REPLY_TOPIC header, copying the request's CORRELATION_ID onto the reply. In Spring, a @KafkaListener method that returns a value does this automatically via the container factory's reply template.

open as a page

What failure modes does the in-memory pending-request map introduce in a Kafka request-reply client, and how do you keep it bounded and correct?

level: seniorimportance: should knowfreq 35%
basics
~20 s

The pending map holds a future per outstanding request in JVM memory. Risks: entries leak if replies never arrive and there's no timeout, the map grows under load, late or duplicate replies have no owner after eviction, and futures are lost on a process crash. Bound it with per-request timeouts, eviction on completion/timeout, and treat lost requests as RPC timeouts.

open as a page