Key Takeaways
- Kafka partition ordering alone is not enough for session-ordered workloads. When many sessions share partitions, the application layer still needs routing logic to preserve ordering within each session while allowing parallelism across sessions.
- Consistent hashing and per-session workers provide session affinity. Messages from the same session always land on the same worker and are processed sequentially, while unrelated sessions continue in parallel.
- Retries preserve ordering when they happen in-place. A per-session goroutine can retry a failed message with backoff before processing the next message, avoiding message overtaking and removing the need for a separate retry state machine.
- Contiguous watermark commits make concurrent processing recoverable. Offsets are committed only up to the highest contiguous completed offset, preventing a later completed message from causing an earlier unfinished message to be skipped after a crash.
- Production ordering guarantees require operational hardening. Rebalance draining, backpressure, pending-record buffering, stuck-offset detection, idempotency, graceful shutdown, and dead letter queue handling are all needed to preserve ordering under real failure conditions.
Introduction
In conversational AI, message ordering carries semantic weight. When a user sends "Book me a flight to London" followed by "Actually, make it Paris," the system must process those messages in sequence. If the second message overtakes the first, the AI responds to the wrong context. The user gets a London booking confirmation instead of a Paris redirect. At scale, with thousands of concurrent sessions flowing through a multi-stage pipeline, even rare ordering violations degrade the product.
I led the design and implementation of the core Kafka messaging layer in Go for a real-time conversational AI platform serving enterprise customers. Working closely with my team, we built a pipeline where every message flows through four processing stages: ingest, natural language understanding (NLU), large language model (LLM) inference, and delivery. The ordering requirement is strict: messages within a session are processed in sequence across all four stages, regardless of how long an individual stage takes.
The system has processed over 40 million messages in production with no observed ordering violations, across thousands of concurrent chat sessions. The ordering guarantee here does not come from staying under some safe throughput. It is structural, built from consistent hashing and one goroutine per session. We have since pushed send-side throughput past 100,000 messages/second in testing with zero send errors.
Architecture Overview
The core challenge is that Kafka guarantees ordering within a partition, but a single partition carries interleaved traffic from many sessions. To process sessions in parallel while preserving within-session order, the application layer needs its own routing.
Our first implementation used a flat, fixed worker pool with sessions hashed to workers. Under sustained load, one retrying session could block unrelated sessions on the same worker, while retry coordination introduced growing in-memory queues. We replaced it with a two-level worker hierarchy.

Figure 1. Two-level worker hierarchy for session-ordered processing (Image source: created by author)
The current design solves both problems with a two-level hierarchy.
Level 1: SessionWorkers (dispatch layer)
A pool of lightweight dispatch workers receives records from Kafka and routes them based on the session ID. Consistent hashing ensures that messages from the same session always reach the same worker, while different sessions can be distributed across workers for parallel processing.
Level 2: Per-session goroutines (processing layer)
Each active session is handled by its own goroutine, which processes messages one at a time in arrival order. These goroutines are created on demand and removed after a period of inactivity, keeping resource usage proportional to the number of active sessions.

Every session goroutine processes one message at a time. Parallelism across sessions is unlimited by design; serialization within a session is guaranteed by construction.
A natural question is why not rely purely on Kafka partitions for ordering. Partition-per-session is not viable because with thousands of concurrent sessions, the partition count would exceed operational limits and degrade broker performance. A single-threaded-per-partition model is inefficient when partitions carry interleaved session traffic, because session-level sequencing still requires application-level coordination. Consistent hashing provides strict within-session ordering with high cross-session parallelism, without unbounded partition growth.
Retry Without Breaking Order
The first retry implementation used a separate retry goroutine and queued later messages in memory, which made the state machine complex and introduced an unbounded memory risk.
The current model is simpler. Retries happen in-place inside the session goroutine. For each message, the session goroutine attempts to process it. If processing succeeds, the message is marked complete, and the next message in the session is processed immediately. If the failure is transient (for example, a downstream service timeout), the goroutine waits using exponential backoff with jitter before retrying the same message.
Because the goroutine is blocked during the retry, later messages for that session remain queued and cannot overtake the failed message. If the failure is non-recoverable, such as invalid message data, the message is routed directly to the dead-letter queue (DLQ) without retry. Once the failed message has either succeeded or been handled by the DLQ, processing resumes with the next message in the session.
Because the session goroutine sleeps during backoff, message 4 literally cannot execute before message 3 finishes. There is no queue to drain, no retry channel, no state to reset. The goroutine blocks, and new messages pile up in the session channel behind it, exactly where they should be.

