๐ฅ Kafka Consumer API & Monitoring
๐ฃ The Pull-Based Model
What Existed Previously: In legacy messaging systems (like RabbitMQ), the server pushes messages to the consumers as fast as possible.
Problems Faced: If the consumer is a slow machine, it gets completely overwhelmed by the "firehose" of pushed messages and eventually crashes (Out of Memory).
How Present Technology Solves It:
Kafka uses a Pull-Based Model. The broker just holds the data. The Consumer actively calls a poll() function to fetch data only when it is ready for more. This naturally controls the ingestion rate and prevents the consumer from ever crashing due to a flood of data.
๐พ Offset Management (Commits)
Anatomy Breakdown: When a consumer reads a message, it needs to tell Kafka, "I have successfully processed up to offset 45." This is called Committing the Offset.
- Auto-Commit: The Consumer API automatically commits its offset in the background every few seconds.
- Risk: If the consumer pulls data, Auto-Commits instantly, but then crashes before processing the data, that data is skipped forever!
- Manual Commit: You turn off Auto-Commit and explicitly write
consumer.commit()in your code only after you have successfully saved the data to your database. This is slower but guarantees zero data loss.
โ๏ธ Consumer Rebalancing
Follow the Data:
- You have a Topic with 4 Partitions.
- You have a Consumer Group with 2 Consumers (Node A and Node B).
- Kafka automatically assigns 2 partitions to Node A, and 2 partitions to Node B.
- Suddenly, Node B crashes.
- Kafka detects the crash and triggers a Rebalance. It re-assigns all 4 partitions to Node A so processing can continue!
During a rebalance, ALL consumers in the group are temporarily paused. If you use a ConsumerRebalanceListener, you can trigger a callback to manually commit your offsets right before your partitions are taken away!
๐ Monitoring Kafka
To ensure your cluster is healthy, you must monitor specific metrics.
The Golden Metric: Consumer Lag Consumer Lag is the difference between the latest message produced and the latest message consumed. If the producer is at offset 100, and your consumer is stuck at offset 50, your Lag is 50. If Lag keeps growing, your consumer is too slow, and you need to add more consumer nodes to the group!
Monitoring Stack:
- Prometheus: Automatically scrapes Kafka broker metrics (like throughput and lag) every few seconds.
- Grafana: Visualizes those metrics in beautiful real-time dashboards so you can set up alerts before outages happen.
๐งช Practice Drill
// Try answering these:
Q1. Why does Kafka use a Pull-based model instead of pushing data to consumers?
Q2. Which offset commit strategy guarantees you won't lose data if your consumer crashes mid-processing?
Q3. What happens to the partitions if a consumer node suddenly leaves a consumer group?
Q4. If your producer is generating 1,000 messages per second, but your consumer is only processing 500 messages per second, what metric will immediately start rising?
๐ก Click for Solutions
A1. It prevents consumers from being overwhelmed. The consumer dictates its own pace by pulling data only when it is ready.
A2. Manual Commits (Committing the offset after the data is fully processed and saved).
A3. A Rebalance occurs, and Kafka automatically assigns those orphaned partitions to the remaining healthy consumers in the group.
A4. Consumer Lag.
โ ๐ค Spark Streaming & Kafka Integration | Next Topic โ ๐๏ธ Advanced Streaming Architectures