Learn Labs
4. Kafka Consumers: Reading Data from Kafka

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 than max.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:

  1. "wrap the consumer logic in its own object and then use Java's ExecutorService to start multiple threads, each with its own consumer."
  2. "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) vs poll(Duration) semantics changed

The old poll(long) signature is deprecated; the new API is poll(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 to poll(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."


On this page