Relying solely on Kafka’s ordering within a partition is not enough when thousands of sessions share the same partition. In conversational and sequential-task systems, a corrective message such as “make the destination Paris” may arrive after “book a flight to London,” but executing it before the first message corrupts the context. The solution described by Joshua Oluikpe in InfoQ moves part of the ordering responsibility to the Go-based application layer.
Ordering Within a Session and Parallelism Between Sessions
The design relies on two levels of workers. The distribution layer receives Kafka records and routes them according to the session ID using consistent hashing, so messages from the same session reach the same worker, while different sessions can be distributed across multiple workers. Each active session then gets a dedicated goroutine that processes one message at a time in arrival order.
This prevents a slow session from blocking other sessions, as can happen in a flat worker pool when retrying a single session occupies the worker entirely. Goroutines are created as needed and removed after a period of inactivity, making resource consumption depend on the number of active sessions rather than the total number of possible sessions.
Retrying Without Skipping Messages
Retries take place within the session’s own goroutine. When a temporary error occurs, such as a downstream service timeout, the worker waits using exponential backoff with some jitter and then processes the same message again before moving on to subsequent messages. Because later messages remain in the channel behind the failed message, the fourth message cannot overtake the third.
Non-recoverable errors, such as invalid data, go directly to the dead-letter queue (DLQ). This division between temporary and terminal errors eliminates the need for a separate state machine to manage retries, but it also means that a particular session’s messages may accumulate during the backoff period.
Safe Committing and Recovery After a Crash
When multiple sessions process their messages in parallel, it is not safe to commit only the highest completed offset. If offset 104 completes before 102 and the system commits 104, a consumer crash could cause 102 to be skipped permanently. The system therefore tracks, for each partition, offsets that are in progress and those that have completed, and advances the commit point only to the highest consecutively completed offset.
As a result, a restart may begin at an earlier offset and reprocess some messages. The design uses a stable event ID to help eliminate duplicates, but it does not claim to achieve “exactly once” execution; external side effects still depend on downstream services being able to handle duplicates in an idempotent manner. In rare cases, a stuck gap can be skipped after a specified timeout to restore partition progress, with the event recorded as an operational failure and the message routed to the DLQ. This is an explicit trade-off between continuity and processing completeness.
What Is Required for Actual Production?
Ordering guarantees are not complete with the core logic alone. During partition rebalancing, the consumer stops receiving new records, waits for in-flight work to finish, and then commits the latest contiguous watermark before transferring ownership. When worker channels become full, it temporarily stops fetching messages instead of dropping them, while records already fetched are stored in per-partition buffers so newer messages cannot move ahead of them.
The system also needs to detect stuck offsets individually, propagate tracing across asynchronous Kafka boundaries, retain the lag signals needed for autoscaling while fetching is paused, and drain work and write to the DLQ before shutdown. The producer also applies consistent hashing on the sending side to ensure message ordering during asynchronous production.
Results and Limitations
In a synthetic test involving 50,000 messages and 10 sessions, with transient errors injected into 10% of the messages, the system recorded 14,027 messages per second with no ordering violations and a median latency of 11.7 milliseconds. In a load test with 1,000 sessions, the effective send rate reached 48,805 messages per second out of a target of 50,000, with no send errors. However, tail latency exceeded the 200-millisecond budget because of the message-rate quota in the API Gateway, not because of the session-ordering mechanism.
The author says that the system processed more than 40 million production messages without noticeable ordering violations, with only 460 messages routed to the DLQ, or approximately 0.002%. However, this figure is based on operational monitoring rather than a comprehensive check of every message sequence, and the higher load test did not independently verify ordering for every session. The design’s practical value therefore lies in turning ordering from a Kafka-level assumption into a property enforced by the application, while retaining the need for explicit sequence tests and idempotency guarantees in downstream services.