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.
| Term | Description |
|---|---|
| Consumer group | A set of consumers sharing work — each partition goes to at most one member at a time. |
| Offset | Position in a partition log — where this consumer last read. |
| Poll loop | consumer.poll() fetches batches of records; keep processing faster than max.poll.interval.ms. |
| Commit | Saving offset to Kafka after processing — sync or async, manual or auto. |
| Lag | Difference 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
- Subscribe — Consumer joins a group and receives partition assignments.
- Poll — Fetch records in batches with consumer.poll().
- Process — Run business logic for each record.
- Commit — Save offset after successful side effects.
- 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 / 5Subscribe
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.