skip to content

Spring for Apache Kafka and Spring Cloud Stream

Spring's Kafka layer: KafkaTemplate, @KafkaListener containers, ack modes, and Spring Cloud Stream bindings. Extremely common in JVM interviews, because most Spring services meet Kafka through exactly this.

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

questions

6

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

level: juniorimportance: must knowfreq 72%

answer

  1. wrapper over KafkaProducer
  2. send(topic,...) vs sendDefault uses defaultTopic
  3. returns CompletableFuture<SendResult>
  4. async + thread-safe single bean
  5. built from ProducerFactory

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.

solid answer

~40 s

KafkaTemplate is a thin Spring wrapper over the native KafkaProducer that simplifies producing records and integrates with Spring (auto-config, message converters, observability). send(...) has overloads taking topic, optional partition, key, value, or a full ProducerRecord/Message<?>; sendDefault(...) drops the topic argument and uses the defaultTopic property set on the template. Sends are asynchronous: send() returns a CompletableFuture<SendResult> (ListenableFuture before 3.0). You add callbacks via whenComplete or block with get() if you need synchronous confirmation. The template is backed by a ProducerFactory that builds and pools the underlying producer. Because the Kafka producer is thread-safe, a single KafkaTemplate bean is shared across threads. You can also enable transactions on the factory to get send-within-transaction semantics via executeInTransaction.

code

java · 15 lines
java
kafkaTemplate.setDefaultTopic("orders");

// explicit topic
kafkaTemplate.send("orders", order.id(), order)
    .whenComplete((result, ex) -> {
        if (ex == null) {
            var md = result.getRecordMetadata();
            log.info("sent to {}-{}@{}", md.topic(), md.partition(), md.offset());
        } else {
            log.error("send failed", ex);
        }
    });

// default topic
kafkaTemplate.sendDefault(order.id(), order);

go deeper

for a junior

Know KafkaTemplate produces records; send names a topic, sendDefault uses a preset default topic.

for a middle

Explain the CompletableFuture<SendResult>, async batching, and thread-safe shared bean.

for a senior

Discuss ProducerFactory backing, message converters, transactions, and acks interplay with the future.

for a principal

Reason about throughput tradeoffs of blocking on get(), idempotence/ordering, and template configuration patterns at scale.

**Background.** Apache Kafka is a distributed log: producers append *records* (key + value + headers) to *topics*, which are split into *partitions*. The raw Java client exposes `KafkaProducer`, which you must configure, call `send()` on, and close. Spring for Apache Kafka (the `spring-kafka` project) adds a higher-level abstraction so you write less boilerplate and integrate with Spring features. **KafkaTemplate.** This is the central producing component. It is constructed from a `ProducerFactory<K,V>` (which holds the producer config map: `bootstrap.servers`, key/value serializers, etc.). The template borrows a producer from the factory, serializes your key/value, builds a `ProducerRecord`, and hands it to the underlying client. Because `KafkaProducer` is thread-safe, one `KafkaTemplate` bean is safely shared by all threads — you do **not** create one per request. **send() vs sendDefault().** - `send(String topic, K key, V value)` (and overloads with partition, timestamp, a `ProducerRecord`, or a Spring `Message<?>`) explicitly names the destination topic each call. - `sendDefault(K key, V value)` omits the topic and uses the `defaultTopic` you set on the template (e.g. `template.setDefaultTopic("orders")` or via the constructor). It is pure convenience for templates dedicated to a single topic; if no default topic is set, calling it throws. **Asynchrony and results.** `send()` returns a `CompletableFuture<SendResult<K,V>>` (since spring-kafka 3.0; earlier versions returned `ListenableFuture`). The `SendResult` carries the `ProducerRecord` and the broker `RecordMetadata` (partition, offset, timestamp). Producing is buffered and batched by the client, so the future completes only after the broker acknowledges per the `acks` setting. To react, use `future.whenComplete((result, ex) -> ...)`; to force synchronous behavior, call `future.get()` (blocks and surfaces failures), though that sacrifices throughput. **Edge cases.** Serializer exceptions and buffer-full conditions surface through the future or as exceptions on the calling thread. For ordering guarantees with retries, you rely on idempotence (`enable.idempotence=true`) rather than the template itself. Transactions are enabled by giving the `ProducerFactory` a `transactionIdPrefix`, after which `KafkaTemplate` exposes `executeInTransaction` and participates in Spring-managed transactions.

  • Is send() synchronous or asynchronous, and how do you wait for the broker ack?
    Asynchronous. It returns a CompletableFuture<SendResult>; attach whenComplete for a callback, or call get() to block until the broker acknowledges (the future completes per the acks setting).
  • Why is it fine to share one KafkaTemplate across many threads?
    Because the underlying KafkaProducer is thread-safe and designed to be shared; spring-kafka pools/reuses it via the ProducerFactory, so a single template bean serves all callers.

saying these in an interview costs you the question

  • Claiming send() blocks until the record is persisted (it returns a future)
  • Saying you should create a new KafkaTemplate per message/thread
  • Confusing sendDefault with sending to a partition's default broker — it just uses a preset topic

context

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

How does Spring Cloud Stream model Kafka messaging with functional bindings, and how does it relate to the spring-kafka binder?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Spring Cloud Stream lets you write messaging as plain java.util.function beans (Supplier/Function/Consumer). It auto-binds them to Kafka topics via the Kafka binder. You declare functions in spring.cloud.function.definition and map them to topics with spring.cloud.stream.bindings, keeping business code broker-agnostic.

open as a page