Is the Kafka KafkaConsumer thread-safe, and what is the canonical threading model for consuming from Kafka?
answer
- KafkaConsumer NOT thread-safe; producer IS
- ConcurrentModificationException via acquire/release
- wakeup() = only safe cross-thread call
- one-consumer-per-thread, same group.id
- parallelism capped by partition count
basics
~10 sNo. The Java KafkaConsumer is not thread-safe — only one thread may call it (except wakeup()). The standard model is one consumer instance per thread, each in its own poll loop.
solid answer
~50 sThe Java KafkaConsumer is explicitly NOT thread-safe: a single instance must be used by exactly one thread, and concurrent access throws ConcurrentModificationException (it checks via a lightweight acquire/release guard). The one exception is wakeup(), which is safe to call from another thread to break a blocking poll(). The canonical model is therefore 'one consumer per thread': you run N consumer instances, each with its own poll loop, all sharing the same group.id so the group coordinator distributes partitions among them. Scaling is bounded by partition count — a group can have at most one active consumer per partition, so extra consumers idle. The producer, by contrast, IS thread-safe and meant to be shared. If you want more processing threads than consumers, you decouple: poll on one thread, hand records to a worker pool — but then you own offset-commit ordering.
go deeper
Know the headline: consumer not thread-safe, one per thread, parallelism limited by partitions.
Explain the acquire/release guard, wakeup() for shutdown, and the partition ceiling.
Contrast with the thread-safe producer, discuss when to decouple polling from processing and the offset-ordering cost.
Frame trade-offs of partition count as the scaling lever vs decoupled worker pools, and the operational impact of rebalances and max.poll.interval.ms.
## What 'thread-safe' means here Kafka's Java client `org.apache.kafka.clients.consumer.KafkaConsumer` keeps mutable state (fetch buffers, the in-flight position, coordinator state) that is not protected by locks for general use. To enforce single-threaded use cheaply, every public method calls an internal `acquire()`/`release()` that records the current thread id; if a second thread enters concurrently it throws `java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access`. So: **one thread per consumer instance**. ## The one safe cross-thread call: wakeup() `poll(Duration)` blocks waiting for data. To shut a consumer down or interrupt it, another thread may call `consumer.wakeup()`. The blocked `poll()` then throws `WakeupException`, which you catch to exit the loop cleanly. `wakeup()` is the *only* method documented as safe to call from a different thread. (`close()` from another thread is not safe.) ## Canonical model: one-consumer-per-thread A **consumer group** is a set of consumers sharing the same `group.id`. The broker-side **group coordinator** runs a rebalance that assigns each topic-partition to exactly one consumer in the group. So if a topic has 12 partitions and you start 4 consumer threads (same group), each gets ~3 partitions. Each thread runs an independent loop: `poll() -> process -> commit`. This gives you parallelism with **no shared mutable consumer state** and naturally preserves per-partition order, because all records of a partition flow through one thread. ## The partition ceiling Parallelism in a single group is capped by partition count: a group can have at most one *active* consumer per partition. Start 20 consumers on a 12-partition topic and 8 sit idle. To raise the ceiling you must add partitions (which is a one-way operation that can break keyed ordering across the change) or decouple processing from polling. ## Producer is different The `KafkaProducer` IS thread-safe and is designed to be shared by many threads — sharing one producer is generally faster than many. Don't confuse the two. ## Edge cases - Calling consumer methods from a Kafka rebalance-listener callback is fine (same thread). - Long processing between polls can exceed `max.poll.interval.ms` and trigger a rebalance that evicts you; the fix is smaller `max.poll.records` or decoupled processing with `pause()`/`resume()`. - Frameworks like Spring Kafka's `ConcurrentMessageListenerContainer` and Kafka Streams hide the loop but still obey one-consumer-per-thread underneath.
- Which KafkaConsumer method is safe to call from another thread, and why would you use it?wakeup(). It interrupts a blocking poll() by causing it to throw WakeupException, used to cleanly break the poll loop for shutdown.
- If you start 10 consumers in one group on a 6-partition topic, what happens?Only 6 consumers get partitions (one each); the other 4 are assigned nothing and stay idle until a rebalance reassigns work, e.g. if an active consumer dies.
saying these in an interview costs you the question
- Saying KafkaConsumer is thread-safe and can be shared across threads.
- Claiming you can add consumers indefinitely to scale beyond partition count.
- Confusing the producer (thread-safe) with the consumer (not).
- Thinking close() is safe to call from another thread (only wakeup() is).