Apache Kafka Tutorial 0/111 lessons ~6 min read Lesson 24

    Consumer API

    What is Consumer API?

    Course progress0%
    Focus
    9 guided sections
    Practice signal
    Examples included
    Career prep
    Foundation builder

    Introduction

    What is Consumer API? The Consumer API polls records, processes them, and commits offsets.

    Understanding the topic

    What happens — Consumer API:

    • Create KafkaConsumer with group.id.
    • Subscribe to topic list.
    • Loop: poll → process → commit.
    • Close consumer on shutdown.
    TermDescription
    Consumer groupA set of consumers sharing work — each partition goes to at most one member at a time.
    OffsetPosition in a partition log — where this consumer last read.
    Poll loopconsumer.poll() fetches batches of records; keep processing faster than max.poll.interval.ms.
    CommitSaving offset to Kafka after processing — sync or async, manual or auto.
    LagDifference between latest offset and consumer offset — key health metric.

    Visual explanation

    Pipeline view:

    text
    Consumer polls records from topic partitions
    Process business logic (DB, API, etc.)
    Commit offset after success
    Monitor consumer lag

    Step-by-step explanation

    1. Subscribe — Consumer joins a group and receives partition assignments.
    2. Poll — Fetch records in batches with consumer.poll().
    3. Process — Run business logic for each record.
    4. Commit — Save offset after successful side effects.
    5. Repeat — Continue polling; rebalance if group membership changes.

    Informative example

    Example:

    java
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "payment-service");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
    try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
    consumer.subscribe(List.of("payment-events"));
    while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> record : records) {
    process(record);
    }
    consumer.commitSync();
    }
    }

    Execution workflow

    1Consumer API workflow
    1 / 5

    Subscribe

    Consumer joins a group and receives partition assignments.

    Best practices

    • Use manual commits for critical side effects.
    • Make consumers idempotent with event IDs or unique constraints.
    • Keep processing under max.poll.interval.ms or use pause/resume patterns.
    • Send poison messages to DLT with failure metadata.

    Common mistakes

    • Committing before database writes complete.
    • Scaling consumers beyond partition count and expecting more throughput.
    • Blocking the poll loop with slow downstream calls.

    Summary

    Consumer API — Commit offsets only after successful processing; make handlers idempotent; monitor lag per partition.

    Ready to mark this lesson complete?Track your journey across the entire course.