The order arrived twice, and once it vanished
Commit timing makes one vanish — rewinding and idempotent processing
This lab runs on a VM
Apache Kafka 4.3.1 runs as a single KRaft node on an Ubuntu VM (at localhost:9092). The helper ship-process imitates a shipping processor: it reads all the orders on standard input and then writes them one at a time to /root/kafka/processed.txt, and unless SHIP_FIXED=1 is set, it dies at order-2. It takes about 4 minutes to come up the first time.
Goal
You reproduce messages "vanishing" after a crash when a consumer commits its offset before processing, confirm with another group that those messages are still in Kafka, and then rewind the group offsets and process them again. You block the duplicates created by the rewind with an idempotent consumer, and see how auto.offset.reset decides the first position of a new group.
Why it matters
The 'message delivery semantics' section of the design document is the script for this lab. If you read, then save the position, then process, a message whose processing died midway never comes again (at-most-once). If you read, then process, then save, it comes again when the process died before the save (at-least-once). A processor that relies on auto-commit (enable.auto.commit defaults to true, with a 5-second interval), like the console consumer, is closer to the first shape. "It vanished once" is usually this, and fixing it turns it into "arrived twice," so the processor has to be idempotent. The document describes this as "the message has a primary key, so the update is idempotent."
Steps
- Create the topic
shipmentswith 1 partition and put in five lines, fromorder-1toorder-5. - With the group
ship-svc, read three events from the beginning and pipe them toship-process(it crashes). Then describeship-svcand save it to/root/kafka/ship-crash.txt. CURRENT-OFFSET should be 3, but/root/kafka/processed.txtshould contain onlyorder-1;order-2andorder-3are the ones that "vanished." - With a new group
audit-svc, read five events from the beginning and save them to/root/kafka/ship-audit.txt. They are all still in Kafka. - Rewind the offsets of
ship-svcto the very beginning (--reset-offsets --to-earliest --execute) and save that output to/root/kafka/ship-reset.txt. - With
SHIP_FIXED=1, haveship-svcread the five events again and pipe them toship-process.processed.txtshould end up with six lines, withorder-1appearing twice: the duplicate that is the price of rewinding. - Create
/root/kafka/dedup.sh. It reads the orders on standard input, and only those that are not in/root/kafka/seen.txtare written to/root/kafka/processed-dedup.txtand recorded in seen (if the environment variablesDEDUP_SEENandDEDUP_OUTare set, it uses those paths). Even if you readshipmentsfrom the beginning twice and pipe it, only five lines should remain. - Without
--from-beginning, read for 5 seconds with a new grouplate-svc(0 events), and withauto.offset.reset=earliest, read five events with a new groupearly-svc. In/root/kafka/offset-reset.txt, write two lines:late_count=0andearly_count=5. - In
/root/kafka/consumer-report.txt, write three lines:lost_after_crash=<2단계에서 사라진 건수>,duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>, andunique_orders=<processed-dedup.txt 의 줄 수>. (These are the number lost in step 2, the number of duplicate entries in processed.txt after step 5, and the number of lines in processed-dedup.txt.)
Notes
- Reading with a group:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic shipments --group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process - Rewinding:
kafka-consumer-groups.sh ... --reset-offsets --group ship-svc --topic shipments --to-earliest --execute. As the operations document says, the consumer must be stopped. If you run it without--execute, it only shows the plan. - A new group's first position:
auto.offset.resetin the consumer configuration document; the defaultlatestmeans that if the group has no offset it reads only what comes after now. Change it with--command-property auto.offset.reset=earliest.--from-beginningis a shortcut by which the console tool does the same thing. - Common mistake 1: leaving out
--max-messagesin step 2. If it reads all five events, the "two vanished events" are not reproduced. - Common mistake 2: keeping dedup's state (seen) only in memory. When the process dies the state dies too, and the next run produces duplicates again. Whether it is a file or a DB, it has to remain together with the processing result.
Five shipping orders
Create the topic shipments with 1 partition and put in five lines, from order-1 to order-5.
printf 'order-1\norder-2\norder-3\norder-4\norder-5\n' | kafka-console-producer.sh .... There is one partition, so the order is fully preserved.
If you commit before processing, it vanishes
With the group ship-svc, read three events from the beginning and pipe them to ship-process (it crashes). Then describe ship-svc and save it to /root/kafka/ship-crash.txt. CURRENT-OFFSET should be 3, but /root/kafka/processed.txt should contain only order-1; order-2 and order-3 are the ones that "vanished."
--group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process. The console consumer hands over three events, commits offset 3, and finishes, but the processor died on the second. The next time you read with this group, it starts from 3.
It is still in Kafka
With a new group audit-svc, read five events from the beginning and save them to /root/kafka/ship-audit.txt. They are all still in Kafka.
--group audit-svc --from-beginning --max-messages 5 --timeout-ms 8000 > /root/kafka/ship-audit.txt. Consuming is not deleting: what vanished is not the messages but the position of ship-svc.
Rewind the group offsets
Rewind the offsets of ship-svc to the very beginning (--reset-offsets --to-earliest --execute) and save that output to /root/kafka/ship-reset.txt.
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --group ship-svc --topic shipments --to-earliest --execute. This is possible only because the consumer's position is a single integer; the design document calls it a feature that violates the queue contract but is indispensable.
If you process again, it arrives twice
With SHIP_FIXED=1, have ship-svc read the five events again and pipe them to ship-process. processed.txt should end up with six lines, with order-1 appearing twice: the duplicate that is the price of rewinding.
... --group ship-svc --max-messages 5 --timeout-ms 8000 | SHIP_FIXED=1 ship-process. If the group has an offset, --from-beginning is ignored: it reads from the rewound position (0). The two events that vanished come back, but the already processed order-1 also comes again.
The idempotent consumer
Create /root/kafka/dedup.sh. It reads the orders on standard input, and only those that are not in /root/kafka/seen.txt are written to /root/kafka/processed-dedup.txt and recorded in seen (if the environment variables DEDUP_SEEN and DEDUP_OUT are set, it uses those paths). Even if you read shipments from the beginning twice and pipe it, only five lines should remain.
Check whether you have seen it before with grep -qxF "$line" "$SEEN", and write to the two files only when you have not. Read twice with --from-beginning --max-messages 5 with no group and pipe it. The grader passes temporary paths through the environment variables and feeds it input mixed with duplicates.
Where does a new group start?
Without --from-beginning, read for 5 seconds with a new group late-svc (0 events), and with auto.offset.reset=earliest, read five events with a new group early-svc. In /root/kafka/offset-reset.txt, write two lines: late_count=0 and early_count=5.
--group late-svc --timeout-ms 5000 finishes without reading anything (the default is latest). --group early-svc --command-property auto.offset.reset=earliest --max-messages 5 --timeout-ms 8000 reads five events. Count with wc -l and write the numbers to the file.
Count what vanished and what came twice
In /root/kafka/consumer-report.txt, write three lines: lost_after_crash=<2단계에서 사라진 건수>, duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>, and unique_orders=<processed-dedup.txt 의 줄 수>. (These are the number lost in step 2, the number of duplicate entries in processed.txt after step 5, and the number of lines in processed-dedup.txt.)
The number lost is CURRENT-OFFSET in ship-crash.txt minus the number processed at that time (1), and the number of duplicates is the number of lines in processed.txt minus the number of distinct orders. The grader recounts from the same files.