What does num.recovery.threads.per.data.dir control, and when does increasing it help?
answer
- per data dir, default 1
- total = value × #log.dirs
- matters after unclean shutdown
- scans recovery-point → log end
- startup/shutdown only, not steady-state
basics
~20 sIt sets how many threads, per log directory, Kafka uses to load and recover partition logs at startup and to flush them at shutdown. Raising it speeds up startup recovery on brokers with many partitions, especially after an unclean shutdown.
solid answer
~50 s`num.recovery.threads.per.data.dir` (default 1) controls the thread pool size, *per* `log.dirs` entry, that Kafka uses to (a) load and sanity-check all log segments on startup and (b) flush logs on shutdown. The total recovery parallelism is this value multiplied by the number of data directories. On a clean shutdown, Kafka has flushed everything and written the recovery-point checkpoints, so startup just loads segments quickly. After an *unclean* shutdown (crash/kill -9/power loss), the broker must scan and re-validate every segment from the last recovery point to the end of the log on the affected partitions — this is the expensive case. A broker with thousands of partitions can take many minutes to recover with a single thread. Increasing recovery threads parallelizes this across partitions within each directory, cutting startup time, at the cost of more CPU and disk I/O during recovery. It does not affect steady-state throughput.
go deeper
Know it controls how many threads recover logs when the broker starts up.
Explain it is per data directory, default 1, and speeds up startup after a crash.
Tie it to clean vs unclean shutdown, recovery-point scanning, and the under-replication window during startup.
Set fleet defaults balancing startup SLA against CPU/disk burst, factoring partition counts, instance volatility, and disk parallelism.
## What 'recovery' means here When a Kafka broker starts, the `LogManager` must bring every partition's on-disk log into a usable state before the broker can serve traffic. This involves: - Opening each partition directory and its segment files. - Loading/validating the offset and time indexes. - For partitions whose last flush is uncertain, **scanning records from the recovery point to the end of the log** to verify integrity and rebuild indexes — this is *log recovery*. The symmetric work happens at shutdown: flushing logs and writing checkpoints. ## The config ``` num.recovery.threads.per.data.dir=1 # default ``` It sizes a thread pool **per data directory** (`log.dirs` entry). So effective parallelism = `num.recovery.threads.per.data.dir` × `count(log.dirs)`. The threads work across the partitions in a directory concurrently; a single partition is still recovered by one thread. ## Clean vs unclean shutdown — why it matters - **Clean shutdown:** the broker flushed all data and advanced the `recovery-point-offset-checkpoint` file to the log end. On restart there is essentially nothing to re-scan; loading is fast regardless of thread count, because the recovery point equals the log end offset. - **Unclean shutdown** (SIGKILL, OOM-kill, host crash, power loss): the recovery point lags behind the actual log end. On restart Kafka must re-validate every segment from the recovery point forward on each affected partition. With one thread and thousands of partitions, this serializes into a long startup (potentially many minutes), during which the broker is unavailable and its partitions are under-replicated. ## When raising it helps - Brokers hosting **many partitions** (hundreds to thousands). - Environments where **unclean shutdowns are plausible** (spot/preemptible instances, frequent restarts). - Multiple `log.dirs` already multiply parallelism; tuning the per-dir count adds more within each disk. A common production value is 2-8 (or roughly matched to CPU/disk parallelism). Setting it too high wastes CPU and can saturate disk during the recovery burst without further benefit. ## What it does NOT do - It does not affect normal request handling, produce/fetch latency, or replication throughput — those use other thread pools (`num.io.threads`, `num.network.threads`, replica fetchers). - It does not change durability; it only changes how fast recovery completes. ## Edge cases - With a single data directory, total recovery threads = this value alone. - If recovery is I/O-bound (slow disks), adding threads beyond disk parallelism gives little gain. - Faster recovery shortens the window of under-replication and unavailability, which can be the real operational driver (SLA), not raw speed. ## Bottom line It is a startup/shutdown parallelism knob, per log dir, that matters most after unclean shutdowns on many-partition brokers; raise it to shorten recovery time within the limits of CPU and disk.
- Why is recovery cheap after a clean shutdown but expensive after a crash?A clean shutdown flushes data and advances the recovery-point checkpoint to the log end, so there is nothing to re-scan. After a crash the recovery point lags the log end, so Kafka must re-validate every segment from the recovery point forward on each affected partition.
- Does raising num.recovery.threads.per.data.dir improve steady-state throughput?No. It only affects startup recovery and shutdown flush. Steady-state produce/fetch uses separate pools (num.io.threads, num.network.threads, replica fetcher threads); recovery threads are idle during normal operation.
saying these in an interview costs you the question
- Saying it tunes normal request/throughput performance — it only affects startup/shutdown.
- Thinking the value is the total number of recovery threads, ignoring that it's per data directory.
- Claiming recovery is equally expensive after clean and unclean shutdown.