← All blogs
event-streaming · go · distributed-systems · platform-engineering

Event Stream Partitioning Under Skewed Key Distribution: Rebalancing Mechanics, Consumer Lag, and the Hot Partition Trap

How skewed partition keys collapse Kafka throughput, why rebalancing amplifies lag, and the architectural decisions that prevent it in Go services.

The Partition Is the Unit of Parallelism—Until It Isn't

Kafka's throughput model rests on one invariant: a partition is consumed by exactly one consumer in a group at a time. That invariant is also the source of the most common production failure in event-driven systems—hot partition saturation. When key distribution is skewed, a single partition absorbs a disproportionate fraction of write volume, saturating the consumer assigned to it while adjacent consumers sit idle. The result looks like downstream lag, not a partitioning problem, which delays diagnosis by hours.

This article examines the mechanics of partition skew, the failure modes it triggers during rebalancing, and the architectural decisions that constrain the blast radius in Go consumers.


How Skew Manifests in Production

Partition assignment for a produced message follows a hash function over the message key. The default Kafka partitioner uses murmur2 over the raw key bytes modulo the partition count. For uniformly distributed keys this produces reasonable balance. For real workloads—tenant IDs in a SaaS platform, user IDs dominated by bot traffic, device IDs from a single high-throughput sensor—the key space is not uniform. A single key or a small cohort of keys can generate 60–80% of topic write volume and land consistently on one or two partitions.

The producer does not know partition load. It hashes the key, computes the target partition, and batches records. If that partition's leader is underprovisioned or if the consumer assigned to it cannot drain records at the production rate, lag accumulates. Lag on one partition does not trigger Kafka's internal rebalancing—consumer group rebalancing is membership-driven, not lag-driven. The partition stays assigned to an overwhelmed consumer.


Consumer Group Rebalancing Makes Skew Worse

A rebalance is triggered by membership change: a consumer joins, leaves, or fails its heartbeat. During eager rebalancing (the default prior to cooperative-sticky), all consumers revoke their partitions simultaneously. Every consumer stops processing. Partition ownership is then reassigned from scratch. A saturated consumer that was falling behind now stops entirely during the rebalance window, allowing lag to accumulate further. When ownership resumes, the hot partition is likely reassigned to the same or a similarly provisioned consumer, and the cycle repeats.

Cooperative-sticky rebalancing (incremental rebalancing) mitigates total-stop behavior—only partitions being moved are revoked—but it does not redistribute load based on lag depth. The partition assignment strategy is membership-aware, not throughput-aware. A consumer processing 50,000 records/sec on a hot partition continues to hold that partition through rebalance cycles unless a custom assignor intervenes.

Session Timeout and Heartbeat Interaction

A consumer under heavy load may miss heartbeats. The heartbeat loop runs in a background goroutine in the Kafka client library. If the poll loop—where message processing actually happens—blocks for longer than max.poll.interval.ms, the broker treats the consumer as dead and triggers a rebalance. In Go consumers using confluent-kafka-go or segmentio/kafka-go, the processing goroutine and the poll loop are often the same goroutine path. Slow processing on a hot partition increases poll interval, which triggers a rebalance, which causes total lag accumulation, which increases processing time further. The feedback loop is self-reinforcing.

// Decoupled poll and processing to prevent heartbeat starvation
func runConsumer(ctx context.Context, reader *kafka.Reader, process func(kafka.Message) error) error {
    msgs := make(chan kafka.Message, 512)

    go func() {
        for {
            msg, err := reader.FetchMessage(ctx)
            if err != nil {
                if ctx.Err() != nil {
                    return
                }
                continue
            }
            select {
            case msgs <- msg:
            case <-ctx.Done():
                return
            }
        }
    }()

    for {
        select {
        case msg := <-msgs:
            if err := process(msg); err != nil {
                // handle with dead-letter or retry budget
            }
            if err := reader.CommitMessages(ctx, msg); err != nil {
                return err
            }
        case <-ctx.Done():
            return ctx.Err()
        }
    }
}

Separating the fetch goroutine from the processing goroutine decouples the heartbeat path from processing latency. The fetch goroutine keeps the Kafka connection alive; the processing goroutine can take time without triggering session timeout.


Architectural Levers for Skewed Workloads

1. Key Salting and Virtual Partitions

