What do @KafkaListener and @EnableKafka do, and how do they work together?
answer
- @EnableKafka registers the BeanPostProcessor + EndpointRegistry
- @KafkaListener = method consumer; scanned at startup
- Boot auto-applies @EnableKafka
- Factory builds MessageListenerContainer running poll loop
- 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@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
Know @KafkaListener consumes from a topic and @EnableKafka enables the machinery; in Boot it's automatic.
Explain the BeanPostProcessor + EndpointRegistry + container-factory chain and same-vs-different groupId semantics.
Discuss parameter binding options, custom container factories, and when to declare @EnableKafka explicitly.
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.