Apache Kafka for Data & AI Engineers

Building a Replay Tool — Rewinding a Consumer to a Point in Time


The Requirement

In the last lecture, we compared two ways to reprocess the past. We built the first one: reset the production group's offsets, and let the real consumer redo its own work.

Now let's look at the second situation. This time:

  • Production must keep running. No pause, no downtime.
  • You do not want production to redo anything. You want a separate look at a time window.
  • The result should go somewhere else: a separate table, a file, an analyst's notebook.

For example, the bug fix is deployed, and an analyst asks: "Show me what the corrected logic would have produced for the last two hours, without touching the live pipeline."

So far, every consumer in this chapter either started fresh or continued from where it left off. Today we add a new ability:

Output / Note

Deliberately rewind to any point you choose.

And we do it in code we control.


Why subscribe() Can't Do This

Every consumer you have built so far calls subscribe(). What does subscribe() give up?

Control. It hands the partition assignment to the group coordinator. You never choose your partitions, and you cannot choose where to start, except through the narrow auto.offset.reset case you saw in the offset lecture.

subscribe() vs. assign()subscribe() vs. assign()

To rewind on purpose, you need the other mode: assign().

With assign(), your own code picks the partitions, and it can also pick the start offset for each one. The coordinator is not involved, and there is no rebalancing.

What is the trade-off?

You get full control, but you lose the group's help. Nobody rebalances for you, and if your tool stops, nobody takes over its partitions. That is fine for a one-off replay tool. It is not what you want for a long-running production service.


Let's Build It

Here is what we are building: a standalone replay tool. It is deliberately separate from the live recommendation-service group. It rewinds to a number of minutes ago, reprocesses everything from there up to now, and stops.

Step 1: Use a Different group.id on Purpose