For keys with extreme skew, salting distributes records across multiple physical partitions. A tenant ID that generates 70% of volume is split across N virtual partitions by appending a random suffix (e.g., tenantID:0, tenantID:1, ..., tenantID:N-1). Consumers must aggregate across the virtual partition set, which introduces ordering complexity—you lose per-key ordering guarantees across the salt space. This tradeoff is acceptable for workloads where per-event ordering is not required but throughput is.

For workloads requiring ordered processing (e.g., CDC events for a specific entity), salting is not safe. The ordering guarantee is the constraint that prevents you from distributing the load.

2. Topic Repartitioning as a Deployment Operation

Increasing partition count is a controlled operation that does not redistribute existing data—only new messages route to new partitions. Consumer groups must be rebalanced after partition count changes. Tooling (kafka-reassign-partitions) can move partition leadership across brokers, but consumer assignment rebalances on its own schedule. In production, increasing partition count mid-traffic requires coordinating producer restart (so new partition count is picked up), consumer group stabilization, and monitoring lag depth across all partitions before declaring the operation complete.

Increasing partition count also has cost: more open file descriptors per broker, more replication traffic, more metadata overhead in Zookeeper or KRaft. Doubling partitions to fix a hot partition problem without diagnosing the key distribution first often creates a different skew pattern.

3. Consumer-Side Rate Shaping

A hot partition consumer can apply token bucket rate limiting internally to prevent downstream services from being overwhelmed, even if the consumer itself can drain the partition quickly. This is distinct from the Kafka-level problem—it controls how fast the consumer pushes work into downstream systems (databases, HTTP endpoints, gRPC calls).

type RateLimitedProcessor struct {
    limiter *rate.Limiter
    inner   func(kafka.Message) error
}

func (r *RateLimitedProcessor) Process(ctx context.Context, msg kafka.Message) error {
    if err := r.limiter.Wait(ctx); err != nil {
        return err
    }
    return r.inner(msg)
}

This prevents hot partition throughput from cascading into MongoDB write saturation or Redis connection exhaustion downstream. It does not reduce lag—it trades lag accumulation for downstream stability, which is often the correct tradeoff when the downstream cannot scale horizontally as fast as Kafka can produce.

4. Lag-Aware Autoscaling

KEDA (Kubernetes Event-Driven Autoscaling) supports Kafka lag as a scaling metric. A consumer deployment can scale horizontally based on partition lag, subject to the ceiling that consumer count cannot exceed partition count. For a 12-partition topic with 2 hot partitions, you can scale to 12 consumers maximum. The 2 hot partitions are still each handled by a single consumer—scaling does not help unless partition count also scales or the hot key distribution changes.

The correct mental model: autoscaling helps with uniform lag across all partitions. It cannot parallelize a single partition.


Observability Requirements

Detecting skew requires per-partition lag metrics, not aggregate consumer group lag. Aggregate lag hides hot partitions behind healthy averages. The operational baseline should include:

  • Per-partition end offset minus committed offset, exported as a labeled metric with topic, partition, and consumer_group.
  • Per-partition produce rate (records/sec), derived from broker JMX or Prometheus exporter.
  • Consumer processing time per message, tracked as a histogram in the Go consumer, labeled by source partition where the client library exposes it.

Without per-partition granularity, skew is invisible until service-level latency or error rate is already degraded.


Decision Framework

Before modifying partitioning strategy, determine which constraint is binding:

ConstraintSymptomLever
Single high-volume keyOne partition at 10x avg produce rateKey salting if ordering not required; dedicated topic if ordering required
Slow downstreamConsumer drains fast, DB writes back upConsumer-side rate limiting, downstream capacity scaling
Broker CPU/networkHot partition leader saturatedPartition reassignment to less loaded broker
Consumer processing timePoll interval exceeded, rebalances triggeredDecouple fetch and process goroutines, tune max.poll.interval.ms
Insufficient parallelismLag uniform across partitionsIncrease partition count with coordinated rollout

Hot partition problems are frequently misdiagnosed as consumer scaling problems. The distinction matters because scaling consumers when the problem is partition-count-bounded or key-distribution-bounded wastes resources and delays the correct fix. Measure per-partition lag, identify the skewed key cohort, then select the lever that addresses the actual binding constraint.

Event Stream Partitioning Under Skewed Key Distribution: Rebalancing Mechanics, Consumer Lag, and the Hot Partition Trap | Neeraj Singhi