Per-Entity Ordering with Kafka Partition Keys

This walkthrough applies the lane model from Message Ordering Guarantees to Kafka, within Queue Fundamentals & Architecture. Kafka orders records per partition and assigns each partition to one consumer, so the whole job is choosing a key, keeping the producer from reordering on retry, and keeping the consumer from reordering by being clever with threads.

Problem Statement

An inventory service emits StockReserved, StockReleased, and StockAdjusted events per SKU to a topic with twelve partitions, and a projection worker maintains available stock. The projection occasionally goes negative. Investigation shows two causes: the producer was sending without a key, so events for one SKU were spread round-robin across partitions; and after keys were added, a well-meaning change made the consumer process each poll batch in a thread pool of 16. Both reorder events for the same SKU. You need per-SKU order end to end, while using more than twelve threads of processing capacity.

Prerequisites

  • Kafka 3.x brokers (or a compatible service such as MSK or Confluent Cloud) and a topic you can recreate or whose partition count you can plan.
  • A producer library that exposes enable.idempotence (the Java client, confluent-kafka for Python and Go, kafkajs with idempotent: true).
  • A stable entity id in every event (sku), and handlers that are idempotent, since rebalances cause redelivery.

Step 1 — Key Every Record by the Entity

The default partitioner hashes the record key with murmur2 and takes it modulo the partition count, so every record with the same key lands on the same partition. A record with a null key is spread across partitions by the sticky partitioner, and ordering for that entity is gone.

# producer.py — confluent-kafka
from confluent_kafka import Producer
import json

producer = Producer({
    "bootstrap.servers": "kafka-1:9092,kafka-2:9092,kafka-3:9092",
    "enable.idempotence": True,           # broker drops duplicate retried batches by sequence
    "acks": "all",                        # required by idempotence
    "max.in.flight.requests.per.connection": 5,   # max allowed while keeping order with idempotence
    "linger.ms": 5,                       # small batching window for throughput
    "compression.type": "zstd",
})

def emit(event: dict) -> None:
    producer.produce(
        topic="inventory-events",
        key=event["sku"].encode(),        # the ordering key: one SKU, one partition
        value=json.dumps(event).encode(),
        on_delivery=_on_delivery,
    )
    producer.poll(0)

def _on_delivery(err, msg):
    if err is not None:
        # With idempotence on, the client already retried; a surfaced error is fatal
        # for ordering of this key. Stop emitting for the key and alert.
        raise RuntimeError(f"delivery failed for key {msg.key()}: {err}")

Idempotence is what makes retries safe for order. Without it, a batch that times out and is retried can land after a later batch that succeeded on the first try, swapping two records in the partition. With idempotence the broker tracks a per-producer sequence number and rejects out-of-sequence writes, so the client retries in order.

Key hash selects the partition Three SKUs are hashed by the default murmur2 partitioner. Every record for SKU A lands on partition 3, every record for SKU B on partition 7, and SKU C on partition 3 as well. A record with no key could land on any partition, breaking order for that SKU. murmur2(key) mod 12 key = SKU-A key = SKU-B key = null partition 3 partition 7 any partition always the same partition always the same partition order lost for this entity

Step 2 — Plan the Partition Count Once

Changing the partition count changes hash mod N, so existing keys move. For a short window after the change, new events for a SKU go to its new partition while older events still sit unconsumed on the old one, and the consumer of the new partition can overtake them. Choose the count for the parallelism you expect to need, not for today's load.

# Create the topic with headroom: 48 partitions, replication 3, min ISR 2
kafka-topics.sh --bootstrap-server kafka-1:9092 --create \
  --topic inventory-events \
  --partitions 48 \
  --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config retention.ms=604800000      # 7 days of replay for rebuilding projections

A rule of thumb: partitions ≥ peak consumer instances × 2, so you can double the fleet without repartitioning. If you truly must grow a keyed topic, create a new topic with the larger count, dual-write, let consumers finish the old topic, then switch — a migration, not a config change. Partition sizing in general is covered under queue partitioning strategies.

Step 3 — Consume with One Worker per Key, Not per Partition

Twelve or 48 partitions may be fewer lanes than you want threads. The safe way to add concurrency inside a consumer is to dispatch by key, not by record: records with the same key go to the same worker queue, in offset order, while different keys run in parallel.

