Assignment: ShopStream
Part 3: The solution, step by step
How to use this part
This part walks you through one complete solution. It follows the same eleven steps as Part 2.
- Each step says what you build, why, and how you check it.
- You will see short pieces of code. They show the important ideas. They are not the full programs.
- The full, working code is in a separate download:
shopstream-solution.zip. Use it to compare with your own code, or to unblock yourself.
A good way to work: read one step, close this file, write the code yourself, run it, then compare.
All commands are run from the folder that has your files.
Step 0: Get ready
Start your three brokers. Then:
pip install -r requirements.txt
python setup_topics.py
You should see:
created: shop.orders
created: shop.orders.dlq
If the topics already exist, you will see already exists: instead. That is fine.
To start again from zero at any time:
python setup_topics.py --reset
Step 1: One place for settings
What: a small file, config.py, with the settings that every program shares.
pythonimport os BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "localhost:9092,localhost:9094,localhost:9095") ORDERS_TOPIC = "shop.orders" DLQ_TOPIC = "shop.orders.dlq" ORDERS_DB = "orders.db" SALES_DB = "sales.db" ZONES = ("NORTH", "SOUTH", "EAST", "WEST")
Why: if a topic name or a port changes, you change it in one place.
The zones are a tuple, in a fixed order. The position of a zone in this list will be its partition number. NORTH is 0. WEST is 3.
The starter kit already gives you this file.
Step 2: Check an event
What: a function parse_event() that takes the raw bytes of a message. It returns a clean dictionary, or raises an error with a reason.
Why: both consumers need the same rules. Writing them once keeps them the same.
First, a small error class. It carries the reason code:
pythonclass InvalidEvent(Exception): def __init__(self, reason): super().__init__(reason) self.reason = reason
Then the checks, from the cheapest to the most detailed:
pythondef parse_event(raw): if raw is None: raise InvalidEvent("EMPTY_MESSAGE") try: event = json.loads(raw) except ValueError: raise InvalidEvent("INVALID_JSON") if not isinstance(event, dict): raise InvalidEvent("NOT_AN_OBJECT") for field in ("event_id", "event_type", "order_id", "zone"): if not event.get(field): raise InvalidEvent("MISSING_" + field.upper()) ...
After this, the function checks the event type, the zone, and (for ORDER_CREATED) the customer and the amount.
A small trap: in Python, True is a number. So isinstance(True, int) is True. The amount check must reject booleans first:
pythonif isinstance(amount, bool) or not isinstance(amount, (int, float)) or amount <= 0: raise InvalidEvent("BAD_AMOUNT")
Check it: try a few examples in a Python shell.
pythonparse_event(b"not json") # raises InvalidEvent: INVALID_JSON parse_event(b'{"event_id": "E1"}') # raises InvalidEvent: MISSING_EVENT_TYPE
Nothing here needs Kafka. That is the point. You can test the rules alone.
Step 3: The zone rule
What: one small function.
pythondef zone_partition(zone, partition_count): if zone not in ZONES: return None return ZONES.index(zone) % partition_count
Why None? For an unknown zone we have no answer. None means: let Kafka decide.
Why % partition_count? It keeps the answer inside the number of partitions, even if the topic changes later.
Check it:
pythonzone_partition("NORTH", 4) # 0 zone_partition("WEST", 4) # 3 zone_partition("CENTRAL", 4) # None
Step 4: The shared producer code
What: a module, feed.py, with everything both producers need.
4a. Build the producer
pythondef build_producer(client_id): return Producer({ "bootstrap.servers": BOOTSTRAP, "client.id": client_id, "enable.idempotence": True, "acks": "all", "linger.ms": 50, "batch.size": 65536, "compression.type": "lz4", })
Every setting here matches a row in the Part 2 table.
4b. Ask Kafka how many partitions there are
Do not hard-code 4. Ask the cluster:
pythonmetadata = producer.list_topics(topic, timeout=10) partition_count = len(metadata.topics[topic].partitions)
4c. Prepare one line
A line becomes a key, a value and a zone.
pythondef prepare(line): raw = line.encode("utf-8") try: event = json.loads(line) except ValueError: return None, raw, None # not JSON: send as it is if not isinstance(event, dict): return None, raw, None event.setdefault("event_time", now_iso()) order_id = event.get("order_id") key = order_id.encode("utf-8") if isinstance(order_id, str) and order_id else None zone = event.get("zone") if isinstance(event.get("zone"), str) else None return key, json.dumps(event).encode("utf-8"), zone
Look at the first return. A broken line goes out unchanged, with no key and no zone.
4d. The delivery callback
The callback runs after Kafka answers. It counts what arrived in each partition:
pythondelivered = Counter() failed = [] def on_delivery(err, msg): if err is not None: failed.append(str(err)) print(f"[{label}] FAILED: {err}") else: delivered[msg.partition()] += 1
4e. The send loop
For each line, build the arguments. Add partition only when our rule has an answer:
pythonkey, value, zone = prepare(line) kwargs = {"key": key, "value": value, "on_delivery": on_delivery} partition = zone_partition(zone, partition_count) if partition is not None: kwargs["partition"] = partition
Then send, and handle a full queue:
pythonwhile True: try: producer.produce(ORDERS_TOPIC, **kwargs) break except BufferError: producer.poll(0.5) # let some messages leave, then try again
After each message, call producer.poll(0). This lets the callbacks run.
4f. The end
pythonremaining = producer.flush(30) print(f"[{label}] delivered per partition: {dict(sorted(delivered.items()))}")
If failed is not empty, or remaining is not 0, return a non-zero exit code.
Step 5: The two producers
What: two tiny programs. Each reads its own file and calls the shared code.
python# order_producer.py parser.add_argument("--file", default=os.path.join(HERE, "data", "orders.jsonl")) parser.add_argument("--rate", type=float, default=20) parser.add_argument("--limit", type=int, default=0) args = parser.parse_args() sys.exit(send_file(args.file, "order-service", args.rate, args.limit))
fulfilment_producer.py is the same, with fulfilment.jsonl and the name fulfilment-service.
Check it: run the Order Service at full speed.
python order_producer.py --rate 0
You will see something like:
[order-service] topic shop.orders has 4 partitions
[order-service] queued 232 messages, 0 still waiting after flush
[order-service] delivered per partition: {0: 50, 1: 81, 2: 42, 3: 59}
The four numbers add up to 232, the number of lines in the file.
Look at the numbers. Partition 1 (SOUTH) has the most messages. This is the busy-zone cost we talked about in Part 2.
Now run the Fulfilment Service:
python fulfilment_producer.py --rate 0
The numbers add up to 139.
Both producers use partition 0 to 3 by zone. The few bad lines are spread by Kafka, so your counts may be a little different.
Step 6: The database code
What: a module, db.py, with two small SQLite databases.
6a. The orders database
Two tables: the orders, and the event IDs we have already seen.
sqlCREATE TABLE IF NOT EXISTS orders ( order_id TEXT PRIMARY KEY, zone TEXT NOT NULL, customer_id TEXT, amount INTEGER, status TEXT NOT NULL, status_rank INTEGER NOT NULL ); CREATE TABLE IF NOT EXISTS processed_events ( event_id TEXT PRIMARY KEY );
status_rank holds the status number from Part 2. It is what stops a status from going backwards.
6b. Save one event in one transaction
pythondef apply_order_event(conn, event): with conn: # one transaction seen = conn.execute( "INSERT OR IGNORE INTO processed_events (event_id) VALUES (?)", (event["event_id"],), ) if seen.rowcount == 0: return False # we have seen this event: do nothing conn.execute(UPSERT_ORDER, (...)) return True
Read it slowly:
INSERT OR IGNOREtries to add the event ID. If it was there already,rowcountis 0, and we stop.- If it is new, we update the order.
with conn:makes both steps one transaction. If anything fails, both are undone.
6c. The update rule
sqlINSERT INTO orders (order_id, zone, customer_id, amount, status, status_rank) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(order_id) DO UPDATE SET customer_id = COALESCE(excluded.customer_id, orders.customer_id), amount = COALESCE(excluded.amount, orders.amount), status = CASE WHEN excluded.status_rank > orders.status_rank THEN excluded.status ELSE orders.status END, status_rank = MAX(excluded.status_rank, orders.status_rank)
What each part does:
ON CONFLICT ... DO UPDATE: if the order already has a row, update it. If not, insert a new row.COALESCE(new, old): use the new value, but if it is empty, keep the old one. ASHIPPEDevent has no amount. It must not erase an amount we already have. And a lateORDER_CREATEDfills in the amount that a stub row did not have.CASE WHEN new rank > old rank: the status moves only forward.
6d. The sales database
Same idea, a different table:
sqlCREATE TABLE IF NOT EXISTS sales_by_minute ( zone TEXT NOT NULL, minute TEXT NOT NULL, orders INTEGER NOT NULL, revenue INTEGER NOT NULL, PRIMARY KEY (zone, minute) ); CREATE TABLE IF NOT EXISTS seen_events (event_id TEXT PRIMARY KEY);
And the save step:
sqlINSERT INTO sales_by_minute (zone, minute, orders, revenue) VALUES (?, ?, 1, ?) ON CONFLICT(zone, minute) DO UPDATE SET orders = orders + 1, revenue = revenue + excluded.revenue
To turn a Kafka timestamp (milliseconds) into a minute:
pythondef minute_of(timestamp_ms): moment = datetime.fromtimestamp(timestamp_ms / 1000, tz=timezone.utc) return moment.strftime("%Y-%m-%d %H:%M")
Check it: write a tiny test. Save the same event twice. Count the rows. You must get one row, not two. Save a DELIVERED event and then a SHIPPED event for the same order. The status must stay DELIVERED.
Step 7: The Order Loader
What: the first consumer. This is the biggest step, so we do it in parts.
7a. The consumer settings
pythonconsumer = Consumer({ "bootstrap.servers": BOOTSTRAP, "group.id": "order-loader", "client.id": args.name, "auto.offset.reset": "earliest", "enable.auto.commit": True, "enable.auto.offset.store": False, "partition.assignment.strategy": "cooperative-sticky", })
Match each line with a requirement:
group.id: requirement B1.auto.offset.reset: earliest: a new group starts from the beginning.enable.auto.offset.store: False: we mark the offsets (requirement B10).cooperative-sticky: requirement B2.
The --name option lets you start two copies with different names, so you can tell their log lines apart.
7b. A producer for the dead-letter topic
The loader needs a producer too. It is a normal producer with safe delivery:
pythondlq_producer = Producer({ "bootstrap.servers": BOOTSTRAP, "client.id": args.name + "-dlq", "enable.idempotence": True, "acks": "all", })
7c. Sending a bad message to the DLQ
pythondef send_to_dlq(msg, reason): problems = [] def on_delivery(err, _msg): if err is not None: problems.append(err) source = f"{msg.topic()}[{msg.partition()}]@{msg.offset()}" dlq_producer.produce( DLQ_TOPIC, key=msg.key(), value=msg.value(), headers=[ ("error", reason.encode()), ("source", source.encode()), ("failed_at", datetime.now(timezone.utc).isoformat().encode()), ], on_delivery=on_delivery, ) dlq_producer.flush(30) if problems: raise KafkaException(problems[0])
Three things to notice:
- The original key and value are copied without change.
- The three headers make the message easy to understand later.
flush(30)waits until Kafka confirms. If the copy failed, the function raises an error. The program stops before it marks the offset. This is the safe order again.
Waiting for each message is slow, but bad messages are rare. The simple way is the right way here.
7d. Know when partitions move
pythondef on_assign(c, partitions): log("assigned: " + str(sorted(p.partition for p in partitions))) def on_revoke(c, partitions): log("revoked: " + str(sorted(p.partition for p in partitions))) try: c.commit(asynchronous=False) # save our progress before giving them away except KafkaException as e: if e.args[0].code() != KafkaError._NO_OFFSET: raise consumer.subscribe([ORDERS_TOPIC], on_assign=on_assign, on_revoke=on_revoke)
The try/except is there because a consumer with nothing marked yet has nothing to commit. That is not an error for us.
7e. The loop
This is the heart of the program:
pythonwhile not stopper.stop: msg = consumer.poll(1.0) if msg is None: continue if msg.error(): if msg.error().code() == KafkaError._PARTITION_EOF: continue raise KafkaException(msg.error()) try: event = parse_event(msg.value()) except InvalidEvent as bad: send_to_dlq(msg, bad.reason) else: apply_order_event(db, event) consumer.store_offsets(msg) # only now: mark the message as done
Follow one message through. It is checked. It is saved (database or dead-letter topic). Only then is it marked.
7f. Stopping cleanly
Wrap the loop in try ... finally. In finally, call consumer.close(). This leaves the group and saves the last marks. A small helper catches Ctrl+C and sets stopper.stop:
pythonclass StopFlag: def __init__(self): self.stop = False signal.signal(signal.SIGINT, self._handle) signal.signal(signal.SIGTERM, self._handle) def _handle(self, signum, frame): self.stop = True
Check it: start one loader in a terminal.
python order_loader.py --name loader-1
[loader-1] started. Press Ctrl+C to stop.
[loader-1] assigned: [0, 1, 2, 3]
[loader-1] bad message shop.orders[3]@33 -> DLQ (MISSING_ZONE)
...
Because the events are already in the topic, it starts working at once. Let it run until the messages stop. Then press Ctrl+C. You will see a summary line:
[loader-1] summary: {'saved': 358, 'duplicate': 3, 'dlq': 10}
Count them: 358 + 3 + 10 = 371. That is 232 + 139, every line from both files.
Step 8: The Sales Aggregator
What: the second consumer. It looks very much like the loader. Only a few things are different.
1. A different group:
python"group.id": "sales-aggregator",
2. It counts only ORDER_CREATED. Everything else is ignored, bad messages too:
pythontry: event = parse_event(msg.value()) except InvalidEvent: counts["ignored"] += 1 else: if event["event_type"] != "ORDER_CREATED": counts["ignored"] += 1 else: timestamp_ms = msg.timestamp()[1] # the time Kafka stored add_sale(db, event, timestamp_ms)
3. The clock. msg.timestamp() returns two values: the type of timestamp, and the time in milliseconds. We use the second.
4. The same safe order. Save with add_sale() first. Then consumer.store_offsets(msg).
5. A summary at the end: a small table of orders and revenue per zone.
Check it:
python sales_aggregator.py
When it has caught up, stop it with Ctrl+C. You will see:
Sales so far
ZONE ORDERS REVENUE (INR)
NORTH 28 110,172
SOUTH 40 109,210
EAST 21 56,429
WEST 31 148,769
These numbers are the ones in expected_results.json. Your aggregator and loader read the same topic, but they did not disturb each other. That is two consumer groups at work.
Step 9: Run everything and check
Start again from zero so that you see the whole flow.
python setup_topics.py --reset
Open five terminals.
| Terminal | Command |
|---|---|
| 1 | python order_loader.py --name loader-1 |
| 2 | python order_loader.py --name loader-2 |
| 3 | python sales_aggregator.py |
| 4 | python order_producer.py |
| 5 | python fulfilment_producer.py |
Start 1, 2 and 3 first. Then start 4 and 5.
Watch the loaders. In the end, each one has two partitions, for example [0, 2] and [1, 3]. The numbers can be different on your computer. The idea is the same: four partitions, two readers, two partitions each.
If you started the loaders one after the other, the first one may first show assigned: [0, 1, 2, 3]. When the second one joins, the first shows revoked: for the two partitions it gives away. This is cooperative rebalancing. Only the partitions that must move, move.
When the producers finish and the consumers go quiet, press Ctrl+C in terminals 1, 2 and 3. Then:
python verify_results.py
You should see ALL CHECKS PASSED.
If a check fails
| Message | Likely reason |
|---|---|
number of orders is too low | The loader stopped early, or a store_offsets happened before the save |
orders per status is wrong | The status rank rule is wrong, or COALESCE is missing |
orders and revenue per zone in sales.db is wrong | You counted events other than ORDER_CREATED, or the duplicate check is missing |
| DLQ count is too low | A bad message was not sent to the DLQ, or the DLQ copy was not flushed |
| DLQ count is too high | The aggregator also sent bad messages to the DLQ |
sales.db not found | The aggregator did not run, or it ran in another folder |
Step 10: The failure tests
Now prove the pipeline is strong. Reset first with python setup_topics.py --reset.
Test 1: Stop one loader
Start two loaders, the aggregator and both producers. While data flows, press Ctrl+C in loader 2. In loader 1 you will see the partitions of loader 2 come over. The log shows only the new partitions, for example:
[loader-1] assigned: [0, 2]
Notice something nice. Loader 1 did not stop and restart. With cooperative rebalancing, it kept its own partitions and only received the extra ones.
Test 2: Crash one loader
Do the same, but stop the loader with force: kill -9 <process id>. Start it again.
Do not worry if it waits. Kafka does not know the loader is gone until its session times out. The default is 45 seconds. After that, the group rebalances and work continues.
When all is quiet, run verify_results.py. It must pass. Look at the summary lines of the loaders. You will probably see duplicate counts that are higher than before. Those are messages that came a second time after the crash. Your event_id check ignored them. Nothing was lost, and nothing was counted twice.
This is the safe order paying for itself.
Test 3: Stop a broker
Start the producers with a slow rate, so you have time:
python order_producer.py --rate 5
While it runs, stop one of the three broker containers, the same way you did in the failure labs of the course. The producer keeps going. It may print a short delay while Kafka picks a new leader for some partitions. When it ends, check that it reports no failed messages. Start the broker again.
Why does this work? Replication factor 3, min.insync.replicas=2, acks=all and an idempotent producer. Two brokers are enough to accept a write. The producer retries safely.
Step 11: The replay tool
What: replay_sales.py. It rebuilds sales for a time window, from Kafka.
11a. Read the times
pythondef parse_time(text): moment = datetime.strptime(text, "%Y-%m-%d %H:%M").replace(tzinfo=timezone.utc) return int(moment.timestamp() * 1000)
11b. A consumer that touches nothing
pythonconsumer = Consumer({ "bootstrap.servers": BOOTSTRAP, "group.id": "sales-replay-" + uuid.uuid4().hex[:8], # a throw-away name "enable.auto.commit": False, "enable.partition.eof": True, })
A random group name and no commits. The real consumers are never disturbed.
11c. Find where to start
pythonasked = [TopicPartition(ORDERS_TOPIC, p, start_ms) for p in partition_ids] found = consumer.offsets_for_times(asked, timeout=10)
Here the third argument of TopicPartition is a time, not an offset. In the answer, the same field holds the offset that Kafka found. An offset below 0 means: nothing at or after this time. Skip that partition.
11d. Read the window
pythonconsumer.assign(list(active.values()))
assign() points the consumer straight at these partitions and offsets. There is no group and no rebalance.
Inside the loop, there are two reasons to stop reading a partition:
pythonif msg.error().code() == KafkaError._PARTITION_EOF: done.add(msg.partition()) # no more messages in this partition ... if msg.timestamp()[1] >= end_ms: consumer.pause([TopicPartition(ORDERS_TOPIC, msg.partition())]) done.add(msg.partition()) # we passed the end of the window
The loop ends when every partition is done. The counting rules are the same as the aggregator's.
11e. Compare
pythonsaved = zone_totals(open_sales_db(args.db), minute_of(start_ms), minute_of(end_ms)) print("\nMATCH" if saved == replayed else "\nDIFFERENT: sales.db does not agree with Kafka")
Check it: note the time when you started the producers, in UTC. For example, if you started at 14:05 UTC and finished by 14:08:
python replay_sales.py --from "2026-11-06 14:05" --to "2026-11-06 14:08"
Read 77 messages between 2026-11-06 14:05 and 2026-11-06 14:08 (UTC)
Rebuilt from Kafka
ZONE ORDERS REVENUE (INR)
...
Saved in sales.db for the same window
...
MATCH
Use a window that covers a time when your producers were running. If your producers sent everything in one minute, use that minute.
If it says DIFFERENT, you have found a bug in the aggregator, or in the replay tool. Fix it. This is exactly what the tool is for.
Step 12: Clean up
When you are done:
- Stop the consumers with Ctrl+C.
- Delete
orders.dbandsales.dbif you want a fresh start.python setup_topics.py --resetdoes this for you.
What you practised
| Skill from the course | Where you used it |
|---|---|
| Keys | order_id as key |
| Custom partitioning | One zone, one partition |
Idempotent producer, acks, min.insync.replicas | Safe delivery |
linger.ms, batch.size, compression | Producer tuning |
Callbacks, poll(), flush(), BufferError | The send loop |
| Consumer groups | Two groups, two loaders |
| Cooperative rebalancing | Stopping and starting loaders |
enable.auto.offset.store, store_offsets() | Save first, mark after |
| Dead-letter topic | Bad messages with headers |
Replay with offsets_for_times() and assign() | The replay tool |
Common problems
| Problem | What to try |
|---|---|
Topic shop.orders was not found | Run python setup_topics.py |
| A loader prints nothing after it starts | Wait a few seconds. Joining a group takes a moment. Is another loader already running? |
| A consumer seems stuck after a crash | Wait up to 45 seconds. Kafka must notice the old member is gone |
| The numbers are doubled | You ran the producers twice. Run python setup_topics.py --reset and start again |
database is locked | Another program is holding the database file open. Close it and run again |
A warning line about Failed to acquire idempotence PID ... Coordinator load in progress | Ignore it. It is a warning. The program retries by itself and continues. It usually appears right after the cluster starts |
Many Connection refused lines while you stop a broker | Expected. The library keeps trying until the broker comes back |
| The replay shows 0 messages | Check the time window. Times are UTC, and --to is not included |
Well done
You built a small but real data pipeline. It takes events from two sources, keeps them in order where it matters, survives crashes, parks bad data safely, and can look back in time.
The same ideas run in large companies every day.