Topic 3.1
The Consumer and the Poll Loop
In one line
A consumer subscribes to topics, calls poll() in a loop to fetch batches of records from the partitions it's assigned, deserializes them, processes them, and commits offsets. poll() also drives group membership, so the loop's timing matters as much as its logic.
Think of it like this
A mail sorter who checks the pigeonhole every few seconds (poll), handles whatever is there, and ticks off the last letter handled in a logbook (commit). If they stop checking the pigeonhole for too long, the supervisor assumes they've left and gives their pigeonholes to someone else (rebalance).
Key ideas
- 01
Loop:
subscribe(topics)→while (running) { records = poll(timeout); process(records); commit(); }→close()on shutdown.poll()returns up tomax.poll.records(default 500) records across assigned partitions. - 02
Fetching: consumers fetch from partition leaders (or the nearest replica with rack-aware follower fetching), controlled by
fetch.min.bytes,fetch.max.wait.ms,max.partition.fetch.bytes(1 MB) andfetch.max.bytes(50 MB). Records arrive in offset order per partition. - 03
Where to start:
auto.offset.resetapplies when the group has no committed offset (or it's out of range):latest(default, only new records),earliest(from the beginning of retained data), ornone(throw). - 04
The consumer is not thread-safe: use it from one thread;
wakeup()is the only method safe to call from another thread, used to break out ofpoll()for shutdown. - 05
Graceful shutdown: on SIGTERM, call
wakeup(), catchWakeupExceptionin the loop, commit final offsets synchronously, thenclose(), which leaves the group promptly so partitions are reassigned without waiting for a session timeout.
Code & diagrams
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka1:9092");
p.put(ConsumerConfig.GROUP_ID_CONFIG, "payment-service");
p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
p.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(p);
Runtime.getRuntime().addShutdownHook(new Thread(consumer::wakeup));
try {
consumer.subscribe(List.of("commerce.orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
log.info("topic={} partition={} offset={} key={} ts={}", r.topic(), r.partition(), r.offset(), r.key(), r.timestamp());
handle(r);
}
consumer.commitSync(); // after processing the batch
}
} catch (WakeupException e) {
// shutting down
} finally {
try { consumer.commitSync(); } finally { consumer.close(); }
}When it breaks
Consumer shared across worker threads
What you see
ConcurrentModificationException or subtle corruption of positions; commits for records not yet processed.
Fix & prevent
One consumer per thread, or one polling thread handing work to workers with careful offset tracking (Topic 3.5).
Explain it without notes
Why must the poll loop keep running even while processing is slow?
Practice
Write a consumer that logs topic, partition, offset, key and timestamp for every record, with graceful shutdown on Ctrl+C.
Trade-offs
- ↔
Larger poll batches increase throughput but lengthen processing time per poll, risking poll-interval timeouts.
Done when you can
I can write a correct poll loop with manual commits and graceful shutdown.