<metaproperty="og:description" content="Parallelism: Some scikit-learn estimators and utilities parallelize costly operations using multiple CPU cores. Depending on the type of estimator and sometimes the values of the constructor parame..." />
<metaname="description" content="Parallelism: Some scikit-learn estimators and utilities parallelize costly operations using multiple CPU cores. Depending on the type of estimator and sometimes the values of the constructor parame..." />
<title>10.3. Parallelism and resource management — scikit-learn 1.10.dev0 documentation</title>
<liclass="toctree-l1 has-children"><aclass="reference internal" href="../model_selection.html">3. Model selection and evaluation</a><details><summary><spanclass="toctree-toggle" role="presentation"><iclass="fa-solid fa-chevron-down"></i></span></summary><ul>
<liclass="toctree-l2"><aclass="reference internal" href="../modules/grid_search.html">3.2. Tuning the hyper-parameters of an estimator</a></li>
<liclass="toctree-l2"><aclass="reference internal" href="../modules/classification_threshold.html">3.3. Tuning the decision threshold for class prediction</a></li>
<liclass="toctree-l2"><aclass="reference internal" href="../modules/model_evaluation.html">3.4. Metrics and scoring: quantifying the quality of predictions</a></li>
<liclass="toctree-l2"><aclass="reference internal" href="../modules/learning_curve.html">3.5. Validation curves: plotting scores to evaluate models</a></li>
<liclass="toctree-l2"><aclass="reference internal" href="../datasets/loading_other_datasets.html">9.4. Loading other datasets</a></li>
</ul>
</details></li>
<liclass="toctree-l1 current active has-children"><aclass="reference internal" href="../computing.html">10. Computing with scikit-learn</a><detailsopen="open"><summary><spanclass="toctree-toggle" role="presentation"><iclass="fa-solid fa-chevron-down"></i></span></summary><ulclass="current">
<liclass="toctree-l2"><aclass="reference internal" href="scaling_strategies.html">10.1. Strategies to scale computationally: bigger data</a></li>
<liclass="breadcrumb-item"><ahref="../computing.html" class="nav-link"><spanclass="section-number">10. </span>Computing with scikit-learn</a></li>
<liclass="breadcrumb-item active" aria-current="page"><spanclass="ellipsis"><spanclass="section-number">10.3. </span>Parallelism and resource management</span></li>
</ul>
</nav>
</div>
</div>
</div>
</div>
<divid="searchbox"></div>
<articleclass="bd-article">
<sectionid="parallelism-and-resource-management">
<h1><spanclass="section-number">10.3. </span>Parallelism and resource management<aclass="headerlink" href="#parallelism-and-resource-management" title="Link to this heading">#</a></h1>
<sectionid="parallelism">
<spanid="id1"></span><h2><spanclass="section-number">10.3.1. </span>Parallelism<aclass="headerlink" href="#parallelism" title="Link to this heading">#</a></h2>
<p>Some scikit-learn estimators and utilities parallelize costly operations
using multiple CPU cores.</p>
<p>Depending on the type of estimator and sometimes the values of the
constructor parameters, this is either done:</p>
<ulclass="simple">
<li><p>with higher-level parallelism via <aclass="reference external" href="https://joblib.readthedocs.io/en/latest/">joblib</a>.</p></li>
<li><p>with lower-level parallelism via OpenMP, used in C or Cython code.</p></li>
<li><p>with lower-level parallelism via BLAS, used by NumPy and SciPy for generic operations
on arrays.</p></li>
</ul>
<p>The <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs</span></code> parameters of estimators always controls the amount of parallelism
managed by joblib (processes or threads depending on the joblib backend).
The thread-level parallelism managed by OpenMP in scikit-learn’s own Cython code
or by BLAS & LAPACK libraries used by NumPy and SciPy operations used in scikit-learn
is always controlled by environment variables or <codeclass="docutils literal notranslate"><spanclass="pre">threadpoolctl</span></code> as explained below.
Note that some estimators can leverage all three kinds of parallelism at different
points of their training and prediction methods.</p>
<p>We describe these 3 types of parallelism in the following subsections in more detail.</p>
<h3><spanclass="section-number">10.3.1.1. </span>Higher-level parallelism with joblib<aclass="headerlink" href="#higher-level-parallelism-with-joblib" title="Link to this heading">#</a></h3>
<p>When the underlying implementation uses joblib, the number of workers
(threads or processes) that are spawned in parallel can be controlled via the
<p>When using <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs</span><spanclass="pre">></span><spanclass="pre">1</span></code> (or <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs=-1</span></code>), you may observe a delay
the first time a parallel function is called. This is expected behavior
caused by the overhead of starting the Python worker processes.
Subsequent calls will be faster as they reuse the existing pool of workers.</p>
</div>
<divclass="admonition note">
<pclass="admonition-title">Note</p>
<p>Where (and how) parallelization happens in the estimators using joblib by
specifying <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs</span></code> is currently poorly documented.
Please help us by improving our docs and tackle <aclass="reference external" href="https://github.com/scikit-learn/scikit-learn/issues/14228">issue 14228</a>!</p>
</div>
<p>Joblib is able to support both multi-processing and multi-threading. Whether
joblib chooses to spawn a thread or a process depends on the <strong>backend</strong>
that it’s using.</p>
<p>scikit-learn generally relies on the <codeclass="docutils literal notranslate"><spanclass="pre">loky</span></code> backend, which is joblib’s
default backend. Loky is a multi-processing backend. When doing
multi-processing, in order to avoid duplicating the memory in each process
(which isn’t reasonable with big datasets), joblib will create a <aclass="reference external" href="https://docs.scipy.org/doc/numpy/reference/generated/numpy.memmap.html">memmap</a>
that all processes can share, when the data is bigger than 1MB.</p>
<p>In some specific cases (when the code that is run in parallel releases the
GIL), scikit-learn will indicate to <codeclass="docutils literal notranslate"><spanclass="pre">joblib</span></code> that a multi-threading
backend is preferable.</p>
<p>As a user, you may control the backend that joblib will use (regardless of
what scikit-learn recommends) by using a context manager:</p>
<spanclass="c1"># Your scikit-learn code here</span>
</pre></div>
</div>
<p>Please refer to the <aclass="reference external" href="https://joblib.readthedocs.io/en/latest/parallel.html#thread-based-parallelism-vs-process-based-parallelism">joblib’s docs</a>
for more details.</p>
<p>In practice, whether parallelism is helpful at improving runtime depends on
many factors. It is usually a good idea to experiment rather than assuming
that increasing the number of workers is always a good thing. In some cases
it can be highly detrimental to performance to run multiple copies of some
estimators or functions in parallel (see <aclass="reference internal" href="#oversubscription"><spanclass="std std-ref">oversubscription</span></a> below).</p>
</section>
<sectionid="lower-level-parallelism-with-openmp">
<spanid="id2"></span><h3><spanclass="section-number">10.3.1.2. </span>Lower-level parallelism with OpenMP<aclass="headerlink" href="#lower-level-parallelism-with-openmp" title="Link to this heading">#</a></h3>
<p>OpenMP is used to parallelize code written in Cython or C, relying on
multi-threading exclusively. By default, the implementations using OpenMP
will use as many threads as possible, i.e. as many threads as logical cores.</p>
<p>You can control the exact number of threads that are used either:</p>
<ul>
<li><p>via the <codeclass="docutils literal notranslate"><spanclass="pre">OMP_NUM_THREADS</span></code> environment variable, for instance when:
<li><p>or via <codeclass="docutils literal notranslate"><spanclass="pre">threadpoolctl</span></code> as explained by <aclass="reference external" href="https://github.com/joblib/threadpoolctl/#setting-the-maximum-size-of-thread-pools">this piece of documentation</a>.</p></li>
<h3><spanclass="section-number">10.3.1.3. </span>Parallel NumPy and SciPy routines from numerical libraries<aclass="headerlink" href="#parallel-numpy-and-scipy-routines-from-numerical-libraries" title="Link to this heading">#</a></h3>
<p>scikit-learn relies heavily on NumPy and SciPy, which internally call
multi-threaded linear algebra routines (BLAS & LAPACK) implemented in libraries
such as MKL, OpenBLAS or BLIS.</p>
<p>You can control the exact number of threads used by BLAS for each library
using environment variables, namely:</p>
<ulclass="simple">
<li><p><codeclass="docutils literal notranslate"><spanclass="pre">MKL_NUM_THREADS</span></code> sets the number of threads MKL uses,</p></li>
<li><p><codeclass="docutils literal notranslate"><spanclass="pre">OPENBLAS_NUM_THREADS</span></code> sets the number of threads OpenBLAS uses</p></li>
<li><p><codeclass="docutils literal notranslate"><spanclass="pre">BLIS_NUM_THREADS</span></code> sets the number of threads BLIS uses</p></li>
</ul>
<p>Note that BLAS & LAPACK implementations can also be impacted by
<codeclass="docutils literal notranslate"><spanclass="pre">OMP_NUM_THREADS</span></code>. To check whether this is the case in your environment,
you can inspect how the number of threads effectively used by those libraries
is affected when running the following command in a bash or zsh terminal
for different values of <codeclass="docutils literal notranslate"><spanclass="pre">OMP_NUM_THREADS</span></code>:</p>
<p>At the time of writing (2022), NumPy and SciPy packages which are
distributed on pypi.org (i.e. the ones installed via <codeclass="docutils literal notranslate"><spanclass="pre">pip</span><spanclass="pre">install</span></code>)
and on the conda-forge channel (i.e. the ones installed via
<codeclass="docutils literal notranslate"><spanclass="pre">conda</span><spanclass="pre">install</span><spanclass="pre">--channel</span><spanclass="pre">conda-forge</span></code>) are linked with OpenBLAS, while
NumPy and SciPy packages shipped on the <codeclass="docutils literal notranslate"><spanclass="pre">defaults</span></code> conda
channel from Anaconda.org (i.e. the ones installed via <codeclass="docutils literal notranslate"><spanclass="pre">conda</span><spanclass="pre">install</span></code>)
<spanid="oversubscription"></span><h3><spanclass="section-number">10.3.1.4. </span>Oversubscription: spawning too many threads<aclass="headerlink" href="#oversubscription-spawning-too-many-threads" title="Link to this heading">#</a></h3>
<p>It is generally recommended to avoid using significantly more processes or
threads than the number of CPUs on a machine. Over-subscription happens when
a program is running too many threads at the same time.</p>
<p>Suppose you have a machine with 8 CPUs. Consider a case where you’re running
a <aclass="reference internal" href="../modules/generated/sklearn.model_selection.GridSearchCV.html#sklearn.model_selection.GridSearchCV" title="sklearn.model_selection.GridSearchCV"><codeclass="xref py py-class docutils literal notranslate"><spanclass="pre">GridSearchCV</span></code></a> (parallelized with joblib)
with <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs=8</span></code> over a
(since you have 8 CPUs). That’s a total of <codeclass="docutils literal notranslate"><spanclass="pre">8</span><spanclass="pre">*</span><spanclass="pre">8</span><spanclass="pre">=</span><spanclass="pre">64</span></code> threads, which
leads to oversubscription of threads for physical CPU resources and thus
to scheduling overhead.</p>
<p>Oversubscription can arise in the exact same fashion with parallelized
routines from MKL, OpenBLAS or BLIS that are nested in joblib calls.</p>
<p>Starting from <codeclass="docutils literal notranslate"><spanclass="pre">joblib</span><spanclass="pre">>=</span><spanclass="pre">0.14</span></code>, when the <codeclass="docutils literal notranslate"><spanclass="pre">loky</span></code> backend is used (which
is the default), joblib will tell its child <strong>processes</strong> to limit the
number of threads they can use, so as to avoid oversubscription. In practice
the heuristic that joblib uses is to tell the processes to use <codeclass="docutils literal notranslate"><spanclass="pre">max_threads</span>
<spanclass="pre">=</span><spanclass="pre">n_cpus</span><spanclass="pre">//</span><spanclass="pre">n_jobs</span></code>, via their corresponding environment variable. Back to
our example from above, since the joblib backend of
<aclass="reference internal" href="../modules/generated/sklearn.model_selection.GridSearchCV.html#sklearn.model_selection.GridSearchCV" title="sklearn.model_selection.GridSearchCV"><codeclass="xref py py-class docutils literal notranslate"><spanclass="pre">GridSearchCV</span></code></a> is <codeclass="docutils literal notranslate"><spanclass="pre">loky</span></code>, each process will
only be able to use 1 thread instead of 8, thus mitigating the
oversubscription issue.</p>
<p>Note that:</p>
<ulclass="simple">
<li><p>Manually setting one of the environment variables (<codeclass="docutils literal notranslate"><spanclass="pre">OMP_NUM_THREADS</span></code>,
<codeclass="docutils literal notranslate"><spanclass="pre">MKL_NUM_THREADS</span></code>, <codeclass="docutils literal notranslate"><spanclass="pre">OPENBLAS_NUM_THREADS</span></code>, or <codeclass="docutils literal notranslate"><spanclass="pre">BLIS_NUM_THREADS</span></code>)
will take precedence over what joblib tries to do. The total number of
threads will be <codeclass="docutils literal notranslate"><spanclass="pre">n_jobs</span><spanclass="pre">*</span><spanclass="pre"><LIB>_NUM_THREADS</span></code>. Note that setting this
limit will also impact your computations in the main process, which will
only use <codeclass="docutils literal notranslate"><spanclass="pre"><LIB>_NUM_THREADS</span></code>. Joblib exposes a context manager for
finer control over the number of threads in its workers (see joblib docs
linked below).</p></li>
<li><p>When joblib is configured to use the <codeclass="docutils literal notranslate"><spanclass="pre">threading</span></code> backend, there is no
mechanism to avoid oversubscriptions when calling into parallel native
libraries in the joblib-managed threads.</p></li>
<li><p>All scikit-learn estimators that explicitly rely on OpenMP in their Cython code
always use <codeclass="docutils literal notranslate"><spanclass="pre">threadpoolctl</span></code> internally to automatically adapt the numbers of
threads used by OpenMP and potentially nested BLAS calls so as to avoid
oversubscription.</p></li>
</ul>
<p>You will find additional details about joblib mitigation of oversubscription
in <aclass="reference external" href="https://joblib.readthedocs.io/en/latest/parallel.html#avoiding-over-subscription-of-cpu-resources">joblib documentation</a>.</p>
<p>You will find additional details about parallelism in numerical python libraries
in <aclass="reference external" href="https://thomasjpfan.github.io/parallelism-python-libraries-design/">this document from Thomas J. Fan</a>.</p>