skip to content

An ops team complains their stateful Kafka Streams aggregation can't keep up despite adding more machines. Throughput plateaus. Diagnose the likely cause and the remediation, including the risks.

level: principalimportance: should knowfreq 35%

answer

  1. throughput plateau = partition ceiling
  2. tasks = input/repartition partitions
  3. more machines -> idle slots
  4. repartition re-hashes keys, breaks state/order
  5. over-partition up front for headroom

basics

~20 s

Parallelism is capped at the input topic's partition count. Once instances/threads equal the partition count, extra machines stay idle. Fix it by increasing partitions on the input (or repartition topic) — but that changes key routing and can break ordering/state, so it needs a careful reset/migration.

solid answer

~50 s

The plateau is the partition-count ceiling: task count for the aggregating sub-topology equals the partition count of its input (or the repartition topic feeding it). Once total threads across all instances reach that number, adding machines just creates idle slots. To diagnose, check the input/repartition topic partition count vs the number of active tasks, and look for skew (a few hot keys/partitions). Remediation options: (1) increase the input topic's partitions so more tasks exist — but this re-hashes keys to new partitions, breaking the guarantee that a key always lands in the same partition; existing per-key ordering and state-store locality are disrupted, usually requiring an application reset and state rebuild. (2) Address key skew by improving the key distribution or adding a sub-key. (3) Pre-size partitions generously up front to leave headroom. Because repartitioning is so disruptive, the durable answer is capacity-planning partition counts before launch.

go deeper

for a junior

Recognize that you can't scale past the input partition count.

for a middle

Diagnose by comparing partition count to active tasks/threads and spot idle instances.

for a senior

Distinguish ceiling vs skew, and explain why repartitioning forces a reset and state rebuild.

for a principal

Drive capacity planning: over-partition up front, plan migration/cutover strategy, and weigh ordering/state risks against throughput needs.

## Diagnosis: the partition-count ceiling Kafka Streams parallelism is bounded by **input partition count**. The aggregating sub-topology has **one task per partition** of its input (or of the **repartition topic** that feeds the `groupBy`/aggregate). Total worker slots = sum of `num.stream.threads` across all instances. Once **slots ≥ tasks**, every extra slot is **idle** — so adding machines yields no throughput gain. A flat throughput curve as you scale out is the classic signature. Secondary cause: **partition skew / hot keys**. Even below the ceiling, if a few keys dominate, the partitions holding them become the bottleneck while others idle. One overloaded task caps end-to-end lag regardless of total capacity. ### How to confirm - Compare the input/repartition topic's **partition count** to the number of **active tasks** (and to total threads). If threads ≥ partitions, you're ceiling-bound. - Inspect per-partition consumer lag; uneven lag points to **skew**, even lag at full thread utilization points to the **ceiling**. - Check CPU: ceiling-bound idle instances show low CPU; skew shows one instance/thread hot. ## Remediation ### Option A — increase input partition count More partitions → more tasks → more parallelism. **The catch**: partitions are chosen by `hash(key) % numPartitions`. Changing `numPartitions` **re-routes existing keys to different partitions**. Consequences: - **Per-key ordering** across the boundary is broken (a key's old and new records may sit in different partitions). - **State stores** are partitioned by the old scheme; their changelog data no longer lines up with the new partition→task mapping. The local state is effectively invalid. - You typically must **stop the app, reset it** (`kafka-streams-application-reset`), repartition the topic, and **rebuild state** from source — a planned migration with downtime or a parallel-run cutover. - Upstream producers must also adopt the new partition count consistently. ### Option B — fix skew If the issue is hot keys, increasing partitions won't help (the hot key still goes to one partition). Instead: introduce a **composite/sub-key** (salting) to spread a hot key across partitions, then re-aggregate, or pre-aggregate upstream. ### Option C — vertical headroom / tuning If instances are under-threaded, raise `num.stream.threads` up to the task count first (cheap, no repartition). Tune RocksDB, batching, and `commit.interval.ms` to raise per-task throughput. ## The durable lesson (principal-level) Because repartitioning is so disruptive, **partition count is a capacity-planning decision made before launch**. Pick input partition counts that leave scaling headroom (e.g. a multiple of expected max instances), accepting some per-partition overhead, so you can scale out later without a re-key migration. This is why teams 'over-partition' high-throughput topics deliberately. ## Edge cases - **Stateless** pipelines tolerate repartitioning better (no state to rebuild) but still break per-key ordering. - **Standby replicas / more threads** never raise the ceiling — they're availability and CPU utilization, not added tasks. - A **repartition topic** deep in the topology imposes its own ceiling; sometimes you can bump *its* partition count (via `Repartitioned.numberOfPartitions`) without touching the source, but co-partitioning constraints with joined streams still apply.

  • Why doesn't increasing num.standby.replicas or num.stream.threads fix the plateau?
    Neither adds tasks. Standbys are passive failover copies and never process records; threads only run the existing fixed set of tasks. Once threads equal the task count (= input partitions), both are exhausted as levers — only more partitions add parallelism.
  • Why is increasing partition count on a live stateful topic so risky?
    Partitioning is hash(key) % numPartitions. Changing numPartitions re-routes keys, so a key's history can split across partitions — breaking per-key ordering and invalidating state stores/changelogs that were built under the old mapping. It generally requires an application reset and full state rebuild.
  • How would you avoid this problem at design time?
    Capacity-plan partition counts before launch with headroom (over-partition high-throughput topics) so you can scale out by adding instances/threads later without a disruptive re-key migration.

saying these in an interview costs you the question

  • Recommending just adding more instances/threads when already at the ceiling
  • Increasing partition count on a stateful topic without acknowledging the re-key/state-rebuild risk
  • Confusing the input-partition ceiling with broker replication or with standby replicas

context