skip to content

How does the Spring Kafka ConcurrentMessageListenerContainer 'concurrency' setting map to partitions and threads?

level: middleimportance: must knowfreq 60%

answer

  1. Concurrent = parent spawns N single-threaded children
  2. N consumers, N threads, one group
  3. ceiling = partition count; surplus children idle
  4. concurrency is per-instance — replicas multiply
  5. per-partition order preserved

basics

~10 s

concurrency=N creates N child containers, each with its own consumer and thread sharing the group. Kafka distributes partitions among them, so useful concurrency is capped at the partition count; extra child consumers idle.

solid answer

~40 s

Spring Kafka's ConcurrentMessageListenerContainer is a thin parent that spins up `concurrency` KafkaMessageListenerContainer children, each owning one KafkaConsumer on its own thread, all in the same consumer group. The Kafka group coordinator then distributes partitions across those child consumers via a rebalance. So setting concurrency=4 on a 12-partition topic gives ~3 partitions per thread; setting concurrency higher than the partition count leaves the surplus children with no partitions — they idle. A single listener thread processes records partition-by-partition, preserving per-partition order. Concurrency is a per-instance setting: if you run 3 app instances each with concurrency=4, that's 12 consumers competing for the 12 partitions, so each app gets ~4. You raise the ceiling by adding partitions, not concurrency. @KafkaListener exposes this via the `concurrency` attribute; AckMode and error handling apply per child container.

go deeper

for a junior

Know that concurrency=N means N consumer threads in one group, limited by partitions.

for a middle

Explain the parent/child container structure and the partition-count ceiling with examples.

for a senior

Reason about per-instance multiplication across replicas and when the container model can't add parallelism.

for a principal

Tie concurrency, partition count, and replica autoscaling into a capacity model, and decide when to move to a decoupled worker pool.

## The two container types Spring Kafka has two listener containers: - **KafkaMessageListenerContainer** — single-threaded, one KafkaConsumer, one poll loop. - **ConcurrentMessageListenerContainer** — a parent that creates `concurrency` instances of the single-threaded container. It does NOT add threads to one consumer; it creates *multiple consumers*, each on its own thread, all sharing the `group.id`. ## How concurrency maps to partitions Because all child consumers share the group, the broker-side group coordinator runs a normal rebalance and assigns each topic-partition to exactly one child. Examples on a 12-partition topic: - `concurrency=1` -> 1 consumer owns all 12 partitions. - `concurrency=4` -> 4 consumers, ~3 partitions each. - `concurrency=12` -> 12 consumers, 1 partition each. - `concurrency=20` -> 12 consumers get 1 partition each; **8 idle** with no assignment. So **partition count is the real ceiling.** Over-provisioning concurrency just creates idle consumers (and idle threads). ## Per-partition ordering is preserved A given partition is owned by exactly one child container/thread, and that thread processes its records sequentially in offset order. So ordering is preserved *per partition*, not globally. If you need more processing parallelism than partitions, the container model can't give it — you'd decouple to a worker pool (and take on offset-ordering responsibility). ## Multiple application instances Concurrency is **per JVM/instance**. If you deploy 3 replicas each with `concurrency=4`, you have 12 consumers in the group for 12 partitions -> 1 partition each. Plan total consumers = sum over instances, and keep it <= partition count to avoid idle consumers. This matters in Kubernetes autoscaling: HPA scaling replicas multiplies consumers. ## Configuration ```java @KafkaListener(topics = "orders", concurrency = "4", groupId = "order-svc") public void handle(Order o) { ... } ``` or on the factory: ```java factory.setConcurrency(4); ``` The `AckMode` (e.g. BATCH, RECORD, MANUAL), `CommonErrorHandler`, and back-off all apply within each child container. ## Edge cases - Static assignment (`TopicPartitionOffset`) lets you pin partitions to children, bypassing the rebalance. - Rebalances during scaling momentarily pause processing; cooperative-sticky assignor reduces the impact. - A slow listener can blow `max.poll.interval.ms` and cause the child to be evicted — concurrency does not fix slow per-record processing.

  • You set concurrency=10 but only 4 threads ever do work. Why?
    The topic has 4 partitions. A group assigns at most one consumer per partition, so 4 children get a partition each and the other 6 idle. Add partitions to use more concurrency.
  • If you run 5 pods each with concurrency=3 against a 12-partition topic, how is work distributed?
    15 consumers compete for 12 partitions, so 12 get one partition each and 3 idle. Total consumers across instances must stay <= partition count to avoid idling.

saying these in an interview costs you the question

  • Thinking concurrency adds threads to a single consumer rather than creating multiple consumers.
  • Believing concurrency can scale beyond the partition count.
  • Forgetting concurrency is per-instance, so replicas multiply consumer count.
  • Claiming concurrency gives global ordering — it only preserves per-partition order.

context