Figure 2. In-place retry inside session goroutine (Image source: created by author)
This also eliminates the external retry goroutine that the original design required. The session goroutine handles retries.
Error classification follows two tiers:
- Tier 1, Transient errors: retry with exponential backoff and jitter.
- Tier 2, Poison pills/non-recoverable errors: skip retry and route directly to the DLQ.
Offset Commits and Crash Recovery
With many session goroutines completing out of order, naively committing the latest completed offset is dangerous. If offset 104 completes before 102, committing 104 means a crash would skip 102 entirely.
We solve this with a contiguous watermark managed by a per-partition offset tracking component. For each partition, it tracks InFlight (offsets currently being processed) and Completed (offsets that finished). At regular commit intervals, a commit loop advances the watermark.

Figure 3. Contiguous watermark commit (Image source: created by author)
The key insight is to only commit up to the highest contiguous completed offset. Gaps in in-flight messages act as blockers, preventing any unsafe advancement.
On crash, replay starts at offset 101. Offset 100 was already processed, while offsets 102 and 104 may be reprocessed. The pipeline is designed to make replay safe where possible by carrying a stable event ID for deduplication. This reduces duplicate processing, but external side effects still depend on the idempotency guarantees of the downstream service, so replay is not equivalent to exactly-once execution.
This design intentionally trades minimal replay work for deterministic recovery. In distributed systems, replaying safely is preferable to risking silent data loss.
In rare cases, an interrupted handoff can leave a permanent gap in the commit watermark. To prevent the partition from stalling indefinitely, the coordinator tracks unresolved gaps and, after the configured processing timeout, advances past gaps it can no longer complete. This is a last-resort availability trade-off: the system may sacrifice completeness for a permanently blocked message to restore partition progress, and the event should be treated as an operational failure rather than a normal successful delivery.
Operational Hardening
The four core primitives (consistent hashing, per-session goroutines, contiguous watermark commits, and in-place retry) solve the ordering problem. Production resilience requires several additional layers.
Rebalance Safety
During a rebalance, the consumer first marks the partition as revoking so it doesn't route new records to it. It then waits for in-flight session work to finish, checking every 50 ms and respecting the broker’s rebalance timeout. Once processing is complete, the final contiguous watermark is committed, partition-specific state is cleared, and ownership can safely move to another consumer.

