skip to content

What does num.stream.threads control, and how do threads relate to tasks and instances when scaling a Kafka Streams app?

level: middleimportance: must knowfreq 65%

answer

  1. num.stream.threads = threads per instance
  2. partitions->tasks->threads->instances
  3. scale up vs scale out, same ceiling
  4. max useful threads = total tasks
  5. addStreamThread() at runtime (KIP-663)

basics

~20 s

num.stream.threads sets how many StreamThreads run inside one application instance. Tasks are distributed across all threads of all instances. You scale up (more threads per instance) or out (more instances), but never beyond the total number of tasks.

solid answer

~40 s

num.stream.threads is the per-instance config for the number of StreamThreads. Each StreamThread is an OS thread that runs an independent consumer and executes a subset of the assigned tasks. Total parallelism = sum of all threads across all instances, capped by the total task count. You have two scaling axes: vertical (raise num.stream.threads to use more cores on one machine) and horizontal (start more instances, all in the same application.id consumer group). Both end up at the same place: tasks get spread over the available threads. If threads + instances exceed the task count, extra threads stay idle. A practical heuristic is to set num.stream.threads near the number of cores and then scale out for fault tolerance, keeping total threads <= total tasks.

go deeper

for a junior

Know num.stream.threads = threads per instance, default 1.

for a middle

Explain the partitions->tasks->threads->instances chain and the two scaling axes.

for a senior

Trade off scale-up vs scale-out for fault tolerance and the total-threads <= total-tasks ceiling.

for a principal

Design capacity + failure-domain strategy: instance sizing, runtime thread management (KIP-663), and where the partition ceiling forces a repartition.

## Threads vs tasks vs instances - **Instance**: one running JVM process of your Streams app, all sharing the same `application.id`. - **StreamThread**: a worker thread inside an instance. `num.stream.threads` (default 1) sets how many each instance creates. Each StreamThread owns its own Kafka **consumer** (and producer) and polls/processes independently. - **Task**: the unit of work (owns specific partitions). Tasks are assigned to StreamThreads. The assignment chain is: **partitions → tasks → threads → instances**. Streams' partition assignor (`StreamsPartitionAssignor`) distributes the fixed set of tasks across every StreamThread in the consumer group, balancing load and (since 2.4) honoring stickiness for state locality. ## The two scaling axes 1. **Vertical (scale up)**: increase `num.stream.threads` on an instance to use more CPU cores on one machine. More threads = more tasks can run concurrently on that box. 2. **Horizontal (scale out)**: start more instances with the same `application.id`. They join the consumer group; a rebalance reassigns tasks across all threads of all instances. Both converge on the same total: **total worker slots = sum of num.stream.threads over all instances**. Tasks are dealt out into those slots. ## The hard ceiling Because the number of tasks is fixed by input partition count (see the 'what is a task' question), **the maximum useful number of threads = total task count**. If you have 6 tasks and run 4 instances with 3 threads each (12 thread slots), 6 slots will be idle. Idle threads consume few resources but give no throughput. ## Practical guidance - A common starting point: `num.stream.threads` ≈ number of CPU cores on the instance. - For fault tolerance you usually prefer **more smaller instances** over one big instance with many threads, so a single machine failure loses less. - Changing `num.stream.threads` requires a restart of that instance (it's read at startup), and triggers a rebalance. Note: since KIP-663, threads can be added/removed at runtime via `KafkaStreams#addStreamThread()` / `removeStreamThread()` without a full restart. ## Edge cases - A single StreamThread can run multiple tasks (it loops over them in its poll cycle), so threads < tasks is normal and fine. - An uncaught exception in a StreamThread can kill just that thread (depending on the `StreamsUncaughtExceptionHandler` decision: REPLACE_THREAD keeps the app alive by spawning a new one).

  • You have 12 tasks and want fault tolerance. Would you prefer 1 instance with 12 threads or 4 instances with 3 threads each?
    4 instances with 3 threads each — losing one machine only costs 3 tasks (which rebalance elsewhere) instead of the whole app, and you spread CPU/IO and standby state across more failure domains.
  • Can you change the number of threads without restarting the app?
    Yes, since KIP-663 you can call KafkaStreams#addStreamThread() and removeStreamThread() at runtime; changing the num.stream.threads config value itself still needs a restart of that instance.

saying these in an interview costs you the question

  • Confusing threads with tasks (saying #threads determines #tasks)
  • Claiming more threads always increases throughput regardless of task count
  • Saying num.stream.threads is a cluster-wide/global setting (it's per instance)

context