PrepZone Logo
PrepZone

Poll Loop, max.poll.interval, and Heartbeats

Why long-running VaultCommerce handlers must not block poll() — heartbeats, session timeout, and max.poll.interval.ms.

Why this matters

  • VaultCommerce's image-resize consumer exceeded max.poll.interval — kicked from group every 10 minutes.
  • Heartbeats prove liveness; poll frequency must stay within session.timeout.ms / 3.
  • Long processing belongs off the poll thread — pause/resume or decouple with queue.
  • Classic on-call issue: 'consumer keeps rebalancing'.
Topic partitions
P0
P1
P2
Group: inventory-service
Consumer Areads P0
Consumer Breads P1
Consumer Creads P2
Each partition is consumed by exactly one consumer in the group. Add consumers up to partition count.

poll() contract

consumer.poll(Duration) returns records AND sends heartbeats. Default max.poll.records=500 — processing 500 heavy records before next poll risks timeout. VaultCommerce sets max.poll.records=50 for inventory and processes in a worker pool with pause/resume.

Key points

  • poll() — fetch records and maintain group membership
  • max.poll.interval.ms — upper bound on time between poll() calls
  • session.timeout.ms — heartbeat failure threshold
  • heartbeat.interval.ms — typically session.timeout / 3
  • pause()/resume() — stop fetching while processing backlog locally

Timeouts explained

session.timeout.ms (45s): no heartbeat → member removed. max.poll.interval.ms (300s default): poll too slow → member removed. Heartbeat thread (Kafka 0.10.1+) decouples heartbeats from poll for standard consumers, but max.poll.interval still applies to processing time between polls.

VaultCommerce rollout checklist

Before promoting changes that touch the VaultCommerce order and payment event backbone, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. Never block poll() longer than max.poll.interval. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.

Java
// VaultCommerce — bounded batch with pause during heavy work
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 25);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600_000); // 10 min for batch jobs

while (running) {
    ConsumerRecords<String, OrderPlaced> records = consumer.poll(Duration.ofSeconds(1));
    if (records.isEmpty()) continue;
    Set<TopicPartition> partitions = records.partitions();
    consumer.pause(partitions);
    try {
        for (ConsumerRecord<String, OrderPlaced> r : records) {
            inventoryService.reserve(r.value()); // may take seconds each
        }
        consumer.commitSync();
    } finally {
        consumer.resume(partitions);
    }
}

Quick recall

Everything you need if you only revisit this box.

  1. Never block poll() longer than max.poll.interval.
  2. Use pause/resume for slow handlers.
  3. Lower max.poll.records for heavy processing.

Test yourself

Answer these before moving on — recall is what makes it stick.