Distributed systems
Messages processed out of order across partitions
Written and reviewed by Sahil Srivastav
# two updates for the same user landed on different partitions
partition=2 offset=884471 user=1042 event=address_changed city="Pune" event_ts=2026-09-14T11:02:17.201Z
partition=5 offset=119188 user=1042 event=address_changed city="Mumbai" event_ts=2026-09-14T11:02:17.455Z
# consumer for partition 5 was caught up; consumer for partition 2 was 40s behind
11:02:17.480Z indexer applied user=1042 city="Mumbai"
11:02:57.902Z indexer applied user=1042 city="Pune" <- older event applied last
$ curl -s localhost:9200/users/_doc/1042 | jq ._source.city
"Pune"What this error actually means
Kafka guarantees order within a partition and nothing at all across partitions. Two events for the same entity that land on different partitions are consumed by different consumers, at different speeds, with no coordination between them. Whichever consumer happens to be less busy wins the race, and the last write to the database is not necessarily the newest event.
The result is a corrupted current state with a perfectly intact event log — which is why this is so hard to spot. Replaying the log in order produces the right answer, so nobody believes the pipeline is broken. Only the materialised view, search index, or cache is wrong, and only for the entities that happened to have two rapid updates.
The usual reason two events for one entity split across partitions is the partition key. No key at all means round-robin, so consecutive events for the same user land anywhere. A key chosen at the wrong granularity — `region` when you need `user_id`, or `user_id` when a transfer touches two accounts — has the same effect. And changing a topic’s partition count rehashes every key, so events published before the change sit on a different partition from events published after.
Note that per-partition order is necessary but not sufficient. A consumer that receives an ordered batch and then hands records to a thread pool has thrown the ordering away inside its own process, and this is a very common way to reintroduce the bug while believing the keys protect you.
Causes, most common first
- 1No partition key, or a key coarser than the entity. Producing with a null key distributes records across all partitions. Choosing `tenant_id` when correctness depends on `account_id` has the same effect at a smaller scale. The ordering guarantee you need is always at the granularity of the thing you mutate.
- 2Partition count changed, so the key-to-partition mapping changed. The default partitioner is a hash modulo the partition count. Adding partitions moves most keys, and for a window equal to your retention, the same key exists on both the old and the new partition. Increasing partitions is therefore an ordering event, not just a capacity change.
- 3Concurrent processing inside the consumer. The records arrive ordered and the handler dispatches them to an executor, or uses a reactive pipeline with `flatMap`, or processes a batch in parallel. The broker kept its promise; the consumer broke it. Also produced by any per-record retry that lets a later record overtake an earlier one.
- 4Two producers writing the same entity to different topics. A profile service and an admin tool both emit user updates, on separate topics with separate consumers. There is no ordering relationship between two topics at all, so the outcome depends entirely on relative consumer lag.
- 5Writes applied as blind overwrites. The deeper cause behind all of the above: a handler that does `UPDATE ... SET city = $1` with no notion of which version of the entity it is writing. It cannot tell a new value from an old one, so it cannot refuse a stale one.
When you see it
- A field in a read model reverts to an older value and stays there until the next update
- Only entities with two updates within a few seconds are affected; the rest look perfect
- Replaying the topic from the beginning produces the correct state, so the bug "cannot be reproduced"
- A delete followed by a recreate yields a record that exists with the pre-delete data
- The incidence jumped after a partition-count increase or after the producer key changed
- One consumer instance has visible lag while others are caught up
How to diagnose it
Step 1
Check whether the entity is spread across partitions
This is the decisive test. Consume the topic and group the events for one affected entity by partition. More than one partition for a single entity key proves the producer side; a single partition points you at the consumer.
kafka-console-consumer.sh --bootstrap-server broker:9092 --topic user-events \
--from-beginning --property print.partition=true --property print.key=true \
| grep '"user_id":1042'Step 2
Compare event timestamp with apply timestamp
Log both the producer-assigned event time and the time you applied the write. Any row where a later apply carries an earlier event time is an inversion, and counting them tells you the real exposure rather than the reported exposure.
SELECT id, last_event_ts, updated_at
FROM user_search_index
WHERE last_event_ts < (
SELECT max(event_ts) FROM user_events e WHERE e.user_id = user_search_index.id
);Step 3
Look for a partition-count change in the topic history
If the topic was expanded, the timing of that change will line up with the onset of the symptom. Confirm the current count, then check whether affected keys have events on two partitions spanning that moment.
kafka-topics.sh --bootstrap-server broker:9092 --describe --topic user-eventsStep 4
Rule out the consumer reordering an ordered stream
Grep the handler for thread pools, `CompletableFuture`, `parallelStream`, `flatMap`, and per-record retry loops. If any of them sit between `poll()` and the write, the consumer is the reordering agent regardless of keys.
The fix
Key the producer by the entity whose order you must preserve, so all events for that entity are serialised into one partition. This is the cheap structural fix and it covers the majority of cases — but it buys ordering at the price of a hot partition if one key is disproportionately busy, so check key distribution before committing to it.
Make every write reject stale data with a version guard, so ordering stops being load-bearing. Carry a monotonic version per entity — a database sequence, an LSN, an event sequence number — and apply the update only when the incoming version is greater than the stored one. This is the fix that survives partition changes, consumer concurrency, and multi-topic producers alike.
Do not use wall-clock timestamps as the version unless they come from a single clock. Producer clocks disagree by tens of milliseconds at best and minutes at worst, and two events within that window will be ordered arbitrarily. Where a timestamp is the only option, source it from the database on the write path rather than from the producer host.
If you must process concurrently for throughput, partition the concurrency by the same key you partitioned the topic with: hash the entity id to a fixed set of single-threaded workers. Order is preserved per entity while unrelated entities proceed in parallel.
Expand partition counts deliberately. Either accept a bounded window of disorder and rely on the version guard to absorb it, or drain the topic before the change, or create a new topic with the new partition count and cut consumers over. Adding partitions to a live keyed topic with no version guard reorders production data.
// Blind overwrite: the last event to be applied wins, whichever one that is.
jdbc.update("UPDATE user_index SET city = ? WHERE user_id = ?",
evt.city(), evt.userId());
// Version guard: an older event cannot overwrite a newer one, regardless
// of which partition it arrived on or how many threads are consuming.
int applied = jdbc.update("""
UPDATE user_index
SET city = ?, version = ?, updated_at = now()
WHERE user_id = ?
AND version < ?
""", evt.city(), evt.version(), evt.userId(), evt.version());
if (applied == 0) {
metrics.counter("index.stale_event_dropped").increment();
}
-- for the insert-or-update case, keep the guard in the conflict clause
INSERT INTO user_index (user_id, city, version, updated_at)
VALUES ($1, $2, $3, now())
ON CONFLICT (user_id) DO UPDATE
SET city = EXCLUDED.city,
version = EXCLUDED.version,
updated_at = now()
WHERE user_index.version < EXCLUDED.version;How to stop it coming back
- Record the ordering key and the version field for every event type in a schema registry or contract test, so a new producer cannot omit them
- Count dropped stale events as a metric — a sudden rise means disorder increased, which is an early warning that a key or partition count changed
- Treat a partition-count increase as a migration with a written plan, not as a capacity dial
- Ban unkeyed produces to topics that feed materialised state; a lint rule or a producer wrapper is enough
- When adding concurrency to a consumer, make the key-affinity requirement explicit in the code rather than a comment — a plain executor is the default and it is wrong here
FAQ
Does Kafka guarantee ordering within a topic?
No — only within a partition. A single-partition topic gives you total order and caps your throughput at one consumer; every other configuration gives you per-key order at best, and only if the producer supplies the key.
Is max.in.flight.requests.per.connection relevant here?
Yes, on the producer side. With retries enabled and more than one request in flight, a retried batch can land after a later batch on the same partition. Producer idempotence (`enable.idempotence=true`) preserves order for up to five in-flight requests; without it, set the limit to one if order matters.
Can I sort the events in the consumer with a small buffer?
Only with a time window and a definition of how late an event may be, which means holding writes back and still being wrong for anything later than the window. It is a real technique in stream processing, but it is far more machinery than a version column, and it needs comparable clocks to work at all.
Why did this start after we added partitions?
Because partitioning is hash modulo partition count. Changing the count remaps most keys, so for as long as both the old and new events remain in retention, the same entity has events on two partitions being consumed independently. The disorder is transient but the corrupted state it writes is permanent.
Related
Other errors engineers hit next to this one
- Too many open files
- Out of memory: Killed process
- No space left on device despite free disk space
- Text file busy during executable replacement
- set -e script continues after a failed pipeline
- An unquoted variable turns one argument into several
- OOMKilled — container exit code 137
- CrashLoopBackOff