05 / 05

How would you introduce parallel processing inside a consumer while preserving per-key ordering?

Difficulty: 9/10
Polling and processing

Keep partition ownership on the poll thread, route records to a fixed set of single-threaded workers by key hash, and commit only contiguous completed offsets

The first question is whether I need this at all. The simplest way to scale is more partitions and more consumers, because Kafka already gives per-partition ordering and parallelism. Internal concurrency is justified when the partition count is capped (ordering domain, existing key mapping, metadata cost) or handlers are slow I/O-bound calls where threads mostly wait. If I do need it, the design principle is to separate partition ownership from execution. One thread owns the KafkaConsumer, polls, tracks offsets and handles rebalances. Execution happens elsewhere, and ordering is preserved by routing: hash the record key to one of N single-threaded executors. The same key always goes to the same executor, so its records run sequentially, while different keys run in parallel.

The hard parts are offsets and lifecycle, not the thread pool. Records complete out of order across workers, so I cannot commit the offset of the latest finished record: that would acknowledge earlier unfinished ones. I track in-flight offsets per partition and commit only up to the lowest incomplete offset, a contiguous watermark. I need backpressure: if workers fall behind, queues grow without bound, so I pause() the partitions above a high-water mark and resume() below a low-water mark, continuing to poll so the consumer stays in the group. And I need rebalance handling: in onPartitionsRevoked I wait (bounded) for in-flight work on revoked partitions to finish, commit their watermark, and drop their state, so the next owner does not duplicate or reorder work. The delivery semantic remains at-least-once, so handlers must be idempotent.

javascript
  1. 1

    Trade-off: hash-based routing gives strict per-key order but can create hot workers if keys are skewed, since one heavy key keeps one worker busy while others idle. Mitigate with more workers than cores for I/O-bound work, or accept that a single hot key is inherently serial.

  2. 2

    Trade-off: the watermark means one slow record holds back the commit for its whole partition. That is correct for safety, but a stuck record can cause large replay after a crash, so add handler timeouts and a retry or DLQ policy that never blocks a worker forever.

  3. 3

    Alternative: more partitions and more consumers. Simplest, no custom concurrency, but partition count cannot be reduced and changing it remaps keys. I pick internal concurrency only when that route is blocked or handlers are mostly waiting on I/O.

  4. 4

    Alternative: an existing library such as the open-source Confluent Parallel Consumer implements key-ordered concurrent processing and offset tracking. I would evaluate it before maintaining my own, weighing its maturity, maintenance status and operational fit for my version.

  5. 5

    Alternative: if order does not matter at all, a queue-like model (Kafka share groups, KIP-932, where available in your release and marked production-ready) or a plain worker pool removes the key-routing and watermark complexity.

  6. 6

    Common mistake: committing the offset of the latest completed record, or committing from worker threads. KafkaConsumer is not thread-safe and out-of-order completion makes that commit acknowledge unfinished work.

  7. 7

    Common mistake: forgetting backpressure, so queues grow until memory runs out, or pausing without continuing to poll, which triggers the max.poll.interval.ms removal.

  8. 8

    Common mistake: ignoring rebalances. Without draining and committing on revocation, the new owner reprocesses in-flight records while the old workers are still running them, which breaks ordering for those keys.

  9. 9

    Semantics: this is at-least-once. Handlers must be idempotent, and if the work writes back to Kafka with exactly-once needs, the transactional read-process-write pattern conflicts with out-of-order completion and needs a different design.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.