DEV Community

Cover image for Kafka Partitions and Consumer Groups: How to Prevent Head-of-Line Blocking
DEVANSHU PATIL
DEVANSHU PATIL

Posted on AI-assisted

Kafka Partitions and Consumer Groups: How to Prevent Head-of-Line Blocking

Kafka Partitions and Consumer Groups: How to Prevent Head-of-Line Blocking

Apache Kafka is celebrated for its incredible horizontal throughput, capable of streaming millions of messages per second.

Yet, many engineering teams encounter a frustrating bottleneck in production: an individual consumer suddenly stalls or slows down, and within minutes, the consumer lag spikes into the tens of thousands. Messages destined for other users back up behind the slow task, creating a classic Head-of-Line (HoL) Blocking crisis.

Why does this happen? And how can you design your consumers to maintain high throughput even when individual messages take seconds to process?

The Golden Rule of Kafka Partitions

In Kafka:

  • A topic is split into multiple independent Partitions (the unit of parallelism).
  • A Consumer Group consists of one or more consumer processes cooperating to read from the topic.
  • The Golden Rule: Within a single consumer group, a partition can only be consumed by exactly one consumer at any given time.
Topic: user-notifications (3 Partitions)
  [Partition 0] ----------> [Consumer A]
  [Partition 1] ----------> [Consumer B]
  [Partition 2] ----------> [Consumer C]
Enter fullscreen mode Exit fullscreen mode

The Head-of-Line Blocking Trap

Because Kafka guarantees strict ordering within a single partition, a consumer cannot skip ahead. It reads message $N$, processes it, commits offset $N$, and moves to $N+1$.

Imagine Partition 0 receives:

  1. Msg 101: Fast email notification (takes 10ms).
  2. Msg 102: A massive PDF invoice generation (takes 12,000ms!).
  3. Msg 103: Password reset email (takes 10ms).

While the consumer is trapped generating the heavy PDF for Message 102, Message 103 sits completely blocked on the broker.

The Solution: Decoupled Concurrency (Poller + Internal Worker Pool)

The modern solution is to decouple the Kafka Poller from the business execution workers.

The consumer polls batches quickly from Kafka and dispatches them into an in-memory thread pool or worker queue.

[Kafka Broker]
      | (poll batch)
      v
[Kafka Poller Thread]
      |
      |--> Dispatch to Internal Worker Thread Pool (e.g. 20 Workers)
             |-- Worker 1: [Msg 101] -> Process Fast
             |-- Worker 2: [Msg 102] -> Heavy PDF Generation
             |-- Worker 3: [Msg 103] -> Process Fast
Enter fullscreen mode Exit fullscreen mode

Go Implementation with Worker Pools

package main

import (
    "context"
    "fmt"
    "sync"
    "time"
)

type KafkaMessage struct {
    Partition int
    Offset    int64
    Payload   string
}

func StartDecoupledConsumer(ctx context.Context, incomingMessages <-chan KafkaMessage) {
    const workerCount = 10
    taskQueue := make(chan KafkaMessage, 100)

    var wg sync.WaitGroup

    for i := 1; i <= workerCount; i++ {
        wg.Add(1)
        go func(workerID int) {
            defer wg.Done()
            for msg := range taskQueue {
                processMessage(workerID, msg)
            }
        }(i)
    }

    for {
        select {
        case <-ctx.Done():
            close(taskQueue)
            wg.Wait()
            return
        case msg, ok := <-incomingMessages:
            if !ok {
                return
            }
            taskQueue <- msg
        }
    }
}

func processMessage(workerID int, msg KafkaMessage) {
    fmt.Printf("[Worker %d] Processing Offset %d: %s\n", workerID, msg.Offset, msg.Payload)
    time.Sleep(100 * time.Millisecond)
}
Enter fullscreen mode Exit fullscreen mode

Kafka Tuning Parameters

max.poll.interval.ms=300000     # 5 minutes
max.poll.records=50
session.timeout.ms=45000
heartbeat.interval.ms=15000
Enter fullscreen mode Exit fullscreen mode

Top comments (0)