skip to content

How does @KafkaListener work, and what does the concurrency setting actually control on a MessageListenerContainer?

level: middleimportance: must knowfreq 68%

answer

  1. ConcurrentMessageListenerContainer = N child single-threaded containers
  2. each thread = its own KafkaConsumer, same group
  3. parallelism = min(concurrency, partitions)
  4. per-partition order preserved
  5. endpoint registry + container factory

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.

solid answer

~40 s

@KafkaListener registers a method as a record handler. The KafkaListenerEndpointRegistry, via a ConcurrentKafkaListenerContainerFactory, builds a ConcurrentMessageListenerContainer that wraps N KafkaMessageListenerContainers — one per thread, each owning its own KafkaConsumer and poll loop. concurrency sets N. Each consumer is a member of the same consumer group, so Kafka's group coordinator assigns partitions across them; effective parallelism is min(concurrency, partition count) — extra threads sit idle. Within a single container, records from one partition are processed sequentially, preserving per-partition ordering. The method can receive the deserialized payload, the full ConsumerRecord, headers via @Header, or (with batch mode) a List. Offsets are committed per the container's AckMode. Container lifecycle (start/stop/pause/resume) is managed by Spring and exposed through the registry.

go deeper

for a junior

Know @KafkaListener consumes from a topic and Spring runs the loop; concurrency adds consumer threads.

for a middle

Explain ConcurrentMessageListenerContainer fan-out, min(concurrency, partitions), and ordering.

for a senior

Tie in rebalances, max.poll.interval.ms, container factory config, and batch vs record dispatch.

for a principal

Design partition counts and instance/concurrency layout across a fleet; reason about idle members and rebalance cost.

**The consumer model.** A Kafka *consumer group* is a set of consumers that cooperatively read a topic: each partition is assigned to exactly one consumer in the group, giving horizontal scaling and ordered, exclusive processing per partition. The raw client requires you to write a `while(true) { poll(); process(); }` loop and manage offset commits yourself. **What @KafkaListener does.** Annotating a bean method with `@KafkaListener(topics="orders", groupId="billing")` tells spring-kafka to create and manage that poll loop for you. At startup, `KafkaListenerAnnotationBeanPostProcessor` finds the annotation and registers an *endpoint* with the `KafkaListenerEndpointRegistry`. The registry uses a `ConcurrentKafkaListenerContainerFactory` to build a **container** that drives the loop and dispatches each record to your method via a `MessageHandlerMethodFactory` (which resolves arguments: payload, `@Header`, `ConsumerRecord`, `Acknowledgment`, `Consumer`). **Container types.** - `KafkaMessageListenerContainer` — single-threaded: one consumer, one poll loop. - `ConcurrentMessageListenerContainer` — a fan-out wrapper that creates `concurrency` child `KafkaMessageListenerContainer`s, each on its own thread with its **own** `KafkaConsumer`. **What concurrency controls.** `concurrency=N` (or `factory.setConcurrency(N)`) starts N consumer threads, all in the same group. Kafka's group coordinator then rebalances partitions across these N consumers (and across other application instances). So: - Effective parallelism = `min(N, numberOfPartitions across subscribed topics)`. If a listener subscribes to a 4-partition topic with concurrency=8, four threads get one partition each and four threads are idle. - Per-partition ordering is preserved because one partition is handled by exactly one thread, processing its records sequentially. - Increasing concurrency does **not** subdivide a single partition; throughput beyond partition count requires more partitions. **Batch vs record.** By default each invocation handles one record. Setting `factory.setBatchListener(true)` makes the method receive a `List<...>` per poll — useful for bulk writes. AckMode and error handling differ between the two modes. **Lifecycle and ops.** Containers are Spring `SmartLifecycle` beans. You can name a listener with `id="..."` and later `registry.getListenerContainer(id).pause()/resume()/stop()`. `autoStartup=false` lets you start consumption manually. Rebalances pause assignment briefly; spring-kafka exposes a `ConsumerRebalanceListener` and `ContainerProperties` for tuning poll timeout, `idleEventInterval`, etc. **Edge cases.** Long processing per record can exceed `max.poll.interval.ms`, causing the broker to evict the consumer and trigger a rebalance; mitigate by increasing that timeout, reducing `max.poll.records`, or offloading work. Each concurrent consumer counts toward group membership, so over-provisioning concurrency across many instances can leave many idle members.

  • If a topic has 3 partitions and you set concurrency=6, what happens?
    Only 3 of the 6 consumer threads get a partition each; the other 3 stay idle. Concurrency above partition count yields no extra parallelism — you must add partitions.
  • Does raising concurrency break per-partition ordering?
    No. Each partition is still owned by exactly one consumer thread and processed sequentially, so ordering within a partition holds. Ordering across partitions was never guaranteed.

saying these in an interview costs you the question

  • Saying concurrency parallelizes processing within a single partition
  • Claiming all concurrency threads share one KafkaConsumer (each has its own)
  • Thinking concurrency=10 on a 2-partition topic gives 10x throughput
  • Forgetting that the threads are members of the same consumer group

context