TT Lab
Get started
Learn Learning paths Courses

The order arrived twice, and once it vanished

What Kafka promises — order within a partition and a single-integer position

Continue in TT Lab

One-line summary

A Kafka topic is an append-only log split into partitions; events with the same key pile up in order in the same partition; and a consumer group's "how far have I read" is a single integer per partition (the offset). Consuming does not delete anything; only the retention settings delete. These four sentences are the key to reading the "twice" and "vanished" cases later on.

Why this was needed

People who have used queues mistake Kafka for a queue. In a queue, a message disappears when you take it out, and the broker remembers "who received what" for each message. The 'Consumer Position' section of the design document lists the costs of that approach: if you mark a message as consumed the moment you hand it over, the message is lost when the consumer dies mid-processing; if you wait for an acknowledgment (ack), it gets consumed twice when the consumer processed it but failed to send the ack; and the broker has to hold several pieces of state for every message.

Kafka took a different path. If you split a topic into totally ordered partitions and let exactly one consumer in a consumer group read each partition, the consumer's position becomes a single integer, "the next offset to read." Periodically checkpointing that integer is the whole acknowledgment, which makes it very cheap, and a by-product is rewinding: if there was a bug in your code, you can go back to an old offset and read again. The document notes that this violates the queue contract but is an essential feature for many consumers.

How it works

Events and topics. As the introduction describes, an event has a key, a value, and a timestamp (and optional headers), and a topic is where those events pile up. A topic accepts many producers and many consumers at the same time, and events are not deleted after being consumed. How long they stay is decided by per-topic retention settings (module 4).

Partitions and keys. A topic is divided into several partitions ("buckets") placed across several brokers. A new event is appended to the end of one of them, and events with the same key (for example, an order number) go to the same partition. The partitioner.class entry in the producer configuration document describes the default behavior: if there is a key, the partition is chosen by the hash of the key; if there is no key, it is a sticky approach that sticks to one partition until batch.size is filled. So a small number of events sent without a key pile up in one partition, and if you give a key, each order lines up in a single row.

The only ordering Kafka promises is within a partition. An order's created → paid → shipped is read in that order because it is in the same partition, but across different orders there is no fixed answer to which partition is read first. That is why the unit that needs ordering must be the key.

orders (3 파티션)              컨슈머 그룹 order-svc 의 위치
 P0: order-2 created, order-3 created, order-2 paid   → 오프셋 3
 P1: order-1 created, order-1 paid, order-1 shipped   → 오프셋 3
 P2: order-4 created, order-5 created                 → 오프셋 2

Consumer groups and offsets. A group is a set of consumers tied together by group.id, and the group's position is the offset committed to the broker for each partition. Different groups have different positions: even if order-svc has read to the end, analytics reads again from the beginning. The difference between the log-end offset and the committed offset is the lag, and it is the first number you look at in operations. kafka-consumer-groups.sh --describe shows CURRENT-OFFSET, LOG-END-OFFSET, and LAG for each partition.

The number of partitions can only be increased. The 'Modifying topics' section of the operations document lists three side effects of increasing partitions. Because data is divided by hash(key) % 파티션 수 (the hash of the key modulo the number of partitions), when the number changes the same key can go to a different partition and existing data is not redistributed (so key ordering guarantees can break); an existing consumer with auto.offset.reset=latest can miss messages that arrived before it noticed the new partition; and there is a delay in metadata propagation. And decreasing is not supported.

What it looks like in practice

A common cause of the report "order statuses look flipped" is the key. The order service was sending without a key, or used the event type rather than the order number as the key, or, after the partition count was increased, the old and new events of the same order ended up split across different partitions. In all three cases the broker is healthy and the log is healthy; someone expected something that was never promised.

If the lag graph suddenly climbs, it is one of two things: producers are sending a lot, or the consumer has stopped. If you describe the group and the CONSUMER-ID column is empty, it is the latter. Even if a consumer dies, the committed offset remains on the broker, so when it comes back up it reads from that position; what arrives twice and what disappears at that point is module 3.

What you will do in the next lab

You put keyed events into a topic with 3 partitions and see that the same key piles up in order in the same partition, read the group's offsets and lag, confirm that two groups are independent, and then increase the partitions to 6 and observe for yourself how the key placement changes.