Are KafkaProducer and KafkaConsumer thread-safe? How does that shape how you use each across threads?
answer
- producer = thread-safe, share one
- consumer = NOT thread-safe, one thread
- ConcurrentModificationException on misuse
- wakeup() = only safe cross-thread call
- scale via groups not sharing
basics
~10 sKafkaProducer is thread-safe — share one instance across many threads. KafkaConsumer is NOT thread-safe — it must be used by a single thread; the only safe cross-thread call is wakeup().
solid answer
~40 sThe asymmetry is deliberate. KafkaProducer is fully thread-safe and is designed to be shared: many application threads can call send() on one producer concurrently, which is actually faster than one producer per thread because batching and the shared connection pool work better. KafkaConsumer is explicitly NOT thread-safe — it must be confined to one thread. If a second thread calls into it, the consumer detects the multi-threaded access and throws ConcurrentModificationException ('KafkaConsumer is not safe for multi-threaded access'). The single documented exception is wakeup(), which is thread-safe and exists precisely so another thread (e.g. a shutdown hook) can interrupt a blocking poll(). To scale consumption you don't share a consumer across threads; you run multiple consumers (often one per thread) in the same group, or hand records off to a worker pool after polling.
go deeper
Memorize: producer thread-safe (share it), consumer not thread-safe (one thread each).
Explain the ConcurrentModificationException guard and that wakeup() is the only cross-thread method.
Justify the design via batching/Sender for the producer and hot-path locking cost for the consumer; describe the two scaling patterns.
Set team conventions: shared producer beans, consumer confinement, and a standard wakeup-based shutdown protocol; reason about resource cost at fleet scale.
## The core rule - **`KafkaProducer` is thread-safe** — generally you want ONE producer instance shared across all your application threads. - **`KafkaConsumer` is NOT thread-safe** — one consumer instance belongs to exactly ONE thread. ## Why the producer is shared A producer accumulates records into per-partition batches in an internal buffer (`RecordAccumulator`), and a single background thread (the *Sender*) drains batches to brokers. When many threads call `send()` on the same producer, more records land in the same batches in the same window, so batching is *better* and you amortize the connection pool and metadata. Spinning up one producer per thread wastes memory (each has its own buffer, sender thread, and connections) and produces smaller batches. So: share one producer; only make more if you genuinely need different configs (e.g. different `transactional.id`, `acks`, or compression). ## Why the consumer is single-threaded A consumer maintains a lot of mutable, non-synchronized state: the fetch buffer, current position per partition, the in-flight group membership / heartbeat state, and offset bookkeeping. Making all of that thread-safe would add locking overhead to the hot poll path. Instead the client *detects* concurrent use: before most operations it runs `acquire()` which checks that the calling thread is the one that 'owns' the consumer, and throws `ConcurrentModificationException` with the message 'KafkaConsumer is not safe for multi-threaded access' if not. Note the heartbeat is sent from *within* `poll()` (and via background coordinator logic), so you must keep calling `poll()` regularly from that one thread. ## The one cross-thread method: wakeup() `consumer.wakeup()` is the **only** method safe to call from a different thread. It sets a flag so that a blocking `poll()` returns immediately by throwing `WakeupException`. This is how you cleanly stop a poll loop from a shutdown hook or another control thread — you can't just call `close()` from the other thread, because `close()` is itself not thread-safe. ## Scaling consumption (without sharing) Kafka's scaling model is consumer *groups*, not shared instances: 1. **Thread-per-consumer**: N threads, each with its own `KafkaConsumer`, all in the same `group.id`; the group coordinator assigns partitions across them. Simple and common. 2. **Single-consumer + worker pool**: one thread polls, then hands record batches to a thread pool for processing; you must manage offset commits carefully so you don't commit records still in flight. (The detailed partition-to-thread mapping and ordering trade-offs belong to the consumer-concurrency topic; here the point is just the thread-safety contract.) ## Common mistakes - Sharing a single consumer across a thread pool 'to go faster' → `ConcurrentModificationException`. - Creating a new producer per request/message → resource waste and poor batching. - Calling `close()` from a watchdog thread to stop a poll loop → undefined/unsafe; use `wakeup()` instead.
- Why is sharing one producer often faster than one producer per thread?All threads' records funnel into the same per-partition batches, so batches fill better and fewer, larger requests go to the broker; you also share the connection pool, metadata, and a single Sender thread instead of duplicating them.
- If the consumer isn't thread-safe, how do you signal it to stop from another thread?Call consumer.wakeup() from the other thread. It's the only thread-safe method; it makes the in-progress (or next) poll() throw WakeupException so the poll-owning thread can break its loop and close the consumer itself.
- What exactly happens if a second thread calls poll() on the same consumer?The consumer's acquire() guard detects a different thread and throws ConcurrentModificationException with 'KafkaConsumer is not safe for multi-threaded access'. It does not corrupt state silently — it fails fast.
saying these in an interview costs you the question
- Saying KafkaConsumer is thread-safe
- Saying KafkaProducer is single-threaded / should be one-per-thread
- Using close() from another thread to stop the loop
- Claiming you can call commitSync from a worker thread on a shared consumer