4.5 The poll loop
The timeout parameter controls how long poll() will block if data is not available in the consumer buffer.
Duration timeout = Duration.ofMillis(100);
while (true) { // ①
ConsumerRecords<String, String> records = consumer.poll(timeout); // ②
for (ConsumerRecord<String, String> record : records) { // ③
System.out.printf("topic = %s, partition = %d, offset = %d, " +
"customer = %s, country = %s\n",
record.topic(), record.partition(), record.offset(),
record.key(), record.value());
int updatedCount = 1;
if (custCountryMap.containsKey(record.value())) {
updatedCount = custCountryMap.get(record.value()) + 1;
}
custCountryMap.put(record.value(), updatedCount);
JSONObject json = new JSONObject(custCountryMap);
System.out.println(json.toString());
}
}① The infinite loop. "Consumers are usually long-running applications that continuously poll Kafka for more data."
② The book's own words: "This is the most important line in the chapter."
"The same way that sharks must keep moving or they die, consumers must keep polling Kafka or they will be considered dead and the partitions they are consuming will be handed to another consumer in the group."
The timeout parameter controls how long poll() will block if data is not available in the consumer buffer. If set to 0, or if records are already available, poll() returns immediately; otherwise it waits up to that many ms.
③ Each record contains: topic, partition, offset within the partition, key, value.
What poll() secretly does — the reason exceptions surface here
First poll() with a NEW consumer:
- find the
GroupCoordinator - join the consumer group
- receive a partition assignment
Every poll():
- fetch records
- handle any triggered REBALANCE — including related CALLBACKS
- (if autocommit) check whether it's time to commit, and commit
"This means that almost everything that can go wrong with a consumer, or in the callbacks used in its listeners, is likely to show up as an exception thrown by
poll()."
And the constraint that follows:
"if
poll()is not invoked for longer thanmax.poll.interval.ms, the consumer will be considered dead and evicted from the consumer group, so avoid doing anything that can block for unpredictable intervals inside the poll loop."
Thread safety — the rule
ONE CONSUMER PER THREAD.
- You CAN'T have multiple consumers of the same group in one thread.
- You CAN'T have multiple threads safely use the same consumer.
(Note the contrast with Ch. 3: a producer IS thread-safe and shareable. A consumer is not.)
Two patterns for concurrency:
- "wrap the consumer logic in its own object and then use Java's
ExecutorServiceto start multiple threads, each with its own consumer." - "have one consumer populate a queue of events and have multiple worker threads perform work from this queue." (Pattern documented by Igor Buzatović.)
⚠️ WARNING —
poll(long)vspoll(Duration)semantics changedThe old
poll(long)signature is deprecated; the new API ispoll(Duration). The blocking semantics subtly changed:
Behavior poll(long)(old)Blocks as long as it takes to get the needed metadata from Kafka — even if longer than the timeout poll(Duration)(new)Adheres to the timeout restrictions and does NOT wait for metadata The migration trap: "If you have existing consumer code that uses
poll(0)as a method to force Kafka to get the metadata without consuming any records (a rather common hack), you can't just change it topoll(Duration.ofMillis(0))and expect the same behavior. You'll need to figure out a new way to achieve your goals."The recommended replacement: "Often the solution is placing the logic in the
rebalanceListener.onPartitionAssignment()method, which is guaranteed to get called after you have metadata for the assigned partitions but before records start arriving."