How does the concurrency setting on ConcurrentMessageListenerContainer relate to a topic's partition count?
answer
- concurrency = N child containers, each 1 KafkaConsumer thread
- partition owned by exactly one consumer per group
- effective parallelism = min(concurrency, partitions)
- C > P => idle consumers
- concurrency multiplies across app instances
basics
~20 sConcurrency 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.
solid answer
~40 sA ConcurrentMessageListenerContainer spawns `concurrency` child KafkaMessageListenerContainers, each with its own KafkaConsumer thread joining the same consumer group. Kafka's rule: within a group a partition is owned by exactly one consumer. So if a topic has 10 partitions and concurrency=4, the 10 partitions are spread across 4 consumers; set concurrency=10 and each consumer gets one partition (max parallelism); set concurrency=12 and 2 consumers get no partitions and sit idle. You can never process one partition with more than one thread from the same group — that's how Kafka preserves per-partition ordering. Practical guidance: keep concurrency <= partition count, size partitions for target parallelism up front (repartitioning is disruptive), and remember concurrency multiplies across app instances — 3 pods each with concurrency=4 means 12 consumers competing for the partitions.
code
java · 14 lines@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
ConsumerFactory<String, String> cf) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
factory.setConsumerFactory(cf);
factory.setConcurrency(3); // 3 KafkaConsumer threads in this container
return factory;
}
// Per-listener override; useful gain only up to the topic's partition count
@KafkaListener(topics = "orders", groupId = "order-service", concurrency = "6")
public void handle(ConsumerRecord<String, String> record) {
// one thread per assigned partition preserves per-partition order
}go deeper
Know concurrency = number of consumer threads and it's limited by partitions.
Explain min(concurrency, partitions), child KafkaMessageListenerContainers, and per-partition ordering.
Reason about fleet-wide consumer count, rebalancing cost, and sizing partitions up front.
Design partition topology for total target parallelism, cooperative rebalancing, and ordering guarantees across a scaled deployment.
**What ConcurrentMessageListenerContainer is** The default container type built by `ConcurrentKafkaListenerContainerFactory`. It is a thin manager that creates **N child `KafkaMessageListenerContainer`s**, where N = the `concurrency` value. Each child owns **one `KafkaConsumer`** running its poll loop on its **own thread**. All children join the **same consumer group** (the listener's `groupId`). **Kafka's partition-assignment rule** Within a single consumer group, **each partition is assigned to exactly one consumer** at a time (enforced by the group coordinator / partition assignor). This is the mechanism that guarantees **per-partition ordering** — only one thread ever reads a given partition, so records are processed in offset order for that partition. **Consequences for concurrency vs partitions** (topic with P partitions, container concurrency C): - **C < P**: partitions are distributed across the C consumers; some consumers own multiple partitions. Works fine. - **C == P**: one partition per consumer — maximum parallelism for this app instance. - **C > P**: only P consumers get partitions; the remaining **C − P consumers are idle** (assigned nothing). Wasted threads, no throughput gain. So **effective parallelism is min(C, P)** for a single instance. **Scaling across instances** Concurrency is per-container, and consumers from **all app instances** in the same group share the partition pool. If you run 3 pods each with concurrency=4, that's **12 consumers** competing for the topic's partitions — so again capped by P. Plan partition count for your *total* desired parallelism across the whole fleet, not per instance. **Ordering caveat** Because records for the same key go to the same partition (default partitioner), per-key ordering holds as long as you don't parallelize within a partition. Features like the container's out-of-order handling or a separate thread pool inside the listener can break ordering — avoid unless you understand the trade-off. **Rebalancing** Adding/removing consumers (scaling, deploys, crashes) triggers a **rebalance**: partitions are revoked and reassigned. During a rebalance processing pauses briefly; frequent rebalances (short `max.poll.interval.ms` exceeded by slow processing) hurt throughput. Cooperative rebalancing (CooperativeStickyAssignor) reduces stop-the-world impact. **Setting concurrency** - Factory-wide: `factory.setConcurrency(4)`. - Per listener: `@KafkaListener(topics="t", concurrency="4")` (overrides the factory). **Gotchas** - Setting concurrency higher than partitions to 'go faster' → idle consumers, no benefit. - Repartitioning a topic later is disruptive (changes key→partition mapping, can reorder in-flight work) — size partitions early. - Concurrency is threads, not connections pooled — each is a full KafkaConsumer with its own network client and offset state.
- You set concurrency=8 on a topic with 4 partitions. What happens?Only 4 consumers receive a partition each; the other 4 sit idle with no assignment. Effective parallelism is 4. No throughput gain — you'd need to add partitions (disruptive) to use 8.
- Why does Kafka cap one partition to one consumer per group?To preserve per-partition ordering: a single reader processes a partition's records strictly in offset order. Allowing two consumers on one partition would make ordering and offset commits ambiguous.
saying these in an interview costs you the question
- Thinking concurrency>partitions increases throughput
- Believing multiple threads can read the same partition in one group
- Forgetting concurrency multiplies across app instances/pods
- Assuming you can freely repartition later with no consequences