← Writing

Why My DIPY Array Jobs Were Taking Forever on HPC

My DIPY MP-PCA jobs kept timing out on HPC. The likely culprit was thread oversubscription.

I was running MP-PCA denoising on diffusion MRI data using DIPY on an HPC cluster. The jobs started normally, but the runtimes were wildly inconsistent. Some scans finished in a couple of hours, others took 16 hours, and many more hit the 24-hour walltime limit without finishing or reporting any errors.

It took me longer than I’d like to admit to figure out why. At first, I thought some scans were simply more computationally demanding, or that certain jobs were running on slower cluster nodes. After exploring a few possibilities, I finally traced the problem to thread oversubscription in the underlying numerical libraries.

The fix was surprisingly simple: adding three lines to my job script to limit the number of threads. After that, scans that had previously timed out after 24 hours were finishing in about an hour.

What was actually happening

I was submitting one job per scan as a job array, with each job requesting one CPU core. That seemed reasonable: the scans were independent, and I wanted to process as many as possible in parallel.

The problem was that requesting one CPU core doesn’t necessarily mean the program will use only one thread.

DIPY’s MP-PCA denoising relies on NumPy and SciPy for operations like matrix multiplication and eigenvalue decomposition. These operations can use numerical libraries like OpenBLAS and Intel MKL, which support multithreading.

So even though each of my jobs requested only one core, the underlying libraries were likely creating many additional threads per job. With many jobs running simultaneously on the same node, all those threads end up competing for the same physical cores.

This is called thread oversubscription: too many threads competing for the available CPU cores. The overhead from scheduling and switching between threads can slow everything down dramatically. The jobs were technically running, but making very little progress.

A bit of background on CPUs, cores, and threads

A CPU is the physical processor, which contains multiple cores. An HPC node can have one or more CPUs, often with dozens of cores in total. Each core can run computations independently, so having more cores allows more work to happen in parallel.

A thread is a unit of work within a program. Unlike cores, threads are software, not hardware. A program can create multiple threads to run different parts of a calculation in parallel, but those threads still need CPU cores to execute.

Libraries like OpenBLAS and Intel MKL can use multiple threads to speed up calculations. If the number of threads isn’t explicitly specified, these libraries often determine it automatically based on how many CPU cores they detect as available to the process. On an HPC cluster, this may include cores on the entire node, rather than just those allocated to the job.

The fix

Once I figured out what was happening, the fix was straightforward.

I added these three lines to my job script before launching Python:

export OMP_NUM_THREADS=1
export MKL_NUM_THREADS=1
export OPENBLAS_NUM_THREADS=1

OMP_NUM_THREADS sets the default thread count for OpenMP, MKL_NUM_THREADS controls Intel MKL, and OPENBLAS_NUM_THREADS controls OpenBLAS.

I set all three because the numerical backend being used depends on how NumPy and SciPy were installed, which isn’t always obvious. MKL also gives its own thread-count setting precedence over OMP_NUM_THREADS, so setting both explicitly helps avoid surprises.

Why this makes sense in retrospect

At first, limiting each job to one thread seemed counterintuitive. Wouldn’t more threads make the denoising faster?

The libraries were likely detecting the total cores on the node and creating threads accordingly, regardless of how many cores were actually allocated to each job. With many jobs doing this simultaneously, hundreds of threads would end up competing for the same node-wide pool of cores.

By explicitly limiting the numerical libraries to one thread per job, I avoided that unnecessary overhead. And since I was already processing hundreds of independent scans in parallel through a job array, there was no need for each job to create additional threads.

With the fix in place, scans were finishing in around an hour each, including those that had previously timed out or taken much longer.