skip to content

Spring for Apache Kafka

Spring for Apache Kafka: the producer template, listener containers, error handling and dead-letter topics, transactions and exactly-once, Kafka Streams, and topic provisioning. Kafka is the default event backbone in interviews, so this section carries weight.

part ofSpring Frameworkoverview, primer and where to startread it →
on this pageshow

explore

questions

page 1 of 2

How do you make Spring create a Kafka topic automatically when the application starts?

level: juniorimportance: must knowfreq 45%

answer

  1. NewTopic @Bean
  2. KafkaAdmin scans context at startup
  3. TopicBuilder.name().partitions().replicas()
  4. wraps AdminClient
  5. not broker auto.create.topics.enable

basics

~10 s

Declare a NewTopic as a @Bean (usually built with TopicBuilder). Spring's KafkaAdmin bean scans the context for NewTopic beans at startup and creates any that don't yet exist on the broker.

solid answer

~40 s

Spring for Apache Kafka provides a KafkaAdmin bean (auto-configured by Spring Boot from spring.kafka.bootstrap-servers). At context startup KafkaAdmin looks through the ApplicationContext for every NewTopic bean and, using an internal Kafka AdminClient, creates the ones that are missing on the broker. You just expose a NewTopic @Bean, most idiomatically via TopicBuilder: TopicBuilder.name("orders").partitions(6).replicas(3).build(). Nothing else is required — no explicit AdminClient calls, no imperative creation code. If the topic already exists it is left in place (Spring may add partitions but never fails just because it's there). This is declarative, in-app provisioning; it is not the broker-side auto.create.topics.enable feature, which is a completely separate cluster setting.

code

java · 13 lines
java
@Configuration
class TopicConfig {

    // KafkaAdmin (auto-configured by Spring Boot) finds this bean at
    // startup and creates the topic if it does not already exist.
    @Bean
    NewTopic ordersTopic() {
        return TopicBuilder.name("orders")
                .partitions(6)
                .replicas(3)
                .build();
    }
}

go deeper

for a junior

Know that a NewTopic @Bean plus KafkaAdmin means Spring creates the topic at startup.

for a middle

Explain that KafkaAdmin wraps AdminClient and scans the context; know TopicBuilder and the 1/1 default.

for a senior

Contrast with broker-side auto-create and reason about idempotency and existing-topic behavior.

for a principal

Weigh in-app declarative provisioning against infra-as-code and environment-specific replication.

## The pieces **KafkaAdmin** (`org.springframework.kafka.core.KafkaAdmin`) is a Spring-managed bean that wraps the Kafka client's **AdminClient** (`org.apache.kafka.clients.admin.AdminClient`). AdminClient is the low-level Kafka API for administrative operations (create/delete topics, alter configs, list topics). KafkaAdmin's job is to make topic creation *declarative* so you don't call AdminClient by hand. With **Spring Boot**, `KafkaAutoConfiguration` creates a `KafkaAdmin` for you from `spring.kafka.bootstrap-servers` and any `spring.kafka.admin.*` properties. Without Boot you declare it yourself: ```java @Bean KafkaAdmin kafkaAdmin() { return new KafkaAdmin(Map.of(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")); } ``` **NewTopic** (`org.apache.kafka.clients.admin.NewTopic`) is a plain description of a topic: its name, partition count, replication factor, and topic-level configs. You never create the topic yourself — you just *describe* it as a Spring bean. **TopicBuilder** (`org.springframework.kafka.config.TopicBuilder`) is a fluent helper that produces a `NewTopic`. ## How auto-creation fires During context initialization, KafkaAdmin (via its `initialize()` method) collects **all `NewTopic` (and `KafkaAdmin.NewTopics`) beans** in the context and calls AdminClient to create the missing ones. This happens once, at startup. Topics that already exist are not re-created. ## Minimal example ```java @Configuration class TopicConfig { @Bean NewTopic ordersTopic() { return TopicBuilder.name("orders") .partitions(6) .replicas(3) .build(); } } ``` ## Key gotchas - **Only `NewTopic` *beans* are provisioned.** A `NewTopic` you `new` inside a method but never expose as a bean does nothing. - If no partitions/replicas are set, TopicBuilder defaults to **1 partition, 1 replica** — a classic footgun in multi-broker/prod clusters. - This is *not* `auto.create.topics.enable`. That is a **broker** setting that lazily creates a topic the first time any client references an unknown one, with broker-default partitions/replication. KafkaAdmin auto-creation is explicit, in-application, and gives you control over partitions/replicas/configs. - If the broker is unreachable at startup, by default the app still starts (creation just fails/logs) — it is not fatal unless you opt in.

  • What is the difference between this and the broker's auto.create.topics.enable?
    KafkaAdmin auto-creation is explicit and in-app: you declare NewTopic beans with chosen partitions/replicas/configs, created at startup. auto.create.topics.enable is a broker cluster setting that lazily creates a topic on first client reference using broker defaults, with no per-topic control — and it's out of the application's hands.
  • If you declare a NewTopic but never mark it @Bean, what happens?
    Nothing. KafkaAdmin only scans the ApplicationContext for NewTopic beans, so a NewTopic that isn't a Spring bean is invisible to provisioning.

saying these in an interview costs you the question

  • Thinking the topic is created lazily on first send rather than at context startup
  • Confusing KafkaAdmin auto-creation with broker auto.create.topics.enable
  • Believing you must call AdminClient.createTopics yourself

context

open as a page

What is Spring Kafka's DefaultErrorHandler and how does it retry a failed record?

level: juniorimportance: must knowfreq 70%

basics

~10 s

DefaultErrorHandler is the container's default handler when a @KafkaListener throws. It re-delivers the failed record several times (using a BackOff for delays), and if all attempts fail it hands the record to a recoverer.

open as a page

What do @KafkaListener and @EnableKafka do, and how do they work together?

level: juniorimportance: must knowfreq 78%

basics

~10 s

@KafkaListener marks a method to consume messages from a Kafka topic. @EnableKafka turns on Spring's infrastructure that finds those methods and starts listener containers to feed them records.

open as a page

How do you enable Kafka Streams in a Spring (Boot) application, and what does @EnableKafkaStreams actually wire up?

level: juniorimportance: must knowfreq 62%

basics

~10 s

Put @EnableKafkaStreams on a @Configuration class and provide a KafkaStreamsConfiguration bean named defaultKafkaStreamsConfig that sets the application.id and bootstrap.servers. Spring then creates a StreamsBuilderFactoryBean that builds and starts Kafka Streams.

open as a page

What is KafkaTemplate and how do you use it to send a message to a Kafka topic in Spring?

level: juniorimportance: must knowfreq 78%

basics

~10 s

KafkaTemplate is Spring's helper for producing messages. You inject it and call template.send("topic", key, value). Spring Boot auto-configures it from application.properties, so you just autowire and send.

open as a page

What does setting a Kafka consumer's isolation.level to read_committed do, and why does it matter when producers use transactions?

level: juniorimportance: must knowfreq 45%

basics

~10 s

read_committed makes the consumer skip records from aborted transactions and not read records from transactions that are still open. It only sees data the producer actually committed, so a rolled-back send never becomes visible.

open as a page

What does DeadLetterPublishingRecoverer do, and what does the target topic and message look like?

level: middleimportance: must knowfreq 65%

basics

~20 s

It's a recoverer that, after retries are exhausted, republishes the failed record to a dead-letter topic (default: original topic name + ".DLT") using a KafkaTemplate, adding headers describing the failure so the message isn't lost.

open as a page

How does the concurrency setting on ConcurrentMessageListenerContainer relate to a topic's partition count?

level: middleimportance: must knowfreq 80%

basics

~20 s

Concurrency is how many consumer threads (each its own KafkaConsumer) the container runs. Kafka gives each partition to at most one consumer in a group, so useful concurrency is capped by partition count — extra threads stay idle.

open as a page

How do you define a Kafka Streams topology (a KStream/KTable) as a Spring bean, and who starts it?

level: middleimportance: must knowfreq 55%

basics

~20 s

Write a @Bean method that takes a StreamsBuilder parameter. Spring injects the default StreamsBuilder, you build your KStream/KTable pipeline on it, and the StreamsBuilderFactoryBean builds the topology and auto-starts the KafkaStreams instance after the context refreshes.

open as a page

KafkaTemplate.send() returns a CompletableFuture<SendResult> — what does that mean for correctness, and how do you handle success and failure?

level: middleimportance: must knowfreq 72%

basics

~20 s

send() is asynchronous — it returns a future that completes when Kafka acknowledges the write. You attach whenComplete/thenAccept callbacks to handle the SendResult on success or the exception on failure, instead of assuming send succeeded.

open as a page

How do you enable Kafka transactions in Spring, and what roles do transactional.id (transactionIdPrefix) and KafkaTransactionManager play?

level: middleimportance: must knowfreq 40%

basics

~20 s

Give the producer a transactional.id (in Spring, a transactionIdPrefix on the producer factory) so it can run transactions and be fenced across restarts. Then use KafkaTransactionManager so KafkaTemplate sends inside @Transactional commit or roll back atomically.

open as a page

What is @RetryableTopic non-blocking retry, and how does it differ from DefaultErrorHandler's blocking seek retries?

level: seniorimportance: must knowfreq 60%

basics

~20 s

@RetryableTopic makes retries non-blocking: a failed record is forwarded to separate retry topics (with delays), so the main consumer keeps processing new records instead of being stuck. After the configured attempts, it goes to a DLT. DefaultErrorHandler instead retries in place by seeking, which blocks the partition.

open as a page

What are the Spring Kafka AckModes (RECORD, BATCH, MANUAL, etc.) and how do you use manual acknowledgment?

level: seniorimportance: must knowfreq 74%

basics

~10 s

AckMode controls when the container commits offsets. RECORD commits after each record; BATCH commits after the whole poll batch (default); MANUAL/MANUAL_IMMEDIATE hand you an Acknowledgment object so your code decides when to commit.

open as a page

In a read-process-write pipeline, how does Spring Kafka make the output sends and the consumer offset commit atomic (exactly-once)?

level: seniorimportance: must knowfreq 40%

basics

~20 s

Configure the listener container with a KafkaTransactionManager. For each batch it begins a transaction, runs your listener's sends, then commits the consumed offsets into that same transaction via the producer, and commits once. If anything fails it aborts, so neither the output nor the offset advance.

open as a page

How do you configure partitions, replicas, and topic-level settings (like a compacted topic) using TopicBuilder?

level: middleimportance: should knowfreq 35%

basics

~10 s

Use TopicBuilder.name("t").partitions(n).replicas(m), then add topic configs via .config(key, value) or .compact()/.configs(map). .build() returns a NewTopic bean that KafkaAdmin provisions.

open as a page

How do you send Java objects (POJOs) as JSON with KafkaTemplate, and what should you know about Spring's JsonSerializer?

level: middleimportance: should knowfreq 58%

basics

~10 s

Configure the value-serializer as Spring's JsonSerializer. It uses Jackson to turn your POJO into JSON bytes. Then send(topic, key, pojo) works directly. The consumer side uses JsonDeserializer to reconstruct the object.

open as a page

What happens to KafkaAdmin topic provisioning if the Kafka broker is unavailable when the Spring context starts, and how do you control that behavior?

level: seniorimportance: should knowfreq 22%

basics

~10 s

By default it's non-fatal: KafkaAdmin logs that it couldn't reach the broker and the app still starts, but the declared topics are not created. Set fatalIfBrokerNotAvailable (spring.kafka.admin.fail-fast) to make startup fail instead.

open as a page

At startup, what does KafkaAdmin do when a declared NewTopic already exists on the broker but with different partitions, replication factor, or configs?

level: seniorimportance: should knowfreq 28%

basics

~20 s

It won't recreate the topic. It can increase partitions to match a higher declared count (never decrease). It ignores replication-factor differences. Topic configs are only altered if modifyTopicConfigs is enabled; otherwise differences are left alone.

open as a page

How does DefaultErrorHandler distinguish recoverable from fatal (non-retryable) exceptions, and how do you customize the classification?

level: seniorimportance: should knowfreq 50%

basics

~20 s

It uses an exception classifier. Some exceptions (like deserialization or conversion errors) are fatal by default — retrying can't help — so they skip retries and go straight to the recoverer. You adjust the lists with addRetryableExceptions / addNotRetryableExceptions.

open as a page

What is a batch @KafkaListener, how do you enable it, and how does acknowledgment differ from record listeners?

level: seniorimportance: should knowfreq 58%

basics

~20 s

A batch listener receives a whole List of records from one poll() in a single method call instead of one record at a time. You enable it on the container factory (setBatchListener(true)) and the method takes a List; with manual ack, one Acknowledgment covers the entire batch.

open as a page

In a Spring-managed Kafka Streams app, how do you handle uncaught stream exceptions and customize the client before it starts?

level: seniorimportance: should knowfreq 33%

basics

~20 s

Set a StreamsUncaughtExceptionHandler on the StreamsBuilderFactoryBean to decide whether to replace the thread, shut down the client, or shut down the whole app. For anything you can't set via properties, register a KafkaStreamsCustomizer that receives the KafkaStreams before start().

open as a page

What is StreamsBuilderFactoryBean and how does it relate to StreamsConfig and the KafkaStreams lifecycle?

level: seniorimportance: should knowfreq 40%

basics

~20 s

StreamsBuilderFactoryBean is a Spring FactoryBean that produces the StreamsBuilder and also manages the KafkaStreams lifecycle. It is configured from a KafkaStreamsConfiguration (the StreamsConfig properties) and, as a lifecycle bean, builds the topology and starts/stops Kafka Streams with the context.

open as a page

How does the message key passed to KafkaTemplate.send() affect partitioning, and when would you write a custom Partitioner?

level: seniorimportance: should knowfreq 63%

basics

~20 s

The key decides the partition: Kafka hashes the key and maps it to a partition, so all records with the same key go to the same partition and stay ordered. With a null key, records are spread across partitions. A custom Partitioner overrides this mapping.

open as a page

What is the role of ProducerFactory / DefaultKafkaProducerFactory, and how does it manage the underlying Kafka producer(s)?

level: seniorimportance: should knowfreq 55%

basics

~10 s

ProducerFactory creates the underlying Kafka Producer objects for KafkaTemplate, holding the config and serializers. DefaultKafkaProducerFactory is the standard implementation; by default it creates one shared, thread-safe producer that all sends reuse.

open as a page

When would you use KafkaTemplate.executeInTransaction versus @Transactional with a KafkaTransactionManager?

level: seniorimportance: should knowfreq 30%

basics

~10 s

Use executeInTransaction for a quick local, producer-only transaction without any transaction manager or Spring context. Use @Transactional with KafkaTransactionManager when you want declarative transactions that can coordinate with the listener container or other resources.

open as a page

How would you choose and design a Kafka error/retry/DLT strategy for a service with mixed ordering and throughput requirements?

level: principalimportance: should knowfreq 40%

basics

~20 s

Match the mechanism to the constraint: use blocking seek retries (DefaultErrorHandler) where per-partition ordering matters and outages are short; use non-blocking @RetryableTopic where throughput and long back-offs matter and ordering can be relaxed. Always classify transient vs fatal, and route the unrecoverable to a DLT with replay tooling.

open as a page

What is ContainerProperties and which settings would you tune for reliability and rebalance stability?

level: principalimportance: should knowfreq 45%

basics

~20 s

ContainerProperties is the configuration object on a listener container holding container-level settings: ack mode, poll timeout, ack time/count, consumer rebalance listener, sync/async commits, and the task executor. You tune it for commit behavior and rebalance stability.

open as a page

As a principal engineer, how do you configure KafkaTemplate/ProducerFactory for durable, exactly-once-style delivery, and what are the trade-offs?

level: principalimportance: should knowfreq 40%

basics

~20 s

Set acks=all and enable.idempotence=true on the ProducerFactory for durable, non-duplicating writes. For atomic multi-message sends, add a transactionIdPrefix so the template runs Kafka transactions. The trade-off is lower throughput and higher latency for stronger guarantees.

open as a page

What are the real guarantees and limitations of Kafka exactly-once semantics, and how do fencing, EOSMode, and non-Kafka side effects factor in?

level: principalimportance: should knowfreq 25%

basics

~20 s

Kafka EOS guarantees exactly-once only for read-process-write within Kafka: idempotent producers plus transactions make produce+offset-commit atomic. It does not cover external systems like databases or HTTP calls, so those must be idempotent. Fencing and read_committed complete the picture.

open as a page

Should you rely on Spring's KafkaAdmin auto-creation for provisioning topics in production? Discuss the trade-offs and how you'd structure it.

level: principalimportance: nice to knowfreq 25%

basics

~20 s

It's great for local/dev and convenience, but as the sole production strategy it's weak: it won't change replication factor, tolerates config drift by default, can silently create 1/1 topics, and ties correct topic layout to app startup. Prefer infra-as-code for prod and use KafkaAdmin declaratively.

open as a page

showing 1–30 of 31