skip to content

Java Client Configuration and Lifecycle

Constructing, sharing and shutting down producers and consumers correctly, including thread-safety and wakeup-based shutdown. Interviewers ask because whether KafkaConsumer is thread-safe catches a lot of candidates.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

6

How do you construct a KafkaProducer and a KafkaConsumer in the Java client, and what are the minimum required configuration properties for each?

level: juniorimportance: must knowfreq 75%

answer

  1. Properties -> constructor
  2. producer: 3 serializers config keys
  3. consumer: deserializers + group.id
  4. bootstrap.servers = discovery seed
  5. ProducerConfig/ConsumerConfig constants

basics

~20 s

Build a Properties (or Map) of config, then pass it to the constructor: new KafkaProducer<>(props) or new KafkaConsumer<>(props). Producers need bootstrap.servers plus a key and value serializer; consumers need bootstrap.servers plus a key and value deserializer (and usually group.id).

solid answer

~40 s

Both clients are created by passing a Properties or Map<String,Object> to their constructor. For a producer you must set bootstrap.servers (the broker host:port list used to discover the cluster), key.serializer and value.serializer (classes turning your keys/values into byte[]). For a consumer you set bootstrap.servers, key.deserializer and value.deserializer (byte[] back into objects), and almost always group.id so the consumer joins a group. You can pass serializers/deserializers either as class names in config or as instances to the constructor's second/third argument. Config keys live in ProducerConfig and ConsumerConfig as constants (e.g. ProducerConfig.BOOTSTRAP_SERVERS_CONFIG). Unknown keys log a warning but don't fail. The generic types <K,V> must match the (de)serializers you configure.

code

java · 7 lines
java
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "broker1:9092,broker2:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
    producer.send(new ProducerRecord<>("topic", "key", "value"));
}

go deeper

for a junior

Know the three required producer keys and three required consumer keys, and that you pass a Properties to the constructor.

for a middle

Know serializers can be passed as instances, that group.id is needed for subscribe(), and that bad keys are silently ignored.

for a senior

Discuss using typed config constants, instance-based serializers for DI/schema-registry, and that construction eagerly opens connections.

for a principal

Frame config as a contract: type-safety of generics vs runtime (de)serializers, fail-fast strategies, and standardizing client factory wrappers across teams.

## What a Kafka client is Apache Kafka is a distributed log: producers append records to topics (split into partitions), consumers read them. The Java clients `KafkaProducer<K,V>` and `KafkaConsumer<K,V>` are the official library for doing this. `<K,V>` are the types of the record key and value (e.g. `<String,String>`). ## Construction Both are constructed by handing them configuration. Two equivalent forms: ```java Properties props = new Properties(); props.put("bootstrap.servers", "broker1:9092,broker2:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); KafkaProducer<String,String> producer = new KafkaProducer<>(props); ``` You can also pass a `Map<String,Object>`. The string keys have typed constants in `ProducerConfig`/`ConsumerConfig` (e.g. `ProducerConfig.BOOTSTRAP_SERVERS_CONFIG` == `"bootstrap.servers"`) — prefer the constants to avoid typos, because **misspelled keys are silently ignored** (logged as 'unknown config' but not an error). ## Minimum config **Producer**: `bootstrap.servers`, `key.serializer`, `value.serializer`. A serializer is a class implementing `Serializer<T>` that turns an object into `byte[]` (Kafka only moves bytes). **Consumer**: `bootstrap.servers`, `key.deserializer`, `value.deserializer`. A deserializer (`Deserializer<T>`) turns `byte[]` back into an object. You almost always also set `group.id` — the consumer group name used for coordinated consumption and offset storage; without it you cannot use group subscription or auto-committed offsets. ## bootstrap.servers A comma-separated list of `host:port` for some brokers. The client only uses these to *bootstrap*: it connects, fetches cluster metadata (all brokers, partition leaders), then talks to the right brokers directly. You don't need to list every broker — two or three for redundancy is enough. ## Serializers as instances vs class names You can pass instances instead of class names: ```java new KafkaProducer<>(props, new StringSerializer(), new StringSerializer()); ``` Useful when a serializer needs constructor arguments (e.g. a schema-registry client) that config-by-class-name can't provide. ## Edge cases - Type mismatch between `<K,V>` generics and the configured (de)serializer surfaces as a `ClassCastException` at send/poll time, not at construction. - The constructor eagerly starts background threads/network connections, so construction can be relatively expensive — build once and reuse.

  • What happens if you misspell a config key like 'bootstrap.server'?
    It's silently ignored — Kafka logs a warning about an unknown/unused config but does not throw, so the real key (bootstrap.servers) ends up unset and the client fails to find the cluster. Using the ProducerConfig/ConsumerConfig constants avoids this.
  • Do you need to list every broker in bootstrap.servers?
    No. It's only a discovery seed: the client connects to one, fetches full cluster metadata, then talks to the correct leaders directly. List a few for redundancy in case one seed is down at startup.

saying these in an interview costs you the question

  • Claiming bootstrap.servers must list all brokers
  • Saying a typo'd config throws an exception
  • Forgetting group.id for a consumer that uses subscribe()
  • Thinking Kafka serializes objects automatically without configured serializers

context

open as a page

Describe the consumer poll loop. Why must you call poll() regularly, and what is the role of max.poll.interval.ms?

level: middleimportance: must knowfreq 78%

basics

~20 s

After subscribing, you loop calling poll(timeout) to fetch records, process them, then poll again. poll() also drives group membership/heartbeats. If you take too long between polls (longer than max.poll.interval.ms) the broker assumes you're stuck, removes you from the group, and rebalances your partitions to others.

open as a page

Are KafkaProducer and KafkaConsumer thread-safe? How does that shape how you use each across threads?

level: middleimportance: must knowfreq 80%

basics

~10 s

KafkaProducer 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().

open as a page

Explain the producer's send() async model and the roles of flush() and close(). What can go wrong if you skip them?

level: seniorimportance: must knowfreq 65%

basics

~20 s

send() is asynchronous: it buffers the record and returns a Future immediately; a background Sender thread actually transmits batches. flush() blocks until all buffered records have been sent and acknowledged. close() flushes then releases resources. Skip them and you can lose unsent buffered records on exit.

open as a page

What are serializers and deserializers in the Kafka Java client, and what happens when a deserializer hits a bad (poison) record?

level: middleimportance: should knowfreq 55%

basics

~20 s

Kafka moves only bytes. A Serializer<T> turns your key/value object into byte[] before producing; a Deserializer<T> turns byte[] back into an object when consuming. A bad record that the deserializer can't parse throws inside poll(), and by default the consumer keeps failing on that same offset — a 'poison pill' that blocks progress until you handle it.

open as a page

How do you cleanly shut down a consumer that is blocked in poll()? Explain wakeup() vs close() and the typical shutdown-hook pattern.

level: seniorimportance: should knowfreq 60%

basics

~20 s

From another thread (e.g. a JVM shutdown hook) call consumer.wakeup() — the only thread-safe consumer method. It makes the in-progress poll() throw WakeupException. The poll thread catches it, breaks the loop, and calls close() itself (close() must run on the consumer's own thread).

open as a page