How does the Spring Kafka ConcurrentMessageListenerContainer 'concurrency' setting map to partitions and threads?
answer
- Concurrent = parent spawns N single-threaded children
- N consumers, N threads, one group
- ceiling = partition count; surplus children idle
- concurrency is per-instance — replicas multiply
- per-partition order preserved
basics
~10 sconcurrency=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 sSpring 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
Know that concurrency=N means N consumer threads in one group, limited by partitions.
Explain the parent/child container structure and the partition-count ceiling with examples.
Reason about per-instance multiplication across replicas and when the container model can't add parallelism.
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.