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 pageshowhide
explore
- Java Client Configuration and Lifecycle6 questions
- Error Handling, Retry Topics and DLT6 questions
- Transactional Outbox Pattern6 questions
- Idempotent Consumers and Deduplication5 questions
- Event-Driven Design with Kafka6 questions
- Testing Kafka Applications6 questions
- Consumer Concurrency and Threading Models5 questions
- Reactive and Async Kafka Clients6 questions
- Request-Reply and Correlation Patterns5 questions
questions
62 · 11 sectionsHow do you construct a KafkaProducer and a KafkaConsumer in the Java client, and what are the minimum required configuration properties for each?
basics
~20 sBuild 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).
Describe the consumer poll loop. Why must you call poll() regularly, and what is the role of max.poll.interval.ms?
basics
~20 sAfter 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.
Are KafkaProducer and KafkaConsumer thread-safe? How does that shape how you use each across threads?
basics
~10 sKafkaProducer 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().
Explain the producer's send() async model and the roles of flush() and close(). What can go wrong if you skip them?
basics
~20 ssend() 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.
What are serializers and deserializers in the Kafka Java client, and what happens when a deserializer hits a bad (poison) record?
basics
~20 sKafka 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.
Spring for Apache Kafka and Spring Cloud Stream
all 6 Spring for Apache Kafka and Spring Cloud Stream questions →What is KafkaTemplate in Spring for Apache Kafka, and how do send() and sendDefault() differ?
basics
~10 sKafkaTemplate 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.
How does @KafkaListener work, and what does the concurrency setting actually control on a MessageListenerContainer?
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.
Explain the container AckMode values (RECORD, BATCH, MANUAL/MANUAL_IMMEDIATE) and how each affects offset commits.
basics
~20 sAckMode 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.
What roles do ProducerFactory and ConsumerFactory play, and why configure serializers/deserializers there?
basics
~10 sProducerFactory 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.
How do you enable a batch @KafkaListener, and what changes about delivery and error handling compared to record mode?
basics
~20 sSet 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.
What is DefaultErrorHandler in Spring for Apache Kafka, and what does it do when a listener throws an exception?
basics
~20 sDefaultErrorHandler 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).
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?
basics
~20 sA 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.
Contrast blocking retries (DefaultErrorHandler) with non-blocking retries (@RetryableTopic). When would you choose one over the other?
basics
~20 sBlocking 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.
Compare FixedBackOff and ExponentialBackOff for retry spacing. How do you cap attempts with ExponentialBackOff, and why add jitter?
basics
~20 sFixedBackOff 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.
Walk through configuring @RetryableTopic: attempts, backoff topics, exception classification, and where the DLT fits.
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.
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?
basics
~20 sInstead 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.
Outbox delivery is at-least-once, so consumers may see duplicate Kafka messages. How do you make a consumer idempotent?
basics
~20 sGive 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.
Compare a polling-publisher relay versus a CDC/Debezium relay for moving outbox rows to Kafka. What are the trade-offs?
basics
~20 sA 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.
What columns belong in a well-designed outbox table, and what does each one enable?
basics
~20 sTypical 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.
How do you preserve per-aggregate message ordering when publishing outbox events to Kafka, and where can ordering break?
basics
~20 sKafka 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.
Why does a Kafka consumer often need an application-level deduplication store, and what is the simplest way to build one?
basics
~20 sKafka 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.
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?
basics
~20 sMake 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.
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?
basics
~20 sIf 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.
When should you key your dedup store on a business idempotency key versus (topic, partition, offset)?
basics
~20 sUse 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.
How would you implement a Redis-based dedup store with SETNX, and what TTL/cleanup and durability concerns must you handle?
basics
~20 sUse 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.
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?
basics
~20 sAn 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.
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?
basics
~20 sA 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.
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?
basics
~20 sIn 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.
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.
basics
~20 sUse 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.
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?
basics
~20 sA 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.
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?
basics
~10 sUnit 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.
Compare @EmbeddedKafka and Testcontainers KafkaContainer for integration testing. When would you choose each?
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.
What are MockProducer and MockConsumer, and how would you use them to unit-test producer and consumer logic?
basics
~20 sThey 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.
Why are Awaitility-style polling assertions essential in Kafka integration tests, and how do you write a robust one?
basics
~10 sKafka 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.
How does TopologyTestDriver work for testing Kafka Streams, and what does it let you avoid?
basics
~10 sTopologyTestDriver 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.
Consumer Concurrency and Threading Models
all 5 Consumer Concurrency and Threading Models questions →Is the Kafka KafkaConsumer thread-safe, and what is the canonical threading model for consuming from Kafka?
basics
~10 sNo. 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.
How does the Spring Kafka ConcurrentMessageListenerContainer 'concurrency' setting map to partitions and threads?
basics
~10 sconcurrency=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.
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?
basics
~20 sWorkers 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.
What problem does the Confluent Parallel Consumer solve, and what ordering guarantees do its KEY, PARTITION, and UNORDERED modes provide?
basics
~20 sIt 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.
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?
basics
~20 sAdd 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.
Schema Registry and Serialization Integration
all 5 Schema Registry and Serialization Integration questions →How do you configure a Kafka producer and consumer to use Avro serialization with Confluent Schema Registry?
basics
~10 sSet 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.
A bad record causes a deserialization exception that crashes your @KafkaListener in an infinite loop. How do you handle this poison pill?
basics
~20 sWrap 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.
What does the specific.avro.reader config do, and when do you set it to true?
basics
~10 sspecific.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.
How do you configure Avro Serdes in a Kafka Streams application, and how does it differ from a plain @KafkaListener?
basics
~10 sStreams 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.
What does auto.register.schemas do, and why might you disable it on producers in production?
basics
~20 sauto.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.
What is Reactor Kafka, and how do KafkaSender and KafkaReceiver differ from the plain KafkaProducer/KafkaConsumer?
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.
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?
basics
~20 sBuild 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.
How does manual acknowledgement work in a Reactor Kafka stream, and what is the role of ReceiverOffset.acknowledge() vs commit()?
basics
~20 sEach 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.
What is the Vert.x Kafka client, and how does its flow-control model (pause/resume/fetch) compare to Reactor Kafka?
basics
~20 sThe 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.
How does backpressure propagate in Reactor Kafka, and how is it implemented in terms of consumer pause/resume?
basics
~20 sWhen 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.
What is the request-reply (RPC-style) messaging pattern in Kafka, and what extra pieces does it need on top of plain pub/sub?
basics
~20 sRequest-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.
When should you choose Kafka request-reply (RPC-style) over plain pub/sub event choreography, and what are the architectural trade-offs?
basics
~20 sUse 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.
Walk through how ReplyingKafkaTemplate.sendAndReceive() works end to end, including how the reply container, correlation, and the returned future fit together.
basics
~20 sYou 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.
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?
basics
~20 sThe 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.
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?
basics
~20 sThe 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.