Command Palette

Search for a command to run...

Hectal
PHASE 3Intermediate ~8 min· topic 5 of 5

Topic 3.5

Processing Models: Threads, Pause/Resume, Long Work

In one line

You can process records single-threaded per consumer, with one consumer per thread, or with a worker pool behind one polling thread. Parallelism inside a partition must be per key to keep ordering, offsets must be committed only when contiguous work is done, and long processing needs pause/resume so polling continues.

0/5 · 0%

Think of it like this

A restaurant pass. The expediter (polling thread) takes tickets and hands them to cooks (workers). Two dishes for the same table (key) must be cooked in order by the same cook; different tables can be cooked in parallel. The expediter only calls a table "served" when every earlier ticket for that section is done.

Key ideas

  1. 01

    Single-threaded per consumer: simplest, preserves order per partition; scale by adding consumers up to the partition count.

  2. 02

    One consumer per thread: several KafkaConsumer instances in one process, each single-threaded; same semantics, better CPU use per machine.

  3. 03

    Worker pool: one polling thread dispatches records to workers. To keep per-key order, hash keys to a fixed set of single-threaded queues. Track completion per partition and commit the highest contiguous completed offset. Libraries like Confluent's Parallel Consumer implement this (key-level or unordered parallelism beyond partition count).

  4. 04

    Long-running work (minutes per record): don't block poll(). Hand work off, call consumer.pause(partitions) so poll() returns nothing for them while still keeping the member alive, and resume() when done. Or raise max.poll.interval.ms and lower max.poll.records if the work is bounded.

  5. 05

    Backpressure: bounded worker queues plus pause/resume keep memory stable when downstream is slow.

Code & diagrams

KeyedWorkerPool.javajava

Per-key order preserved by routing each key to one single-threaded lane; commits track contiguous completion.

ExecutorService[] lanes = new ExecutorService[16];
for (int i = 0; i < lanes.length; i++) lanes[i] = Executors.newSingleThreadExecutor();
OffsetTracker tracker = new OffsetTracker();          // per partition: completed offsets -> highest contiguous

while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200));
    for (ConsumerRecord<String, String> r : records) {
        tracker.started(r);
        int lane = Math.floorMod(r.key().hashCode(), lanes.length);
        lanes[lane].submit(() -> { handle(r); tracker.completed(r); });
    }
    if (tracker.inFlight() > 5_000) consumer.pause(consumer.assignment());   // backpressure
    else consumer.resume(consumer.paused());
    Map<TopicPartition, OffsetAndMetadata> safe = tracker.committable();    // highest contiguous + 1
    if (!safe.isEmpty()) consumer.commitAsync(safe, null);
}

Interview problem

The problem

A batch takes 10 minutes to process

A consumer takes 10 minutes to process one batch, but max.poll.interval.ms is 5 minutes. What happens, and how do you design around it?

When it breaks

Records from one partition processed in parallel on an unordered thread pool

What you see

Events for the same order are applied out of order: Shipped before Paid; state machines reject or corrupt state.

Fix & prevent

Route by key to single-threaded lanes, or use a library with key-ordered parallelism.

Explain it without notes

01

How do you process records in parallel within a partition without breaking per-key order or offset safety?

Practice

01

Simulate processing that takes 8 minutes with default settings and observe the rebalance and CommitFailedException; then fix it with pause/resume.

Trade-offs

  • ↔

    Worker pools raise throughput beyond partition count but require careful ordering and offset tracking.

Done when you can

  • I can choose a processing model and keep order and offsets correct.

  • I handle long-running processing without poll-interval rebalances.