python
REPLAY_GROUP_ID = "recommendation-service-replay-tool" consumer = Consumer({ "bootstrap.servers": BOOTSTRAP_SERVERS, # A consumer object always needs a group.id, even when we assign # partitions ourselves. Here it is only a name. "group.id": REPLAY_GROUP_ID, # This tool is read-only. It should never save any position. "enable.auto.commit": False, })

Why a different group.id?

This is on purpose. A replay tool should almost never share a group.id with the production consumer it investigates. If it did, the tool could interfere with the live group's committed position. Give it its own identity.

If we assign partitions ourselves, why do we still need a group at all?

Because the consumer object always requires a group.id, even when it never joins a group. Here it is just a name.

And why turn off auto-commit?

The replay tool must only read. With auto-commit on, it would save its own positions under its own group name. That would not hurt production, but it would leave clutter behind. Turning it off keeps the tool clean.

Step 2: Find the Partitions

python
metadata = consumer.list_topics(TOPIC, timeout=10) partition_ids = sorted(metadata.topics[TOPIC].partitions.keys())

We ask the cluster how many partitions the topic has, instead of writing 3 in the code. If the topic grows later, the tool still works.

Step 3: Turn a Timestamp Into Real Offsets

python
partitions_with_timestamp = [ TopicPartition(TOPIC, p, target_timestamp_ms) for p in partition_ids ] resolved = consumer.offsets_for_times(partitions_with_timestamp, timeout=10)

Do you ever know the offset you want?

Almost never. You do not know "the offset I want is 4,821". What you know is "I want everything from 2 hours ago".

offsets_for_times() bridges that gap. You give it one TopicPartition per partition, with a timestamp in the place of the offset. It gives back the same partitions with the real offset.

Two details matter:

  • The offset returned is the earliest offset whose timestamp is at or after the time you asked for.
  • If no message is that recent, the offset is -1. That simply means "nothing to replay in this partition".

Which timestamp does Kafka compare with?

The timestamp stored on each message. By default, this is the time the producer created the message (the topic setting message.timestamp.type is CreateTime). So "2 hours ago" means when the events happened, as reported by your producers.

The time we compute is in milliseconds since 1970 (time.time() is always UTC-based), so time zones do not matter here.

python
minutes_ago = int(sys.argv[1]) target_timestamp_ms = int((time.time() - minutes_ago * 60) * 1000)

Step 4: Decide Where the Replay Stops

python
end_offsets = {} for tp in resolved: _low, high = consumer.get_watermark_offsets(TopicPartition(TOPIC, tp.partition), timeout=10) end_offsets[tp.partition] = high

Why do we need this?

A live topic never ends. New events keep arriving. If our tool just kept reading, it would never finish. So, before we start, we note where each partition ends right now. The high value is the offset of the next message to be written, so the last existing message is high - 1. Our replay stops at that point. It processes a fixed window: from "N minutes ago" up to "the moment I started".

Step 5: Assign Directly to the Resolved Offsets

python
consumer.assign([tp for tp in resolved if tp.partition in pending])

Here pending is the set of partitions that actually have something to replay. Since each entry in resolved already carries the right offset, assign() alone is enough to start reading from exactly the right place.

What about seek()?

seek() is the other tool. It moves a partition you are already reading to a new offset. We do not need it here, because we pass the offsets inside assign(). But it is worth knowing, and it has one trap. This is a short example, and it is not part of our script:

python
consumer.assign([TopicPartition("clickstream-events", 0)]) msg = consumer.poll(timeout=5.0) # let the partition start being read first consumer.seek(TopicPartition("clickstream-events", 0, 100)) # then jump to offset 100

If you call seek() immediately after assign(), before any poll(), it fails with Local: Erroneous state, because the partition is not being read yet. Call it after the consumer has started reading, or, simpler, pass the offset in assign() as we do.

Step 6: The Same Poll Loop, With a Stopping Point

python
while pending: msg = consumer.poll(timeout=1.0) if msg is None: continue if msg.error(): print(f"Consumer error: {msg.error()}") continue partition = msg.partition() if partition not in pending: continue if msg.offset() < end_offsets[partition]: event = json.loads(msg.value().decode("utf-8")) print( f"[replay] partition {partition} offset {msg.offset()} " f"user={event['user_id']} event={event['event_type']}" ) count += 1 # Reached the last message that existed when we started? if msg.offset() >= end_offsets[partition] - 1: pending.discard(partition)

From here, it is the same poll() loop you have written all chapter. Rewinding only changes where you start, not how you read once you are there. The only new part is the stopping rule: when a partition reaches its last message from the start, we remove it from pending. When pending is empty, the loop ends and the tool prints a summary.


Running It

Step 1: Start Your Cluster

bash
docker compose start docker ps

Step 2: Make Sure There Are Recent Events

The tool replays the last N minutes. If nothing was produced in that time, there is nothing to replay. So, first run one of your clickstream producer scripts from the producer chapter, to have some fresh events.

Step 3: Run the Replay Tool

Download 5.6_replay_by_timestamp.py into your chapter-5 folder and run it from your host terminal. Give it how far back to go, in minutes:

bash
python 5.6_replay_by_timestamp.py 30

First you see the starting offset for each partition. Then you see every event from roughly the last 30 minutes, reprocessed on demand. At the end, you see Replay complete with the number of messages. The tool stops by itself.

What if a partition says "nothing to replay"?

That partition has no message in your time window. It is not an error.

Try two more runs to see the edges:

bash
python 5.6_replay_by_timestamp.py 1
bash
python 5.6_replay_by_timestamp.py 100000

The first is a very small window, so you will probably see "Nothing to replay". The second is a very large window, so you replay everything the topic still holds.

Step 4: Confirm the Live Group Is Untouched

This is worth proving to yourself. Open a broker terminal:

bash
docker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin

Run this before and after a replay:

bash
kafka-consumer-groups.sh --describe --group recommendation-service --bootstrap-server kafka-1:29092

The committed offsets of recommendation-service do not change, because the replay tool never touched that group. You can also look at the replay tool's own group:

bash
kafka-consumer-groups.sh --describe --group recommendation-service-replay-tool --bootstrap-server kafka-1:29092

It has no saved positions at all, since the tool never commits anything.

Step 5: Stop Your Cluster

bash
docker compose stop

Where Does the Replayed Data Go?

Our tool only prints each event. A real replay tool would send the reprocessed result somewhere: another table, another topic, a file for analysis.

That is the whole point of this approach. The replay runs beside production. It does not change production's own output. Compare it with the previous lecture. There, the production consumer itself redid the work, and production's real output was corrected. Here, a separate reader gives you a second, independent view, and production is untouched, unless you deliberately build the tool to write there too.


Common Mistakes at This Stage

Sharing the production group.id. The replay tool could disturb the live group's position. Always use its own name.

Hard-coding the number of partitions. Ask the cluster instead. A topic that grows will silently break a hard-coded loop.

Calling seek() right after assign(). It fails with Erroneous state. Pass the offsets in assign(), or seek after the first poll().

Forgetting that -1 means "nothing there". If you pass -1 straight to assign(), that partition starts at the end and waits for new events. Filter these partitions out, as our tool does.

Trusting the timestamp blindly. It is the time your producers put on the message. If a producer has a wrong clock, the replay window will be wrong too.


Checklist

  • Understood when a separate replay reader is the right choice, compared with resetting the group
  • Understood why subscribe() cannot give you deliberate control over the starting position
  • Understood the trade-off of assign(): full manual control, no automatic rebalancing
  • Understood why a replay tool should use its own group.id and turn off auto-commit
  • Used offsets_for_times() to turn "X minutes ago" into real offsets, and know that -1 means "nothing at or after that time"
  • Understood why the tool records the end offsets first, so the replay stops by itself
  • Know what seek() does, and the trap of calling it right after assign()
  • Ran a real replay and confirmed that the live group's committed offsets did not change

Chapter Wrap-Up: The Consumer API, at a Glance

You started this chapter with a script that printed events. You end it with a consumer you could run in production, and the tools to fix things when they go wrong. Here is the whole path in one table.

TopicThe key ideaThe settings and calls to remember
The poll loopA consumer pulls. poll() returns messages, reports errors, and runs callbackspoll(), msg.error(), close()
Consumer groupsWithin a group, each partition is read by exactly one consumer. More consumers than partitions means idle onesgroup.id, client.id
RebalancingA group change re-runs JoinGroup and SyncGroup. Eager stops everyone. Cooperative moves only what must movepartition.assignment.strategy
OffsetsA committed offset is the next one to read. Mark it done only after success. Dead-letter poison pillsauto.offset.reset, enable.auto.offset.store, store_offsets()
ReprocessingReset the group to make production redo its work. Use a separate reader to look without touching production--reset-offsets, assign(), offsets_for_times()

If you had to build a production consumer tomorrow, this is the checklist you would follow:

  • Choose auto.offset.reset on purpose. Remember that it applies only when there is no valid committed offset.
  • Use cooperative-sticky so a routine deploy does not stop the whole group.
  • Turn off enable.auto.offset.store, and call store_offsets() after each message succeeds.
  • Send messages that can never succeed to a dead-letter topic, and confirm the send worked.
  • Keep the work between poll() calls short, so the consumer stays in the group.
  • Always call close() when the process ends.
  • Make processing safe to repeat, because at-least-once, retries and reprocessing all mean a message can be seen twice.
  • Give every tool its own group.id, so an investigation never disturbs production.

You built, scaled, rebalanced, hardened, and rewound a real Kafka consumer, and you understand what happens behind every line. That completes the Consumer API.