The Consumer API & Poll Loop Internals — Building Your First Consumer
The Requirement
In the producer chapter, you sent clickstream events into Kafka. Now let's ask a simple question.
Who is reading them?
Right now, nobody. So today we build the reader. It will be the first piece of a real-time recommendation engine. It reads each click the moment it happens, so recommendations can update live instead of overnight.
Why This Can't Be a CLI Job Either
You have already read events with kafka-console-consumer.sh. It works well for checking your work. But let's look at what it really does.
CLI vs. an application consuming events on its own
Who is reading the output of the CLI tool?
A human. It prints events on a screen and you look at them.
Can a recommendation engine work like that?
No. There is no person sitting in front of a terminal, updating recommendations by hand. We need code that reads every event and reacts on its own. It must run all day, every day.
It is the same story as producers. We need a client library, this time on the reading side. In Python, that library is confluent-kafka. You already installed it earlier, when we built producers.
Output / NoteGood to know:
confluent-kafkais a thin Python layer on top of librdkafka, a C library that does the real Kafka work. That is why some things happen "in the background" inside your process. We will see this soon.
Let's Build It
Here is what we are building:
Output / NoteA script that plays the ingestion layer of a
RecommendationConsumer. It reads every clickstream event as it arrives and reacts to it. To start simple, "react" means "print it".
We build it in four steps.
Step 1: Configure the Consumer
pythonfrom confluent_kafka import Consumer BOOTSTRAP_SERVERS = "localhost:9092,localhost:9094,localhost:9095" TOPIC = "clickstream-events" GROUP_ID = "recommendation-service" consumer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, "group.id": GROUP_ID, "auto.offset.reset": "earliest", } consumer = Consumer(consumer_config)
There are three settings here. Two of them need a closer look.
What is group.id?
It is the name of the group this consumer belongs to. It is required. Every consumer belongs to a group, even a "group" of one. In the next lecture we will see what a group really does. For now, just remember: no group.id, no working consumer.
What is auto.offset.reset?
Let's ask a question first. Suppose you start a consumer with a brand-new group.id. Kafka has never seen this group. Where should it start reading?
That is exactly what this setting answers.
"earliest"means start from the beginning of the topic."latest"means start from events produced from now on.
What if I don't set it?
Then the default is "latest". Your consumer would start, show nothing, and you would think it is broken. So it is a good habit to always set this on purpose.
Does this setting apply every time the consumer starts?
No. It only matters when the group has no saved position yet. Once the group has saved its position, that position is used instead. We will prove this in a minute, in the "run it twice" test.
Step 2: Subscribe to the Topic
pythonconsumer.subscribe([TOPIC], on_assign=on_assign)
Notice that subscribe() takes a list. A consumer can subscribe to many topics at once.
There is also on_assign=on_assign. This is a small function we wrote, and it prints which partitions this consumer received:
pythondef on_assign(consumer, partitions): """A first peek at partition assignment. We go deep on this soon.""" print(f"Partitions assigned to this consumer: {[p.partition for p in partitions]}")
Don't worry about it now. It is only here so you can see something happen when the consumer joins. We explain it fully in the next lecture.
Step 3: The Poll Loop
This is the heart of every Kafka consumer you will ever write.
pythonwhile True: msg = consumer.poll(timeout=1.0) if msg is None: continue if msg.error(): if msg.error().fatal(): print(f"Fatal consumer error, stopping: {msg.error()}") break print(f"Consumer error (non-fatal): {msg.error()}") continue event = json.loads(msg.value().decode("utf-8")) print( f"[partition {msg.partition()} offset {msg.offset()}] " f"user={event['user_id']} event={event['event_type']}" ) count += 1
(The full file also has import json at the top and count = 0 before the loop.)
Let's go through it slowly.
Does Kafka push data to the consumer?
No. A consumer pulls. It keeps calling poll(), again and again, and asks: "Do you have something for me?"
What does poll() return?
One of three things:
- A message, when data is available.
None, when nothing arrived within the timeout (here, 1 second).- A message object that carries an error instead of data.
Is None a problem?
No. It is normal. It only means "quiet second". That is why the loop just says continue and polls again.
And what about the error case?
Look at msg.error(). If it is not empty, this object is not a real event. Do not try to read its value. Most errors are temporary, and the client retries on its own, so we just print them and continue. But if msg.error().fatal() is true, the client cannot recover, so we stop the loop.
Output / NoteTip: If your cluster is not running, you will not see a nice Python error. You will see red log lines like
Connection refusedon your screen. If you see those, checkdocker psfirst.
What does the last part do?
It decodes the message. In the producer chapter, we did json.dumps() and encoded the text to bytes. Here we do the reverse: decode the bytes, then json.loads(). Then we print the partition and offset, the same ones you saw from the producer side.
Wait, why does poll() do more than return messages?
Good catch. In the Python client, poll() has three jobs:
- It returns the next message.
- It reports errors and events.
- It runs your callbacks, such as
on_assign.
So your on_assign function does not run on its own. It runs inside your poll() call. Remember this. It will matter a lot in the next two lectures.
Step 4: Shut Down Cleanly
pythonfinally: consumer.close() print(f"Done. This run consumed {count} messages.")
Why do we always call close()?
Because close() does three things:
- It stops consuming.
- It commits the final offsets (with the default settings).
- It tells the group: "I am leaving on purpose."
What happens if we skip it?
The final offsets are not committed, and the group does not know you left. It waits until your session times out (45 seconds by default) before it notices. Until then, your partitions sit unread.
Running It
Step 1: Start Your Cluster
bashdocker compose start docker ps
Step 2: Run the Consumer
Download 5.1_recommendation_consumer.py into a new chapter-5 folder. confluent-kafka is already installed from earlier.
bashpython 5.1_recommendation_consumer.py
If clickstream-events already has data from your producer runs, you will first see the assigned partitions, and then all the old events scroll by. That is auto.offset.reset = earliest at work.
Want to see events arrive live? Open a second terminal and run one of your clickstream producer scripts from the producer chapter. Each new event should appear here within about a second.
Now press Ctrl+C. Look at the last line. It tells you how many messages this run consumed. Remember that number.
Step 3: Run It Again (Important!)
Run the very same command a second time:
bashpython 5.1_recommendation_consumer.py
What do you see?
Almost nothing. You see the assigned partitions, and then silence. Press Ctrl+C, and the last line says This run consumed 0 messages.
Is the consumer broken?
No. It is working exactly as designed. Think about it:
- On the first run, this group had no saved position. So
auto.offset.resetstarted us at the beginning. - When we pressed
Ctrl+C,close()saved our position in Kafka. - On the second run, the group had a saved position. So Kafka said "you already read everything up to here" and
auto.offset.resetwas ignored.
Let's confirm this from the cluster's side. Open a broker terminal:
bashdocker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin
bashkafka-consumer-groups.sh --describe --group recommendation-service --bootstrap-server kafka-1:29092
Look at the CURRENT-OFFSET and LAG columns. The lag should be 0 for every partition. The group is fully caught up. (No consumer is running now, so the consumer columns will show -. That is fine.)
Now produce a few new events with a producer script, and run the consumer once more. This time you will see only the new events.
Output / NoteOptional experiment: Change
GROUP_IDin the script to a brand-new name, like"my-experiment", and run it. Since this group has no saved position, you will see the whole history again. Change it back to"recommendation-service"afterwards, because the next lectures use that name.
This is a preview of two big topics. Soon we will look at offsets in detail, and we will also learn how to push a group's position back on purpose.
Step 4: Stop Your Cluster
bashdocker compose stop
Now — What's Actually Happening Inside That Loop?
You watched a consumer read data continuously, and you never told it when to fetch from the broker. Let's look at its whole life: start, running, and stop.
Consumer lifecycle: startup, steady state, shutdown
Step 1: Startup
What happens when you write Consumer(config) and subscribe()?
The client starts background threads inside your Python process. They connect to the brokers, find the group coordinator, and start joining the group. All this happens in the background, while your code moves on to the next line.
Do you have your partitions at this point?
Not yet. The join takes a moment. Your partitions are handed over when poll() runs your on_assign callback. That is why, in your terminal, the "Partitions assigned" line appears just before the first message and not before.
So the correct picture is: the work starts at subscribe(), and the result reaches you inside poll().
Step 2: Steady State
This is where your script spends almost all of its time. Two things happen side by side.
Where do messages come from when you call poll()?
Not from the network. Background fetching keeps asking the brokers for batches of messages and stores them in a local queue inside your process. Your poll() call just takes the next message from that queue. So most of the time, the data is already there when you ask for it.
How does Kafka know your consumer is still alive?
Through heartbeats. A background timer sends a small "I am alive" signal to the group coordinator. By default, this happens every 3 seconds (heartbeat.interval.ms). If the coordinator hears nothing for 45 seconds (session.timeout.ms), it decides your consumer is dead and gives its partitions to someone else.
Does calling poll() send the heartbeat?
No. Heartbeats run on their own timer. Not calling poll() does not stop them.
So can a consumer stay alive forever without calling poll()?
No, and this is the second clock. There is another setting, max.poll.interval.ms. Its default is 5 minutes. If your code goes longer than that between two poll() calls, the client decides your processing loop is stuck. It leaves the group, and the partitions move to another consumer.
Let's put the two clocks side by side:
| Clock | Setting | Default | What it proves |
|---|---|---|---|
| Heartbeat | session.timeout.ms | 45 seconds | The process is alive |
| Poll gap | max.poll.interval.ms | 5 minutes | Your loop is making progress |
Why does this matter for us?
Because of a very common bug: "my consumer keeps leaving the group." Usually the reason is slow work inside the loop, such as a slow database call or a long calculation. The process is alive and heartbeats are fine, but poll() is not called often enough. Keep this in mind whenever you add processing to the loop.
Step 3: Shutdown
What happens at close()?
We saw it already: stop consuming, commit the final offsets, and leave the group. Because the group is told right away, it can give your partitions to another consumer immediately. It does not have to wait for a timeout and guess whether you crashed.
One more thing to notice.
By default, offsets are committed by a background timer (every 5 seconds), not by your code. That is why close() could save your position, and it is also a risk. The timer does not know if you finished processing a message. We will look at this danger properly when we cover offset management.
Common Mistakes at This Stage
Not setting auto.offset.reset on purpose. The default is latest. A new group then sees nothing until new events arrive.
Forgetting close(). The group needs up to 45 seconds to notice you left, and your last offsets are not committed.
Treating every poll() result as a message. Always check msg.error() before reading msg.value().
Doing slow work inside the loop. Long gaps between poll() calls can push the consumer out of the group.
Checklist
- Understood why a recommendation engine needs a consumer client, not a CLI tool
- Built the consumer step by step, and understood
group.idandauto.offset.reset(and its default) - Understood that a consumer pulls with
poll()and thatNoneis normal - Understood that
poll()returns messages, reports errors, and runs callbacks - Ran the consumer twice, and can explain why the second run printed nothing
- Walked through the lifecycle: startup, steady state, shutdown
- Understood the two clocks: heartbeats (
session.timeout.ms) and poll gap (max.poll.interval.ms) - Understood what
close()does, and why we always call it
You built a real consumer, and you know what it does between every line. Next: we run the same consumer three times with the same group.id, and watch Kafka split the partitions between them.