Figure 4. Safe partition handoff during Kafka consumer rebalancing. (Image source: created by author)
Backpressure
When workers are saturated, the system does not drop messages or crash. It pauses fetching. Message routing uses a non-blocking channel send. If the worker's dispatch channel is full, it triggers backpressure: the partition is paused and a monitor goroutine polls capacity every second, resuming when the load drops below the configured resume threshold
Pausing alone is not enough. Pausing message fetching stops new records from arriving from the broker, but records may already have been pulled from the current batch. These records are buffered per partition and remain marked as in-flight, preventing them from being lost or advancing the watermark prematurely.
A background monitor drains the buffer as capacity becomes available. Fetching resumes only after the buffer is empty and the triggering worker has sufficient capacity, preventing newer records from overtaking older buffered messages. This improves session isolation but does not eliminate partition-level coupling: pausing a partition can temporarily delay unrelated sessions assigned to it.
Stuck Offset Detection
A session goroutine that blocks indefinitely, whether due to a non-terminating downstream call or a deadlock, will freeze the commit watermark for its partition. Every subsequent offset on that partition becomes uncommittable.
The first detector treated a continuously non-empty set of in-flight offsets as a stuck partition, which produced false positives under sustained traffic. The current approach tracks each offset individually and only flags one when its processing time exceeds the configured timeout. Stuck offsets can then be force-completed and recorded in the DLQ so the partition can continue making progress.
Tracing and Autoscaling Signals
Production observability required two additional safeguards: trace links across asynchronous Kafka boundaries and independent lag reporting while partitions were paused. This keeps tracing accurate and ensures autoscaling still receives useful lag signals during backpressure.
Session-Ordered Producing
Ordering guarantees do not stop at the consumer. The producer also needs session affinity when sending asynchronously.
Asynchronous production uses the same consistent-hashing strategy as the consumer, ensuring messages from the same session are routed to the same producer worker and sent in FIFO order. Failed sends escalate to the broker DLQ, while shutdown drains and flushes pending writes before the Kafka client closes.
Graceful Shutdown and Deployment Safety
During shutdown, the consumer stops fetching new records and gives active session workers time to finish. Once draining completes or the configured timeout is reached, it commits the final contiguous watermarks and pending DLQ writes before the Kafka clients close. Any uncommitted offsets are replayed after restart, which is safe because processing is idempotent.
Benchmarks and Validation
We validated the architecture in two stages: a synthetic benchmark to test ordering, retry behaviour, throughput, and latency under controlled failures, followed by a production-scale load test against the deployed pipeline.
Test Environment
The benchmarks below ran on AWS EC2 c5.2xlarge instances (8 vCPU, 16GB RAM), a 3-broker Kafka cluster with 10 partitions per topic, and a replication factor of 3. Each consumer instance handled 4 partitions across 3 consumer instances in the group.
The 200ms target applies to messaging overhead across ingest, routing, hand-off, and delivery; variable LLM inference time is outside that budget.
Synthetic Benchmark Results
The first validation pass used simulated downstream processing rather than the real services, with 10% of messages injected with a transient error to verify that retries did not break ordering under partial failure. The latency figures in Table 1 therefore reflect the benchmark harness and messaging path rather than real LLM inference latency.
| Metric | Result | Target |
| Raw throughput | 14,027 msg/s | >10k msg/s |
| E2E latency p50 | 11.7ms | <200ms |
| E2E latency p95 | 115ms | <200ms |
| E2E latency p99 | 185ms | <200ms |
| Ordering violations | 0 | 0 |
| Session affinity | 100% | 100% |
Table 1. Synthetic benchmark results
Zero violations across 50,000 messages, 10 concurrent sessions, with 10% failure injection.
Production-Scale Load Testing
The synthetic benchmark validated correctness, but not production-scale behaviour. We therefore ran a WebSocket-driven test against the deployed pipeline with 1,000 concurrent sessions, pushing the send rate toward 100,000 msg/s.
At a 50,000 msg/s target, with 1,000 concurrent sessions over a bit over 4 minutes:
| Metric | Result |
| Target send rate | 50,000 msg/s |
| Actual send rate | 48,805 msg/s |
| Messages sent | 11,959,510 |
| Send errors | 0 (0.00%) |
| E2E latency p50 | 446.8ms |
| E2E latency p95 | 2,633.9ms |
| E2E latency p99 | 2,843.1ms |
Table 2. Production-scale load testing results
At a 50,000 msg/s target, the deployed pipeline sustained 48,805 msg/s with zero send errors. Tail latency exceeded the 200ms messaging budget because the WebSocket path was operating beyond the API Gateway message-rate quota, rather than because of the session-ordering mechanism. Although the run ended early, it demonstrated additional send-side headroom beyond the 50,000 msg/s test.
The high-rate load tests measured throughput, latency, and error rate, but did not independently verify per-session ordering. Explicit ordering validation was performed in the synthetic benchmark, where zero violations were recorded at 14,027 msg/s with 10% failure injection. Re-running the same sequence-monotonicity check at full load-test scale remains a validation step.
Production Performance
The synthetic benchmark validated the design. Production confirmed it operationally. Since the production rollout, we've processed over 40 million messages with no observed ordering violations under live traffic and real-world conditions, including variable LLM inference times, network jitter, and organic load spikes. This production claim is based on operational observation and monitoring, not an exhaustive per-message sequence check. Only 460 of those messages were ever routed to the DLQ (poison pills, non-recoverable errors, and stuck-offset force-completions combined), a rate of about 0.002%.
How We Validate Ordering
Ordering violations are detected by assigning each message a monotonically increasing sequence number per session at the producer. The consumer-side test harness records the sequence in which messages are processed per session :

Session affinity is validated separately: every message records which worker processed it, and the test asserts that all messages for a given session were handled by exactly one worker.
.In production, consumer lag, processing duration, retries, DLQ counts, pending-buffer depth, and backpressure events provide operational signals for watermark stalls or unexpected session-to-worker reassignment. These are useful indicators, but they are not a direct sequence-level ordering check.
Why Not Redis Streams?
Before settling on Kafka with a custom consumer layer, we considered Redis Streams. Redis Streams offers consumer groups, message acknowledgment, and reasonable ordering guarantees. For simpler workloads, it can be a good fit.
For our workload, Redis Streams would still require application-level coordination for session affinity and additional bookkeeping for recovery. Kafka's partition and offset model, combined with session-aware routing in the application layer, provided a stronger foundation for strict per-session ordering, replay, and high parallelism across sessions.
Conclusion
Strict session ordering requires more than Kafka's partition-level guarantee. In this design, consistent hashing provides session affinity, per-session goroutines serialize processing and retries, contiguous watermarks make concurrent processing recoverable, and rebalance draining protects partition handoff.
Production hardening around backpressure, stuck work, and graceful shutdown turned those architectural guarantees into something that survives real operating conditions. The system has processed more than 40 million production messages with zero observed ordering violations.
These patterns apply anywhere message order carries semantic meaning, including conversational systems, workflow engines, and transaction processing.