# consumer.py — per-key parallelism inside one partition assignment
import hashlib
import queue
import threading
from confluent_kafka import Consumer

WORKERS = 16
lanes = [queue.Queue(maxsize=500) for _ in range(WORKERS)]    # bounded: backpressure
done_offsets: dict[tuple[str, int], int] = {}
lock = threading.Lock()

def lane_for(key: bytes) -> int:
    return int(hashlib.blake2b(key, digest_size=4).hexdigest(), 16) % WORKERS

def worker(q: queue.Queue) -> None:
    while True:
        msg = q.get()
        apply_projection(msg)                # sequential within this lane
        with lock:
            tp = (msg.topic(), msg.partition())
            done_offsets[tp] = max(done_offsets.get(tp, -1), msg.offset())
        q.task_done()

for q in lanes:
    threading.Thread(target=worker, args=(q,), daemon=True).start()

consumer = Consumer({
    "bootstrap.servers": "kafka-1:9092",
    "group.id": "inventory-projector",
    "enable.auto.commit": False,             # commit only what is safely processed
    "partition.assignment.strategy": "cooperative-sticky",
    "max.poll.interval.ms": 300000,
})
consumer.subscribe(["inventory-events"])

while True:
    msg = consumer.poll(1.0)
    if msg is None or msg.error():
        continue
    lanes[lane_for(msg.key())].put(msg)      # same key -> same lane -> same order

The done_offsets bookkeeping above is deliberately simplified: with several lanes consuming one partition, offset 105 can finish before offset 101. Committing 105 would mark 101 as processed. Step 4 fixes that.

Fan out by key inside a partition A partition delivers records for several SKUs in offset order. The consumer hashes each record's key to one of several in-memory lanes. Each lane is processed by one thread, so SKU A's records remain in order on lane 1 while SKU B's run concurrently on lane 2. More threads than partitions, order intact partition 3 A@100, B@101, A@102 C@103, B@104 lane 1: A@100, A@102 lane 2: B@101, B@104 lane 3: C@103 commit watermark highest offset with no unfinished offset below Lanes finish out of offset order, so the commit must track the contiguous finished prefix.

Step 4 — Commit Only the Contiguous Finished Prefix

Track, per partition, the set of offsets dispatched and finished. The safe commit point is the lowest unfinished offset: everything below it is done.

# commit_tracker.py
import heapq

class PartitionTracker:
    def __init__(self) -> None:
        self.in_flight: list[int] = []        # min-heap of dispatched-but-unfinished
        self.finished: set[int] = set()
        self.highest_dispatched = -1

    def dispatched(self, offset: int) -> None:
        heapq.heappush(self.in_flight, offset)
        self.highest_dispatched = max(self.highest_dispatched, offset)

    def finished_offset(self, offset: int) -> None:
        self.finished.add(offset)
        while self.in_flight and self.in_flight[0] in self.finished:
            self.finished.discard(heapq.heappop(self.in_flight))

    def committable(self) -> int | None:
        # Kafka commits the NEXT offset to read
        if self.in_flight:
            return self.in_flight[0]
        return self.highest_dispatched + 1 if self.highest_dispatched >= 0 else None

Commit every few seconds from the poll loop with consumer.commit(offsets=[TopicPartition(t, p, tracker.committable())], asynchronous=False). If the process dies, it replays from the watermark — some records are processed twice, which the idempotent handler absorbs, but none are skipped. Libraries such as Confluent's Parallel Consumer implement exactly this key-ordered mode if you would rather not maintain it.

Step 5 — Survive Rebalances Without Reordering

When a partition is revoked, in-flight lanes may still hold its records. If the new owner starts before they finish, two processes apply the same SKU's events concurrently. Drain before releasing:

def on_revoke(consumer, partitions):
    for tp in partitions:
        wait_for_lanes_to_finish(tp)            # block until this partition's records are applied
        off = trackers[(tp.topic, tp.partition)].committable()
        if off is not None:
            consumer.commit(offsets=[TopicPartition(tp.topic, tp.partition, off)],
                            asynchronous=False)
        trackers.pop((tp.topic, tp.partition), None)

consumer.subscribe(["inventory-events"], on_revoke=on_revoke)

