skip to content

@KafkaListener & Containers

@KafkaListener runs in a container whose concurrency should be matched to the partition count, with acknowledge modes and optional batch listening. The concurrency-versus-partitions question is asked in almost every Kafka interview.

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

explore

questions

5

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

level: juniorimportance: must knowfreq 78%

answer

  1. @EnableKafka registers the BeanPostProcessor + EndpointRegistry
  2. @KafkaListener = method consumer; scanned at startup
  3. Boot auto-applies @EnableKafka
  4. Factory builds MessageListenerContainer running poll loop
  5. same groupId = split partitions; different = fan-out

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.

solid answer

~40 s

@EnableKafka (on a @Configuration class, or auto-applied by Spring Boot) registers a KafkaListenerAnnotationBeanPostProcessor plus a KafkaListenerEndpointRegistry. At startup that post-processor scans beans for @KafkaListener methods and, for each, builds a MessageListenerContainer from a ConcurrentKafkaListenerContainerFactory. The container runs a poll loop on a background thread, deserializes each ConsumerRecord, and invokes your annotated method — passing the payload (converted to your parameter type), or the whole ConsumerRecord, headers, key, etc. @KafkaListener attributes let you set topics, groupId, concurrency, and containerFactory. In Boot you rarely write @EnableKafka yourself: KafkaAnnotationDrivenConfiguration adds it and auto-configures the factory from spring.kafka.* properties. Multiple methods can share a group or use distinct groups for independent consumption.

code

java · 22 lines
java
@Configuration
@EnableKafka // auto-applied by Spring Boot; explicit here for a plain Spring app
public class KafkaConfig {

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory) {
        var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
        factory.setConsumerFactory(consumerFactory);
        return factory;
    }
}

@Component
public class OrderConsumer {

    @KafkaListener(topics = "orders", groupId = "order-service")
    public void handle(String payload,
                       @Header(KafkaHeaders.RECEIVED_KEY) String key) {
        // process the record
    }
}

go deeper

for a junior

Know @KafkaListener consumes from a topic and @EnableKafka enables the machinery; in Boot it's automatic.

for a middle

Explain the BeanPostProcessor + EndpointRegistry + container-factory chain and same-vs-different groupId semantics.

for a senior

Discuss parameter binding options, custom container factories, and when to declare @EnableKafka explicitly.

for a principal

Frame group/topology decisions (fan-out vs competing consumers), factory customization, and lifecycle via the EndpointRegistry.

**The two annotations** - **@EnableKafka**: a class-level annotation that bootstraps Spring's annotation-driven Kafka listener infrastructure. It imports configuration that registers two key beans: a **KafkaListenerAnnotationBeanPostProcessor** (the thing that discovers `@KafkaListener` methods) and a **KafkaListenerEndpointRegistry** (the registry that holds and lifecycle-manages the resulting containers). Without this infrastructure, `@KafkaListener` annotations are inert. In **Spring Boot**, `KafkaAnnotationDrivenConfiguration` applies `@EnableKafka` for you when `spring-kafka` is on the classpath, so you seldom write it explicitly — you add it manually only in a non-Boot Spring app. - **@KafkaListener**: a method (or class) level annotation declaring that the method consumes records. Common attributes: `topics` (or `topicPattern`/`topicPartitions`), `groupId`, `id` (bean name of the container in the registry), `concurrency`, `containerFactory`, `properties`, and `autoStartup`. **What happens at startup** 1. `@EnableKafka` registers the bean post-processor. 2. The BPP scans every bean for `@KafkaListener` methods. 3. For each, it creates a **KafkaListenerEndpoint** and asks a **KafkaListenerContainerFactory** (by default the bean named `kafkaListenerContainerFactory`, normally a **ConcurrentKafkaListenerContainerFactory**) to build a **MessageListenerContainer**. 4. The container is registered in the `KafkaListenerEndpointRegistry` and, if `autoStartup=true` (default), started with the application context. **The runtime loop** Each container owns one or more `KafkaConsumer` instances. On a dedicated thread it calls `consumer.poll(...)`, then for every `ConsumerRecord` invokes a **MessageListenerAdapter** that binds record data to your method parameters. A **MessageConverter** / **Deserializer** turns raw bytes into your payload type. You can declare parameters like the payload directly, `@Payload`, `@Header(KafkaHeaders.RECEIVED_KEY)`, `@Headers Map`, the raw `ConsumerRecord<K,V>`, `Acknowledgment` (for manual ack), or `Consumer<K,V>`. **Grouping and scaling** - `groupId` sets the consumer group; Kafka distributes a topic's partitions across all consumers sharing a group. - Two listeners with the **same** groupId on the **same** topic split partitions (competing consumers); with **different** groupIds each gets every record (fan-out). **Minimal Boot config** — set `spring.kafka.bootstrap-servers` and `spring.kafka.consumer.group-id`; Boot wires the factory. Add a custom `ConcurrentKafkaListenerContainerFactory` bean only to override deserializers, error handlers, ack mode, concurrency, etc. **Gotchas** - Forgetting `@EnableKafka` in a **non-Boot** app → listeners silently never start. - The default container factory bean name is `kafkaListenerContainerFactory`; if you define a differently named factory you must reference it via `containerFactory` on the annotation. - `@KafkaListener` methods should be on **Spring-managed beans**; the BPP only scans beans in the context.

  • Do you need to write @EnableKafka in a Spring Boot app?
    No. When spring-kafka is on the classpath, Boot's KafkaAnnotationDrivenConfiguration applies it automatically and auto-configures a ConcurrentKafkaListenerContainerFactory from spring.kafka.* properties. You add it manually only in a non-Boot Spring app.
  • If two @KafkaListener methods use the same groupId and topic, does each get every message?
    No. Sharing a groupId makes them competing consumers in one group — Kafka splits the topic's partitions between them, so each record goes to only one. Use different groupIds for fan-out where both receive every record.

context

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

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

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

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