With the cooperative-sticky assignor only the partitions that actually move are revoked, so the rest of the consumer keeps working through a rebalance. Keep max.poll.interval.ms comfortably above the time to drain the bounded lanes, or the drain itself triggers another rebalance — the loop described in rebalancing Kafka consumer groups without stalls.

Step 6 — Measure Per-Partition Lag and Key Skew

Ordering converts a throughput problem into a distribution problem. Total consumer lag can look fine while one partition, carrying a handful of hot SKUs, falls minutes behind — and every event for those SKUs waits with it. Two measurements catch this early: lag per partition, and the share of production landing on the busiest partition.

# Seconds of lag on the worst partition (kafka-lag-exporter / Burrow style metric)
max by (partition) (
  kafka_consumergroup_group_lag_seconds{group="inventory-projector", topic="inventory-events"}
)

# Skew ratio: busiest partition's produce rate divided by the average partition's.
# 1.0 is perfectly even; above ~3 a few keys dominate.
max(rate(kafka_log_log_logendoffset{topic="inventory-events"}[15m]))
  /
avg(rate(kafka_log_log_logendoffset{topic="inventory-events"}[15m]))

Alert on the per-partition number, not the sum. A threshold such as "any partition more than 120 seconds behind for 10 minutes" pages on the stuck lane that a total-lag alert averages away.

When skew is the problem, adding consumers does nothing — one partition is still consumed by one member. The remedies are on the key side: make the key finer where business rules allow (warehouse:sku instead of sku if stock is tracked per warehouse), or route a known hot entity to a dedicated topic with its own consumer. Per-key dispatch inside the consumer (Step 3) helps only if the hot partition carries several busy keys; a single hot key is sequential by definition.

One hot partition sets the lag Eight partitions are shown with their produce rates. Seven carry between 180 and 260 records per second. Partition 5 carries 900 records per second because of a few very active SKUs, so its lag grows while the others stay near zero, and total lag hides the problem. Records per second by partition p0 p1 p2 p3 p4 p5 p6 p7 900/s: three hot SKUs 180 to 260/s each

Verification

Replay a synthetic stream with interleaved keys and random handler delays, then assert per-key order and no gaps:

# verify_order.py — run against a staging topic
from collections import defaultdict
seen = defaultdict(list)

def apply_projection(msg):
    time.sleep(random.uniform(0, 0.02))
    ev = json.loads(msg.value())
    seen[ev["sku"]].append(ev["seq"])

# after draining 100k events over 2,000 SKUs:
for sku, seqs in seen.items():
    collapsed = [s for i, s in enumerate(seqs) if i == 0 or s != seqs[i - 1]]
    assert collapsed == list(range(1, len(collapsed) + 1)), sku

In production, track lag per partition with kafka_consumergroup_lag and a skew ratio (busiest partition's produce rate over the total). A skew far above 1/48 means a few hot SKUs dominate one partition.

Gotchas & Edge Cases

Custom partitioners across languages. The Java client and librdkafka use different default hash functions unless configured (partitioner: murmur2_random in librdkafka matches Java). A service written in Go and one in Java producing to the same topic with defaults can send the same key to different partitions.

Transactions do not add ordering. Kafka transactions give atomic writes across partitions and exactly-once read-process-write; they do not order records across partitions. See Kafka exactly-once semantics with transactions for what they do provide.

Retrying inside a lane. A handler that fails should retry in place with bounded backoff. Publishing the record to a retry topic moves it behind later records for the same key; that is acceptable only if the projection is version-checked.

Tombstones and compaction. On a compacted topic, older records for a key disappear. A consumer that rebuilds from the beginning sees only the latest value per key — ordered, but not the full history.

FAQ

Is one partition per entity a good idea? No. Partitions are a broker resource with per-partition overhead in memory, file handles, and leader elections. Many entities share each partition; the hash keeps each entity on one of them.

Does Kafka guarantee order across partitions? No. There is no global order on a multi-partition topic. If you need total order, use a single partition and accept a single consumer's throughput.

Can I use a compound key such as tenant plus order id? Yes, and it is often the right choice: f"{tenant}:{order_id}" orders per order while spreading a large tenant across many partitions. Only use a coarser key when events for different orders in the same tenant genuinely depend on each other.

Related