Parallel Python#

Several quite different things go by the name of parallel Python. Distributing work with dask needs nothing special from your environment. Message passing with mpi4py does: the version you get from conda-forge or from PyPI will not run under srun on Levante, and this page shows what to install instead. Writing one file from many ranks needs the same treatment for h5py and netCDF4, see Parallel netCDF and HDF5.

Before you start either, see Using an environment in a batch job for how to launch a Python process from a batch script – in particular, do not put uv run, pixi run or micromamba run behind srun.

Warning

from mpi4py import MPI – not the bare import mpi4py – and any parallel-enabled h5py/netCDF4 built against it, call MPI_Init as a side effect of that import, before your code does anything. That is harmless almost everywhere: a login node, a Jupyter kernel, a plain job allocation. The exception is an interactive allocation entered with salloc or srun --pty: unlike a plain allocation, that starts a live Slurm job step, and SLURM_STEP_ID being set is what triggers this. Open MPI sees that variable, assumes it should attach to the step’s PMIx context, and a process started directly, without srun, was not launched as part of it and aborts:

PMI_Init [pmix_s1.c:171:s1_init]: PMI is not initialized
The application appears to have been direct launched using "srun",
but OMPI was not built with SLURM's PMI support and therefore cannot
execute.

Inside an interactive allocation, go through srun, even for a single task:

srun -n 1 $PY -c "import netCDF4"

A JupyterHub session’s job allocation has no such step – SLURM_JOBID is set, SLURM_STEP_ID is not – so this does not affect Jupyter at all; import netCDF4 in a notebook works the same as on a login node.

Dask#

dask and dask.distributed are ordinary packages, so uv add dask distributed, pixi add dask distributed or micromamba install dask distributed is all the environment needs. What is Levante-specific is how you start the workers. Two older posts still describe the setup:

mpi4py#

Why the packaged versions do not work#

mpi4py is a binding, so it needs an MPI library, and both package repositories bring their own: conda-forge pulls in mpich, and the PyPI wheels have an MPI bundled inside them. Neither speaks the PMIx interface that Levante’s Slurm uses to start MPI processes, so a job launched with srun aborts immediately, once per rank:

Abort(672810256): Fatal error in internal_Init_thread: Internal MPI error!, error stack:
internal_Init_thread(71)........: MPI_Init_thread(...) failed
MPII_Init_thread(204)...........:
MPIR_pmi_init(225)..............:
check_MPIR_CVAR_PMI_VERSION(167):  Runtime environment uses unsupported PMI version PMIx. Aborting.

This is at least a clean failure – you do not silently end up with a set of independent single-rank processes. The way out is to build mpi4py yourself against one of Levante’s own MPI installations, which do speak PMIx.

Building mpi4py against Levante’s OpenMPI#

Load the OpenMPI module first: the build picks up mpicc from your PATH. No compiler module is needed, the system gcc is enough.

In both recipes below, the no-binary part is the essential one: without it you get a pre-built wheel with its own MPI inside, and you are back at the error above. Compiling mpi4py from source takes a couple of minutes.

With uv:

module load openmpi/4.1.2-gcc-11.2.0

uv init --python 3.13 my_project
cd my_project
MPICC=$(which mpicc) uv add --no-binary-package mpi4py mpi4py

To keep that setting with the project rather than on one command line, record it in pyproject.toml, after which a plain uv add mpi4py with the MPI module loaded is enough:

[tool.uv]
no-binary-package = ["mpi4py"]

With pixi, the setting exists only as a manifest entry, so add it to pixi.toml before you add the package:

[pypi-options]
no-binary = ["mpi4py"]
module load openmpi/4.1.2-gcc-11.2.0

pixi init my_project
cd my_project
pixi add python=3.13
# add the [pypi-options] entry to pixi.toml, then
pixi add --pypi mpi4py

micromamba has no way to build a PyPI package into an environment, and conda-forge offers no OpenMPI build here that links against the system library, so use uv or pixi for MPI-parallel Python.

Starting from a ready-made specification#

The specifications from Ready-made environments already contain dask, distributed and dask-jobqueue, so for dask there is nothing to add. As long as the specification you picked does not bring an MPI library of its own – check its package list for mpi4py, mpich or openmpi – mpi4py can be built on top of it:

module load openmpi/4.1.2-gcc-11.2.0
BASE=https://dkrz-sw.gitlab-pages.dkrz.de/python-envs/specs/km-scale

# uv
wget -nc $BASE/pyproject.toml
uv sync
MPICC=$(which mpicc) uv add --no-binary-package mpi4py mpi4py

# pixi
wget -nc $BASE/environment.yaml
pixi init --import environment.yaml
# add the [pypi-options] entry shown above to pixi.toml, then
pixi add --pypi mpi4py

Your project then keeps its own pyproject.toml or pixi.toml and its own lockfile, which no longer match the published specification – that is expected, and the lockfile is what makes your combination reproducible.

Note

This gives you message passing, but not parallel I/O. Writing a single file collectively from several ranks needs h5py and netCDF4 built the same way, see Parallel netCDF and HDF5. The MPI-enabled conda-forge builds of those two packages are not an alternative: they bring conda-forge’s own MPI, which is what the recipe above deliberately avoids.

Checking what you got#

Ask the library which implementation it is, and the extension module where it found it. Both should point at the module you loaded, not into your environment. Run this from a login-node shell; from inside an allocation, see why it needs srun:

uv run python -c "from mpi4py import MPI; print(MPI.Get_library_version())"
# Open MPI v4.1.2, package: Open MPI Distribution, ident: 4.1.2, ...

ldd .venv/lib/python3.13/site-packages/mpi4py/MPI*.so | grep libmpi
# libmpi.so.40 => /sw/spack-levante/openmpi-4.1.2-mnmady/lib/libmpi.so.40

If this says MPICH instead, the no-binary setting did not take effect.

Running it#

Resolve the interpreter once, then let srun start the ranks, as described in Using an environment in a batch job:

cd /home/x/xz0123/my_project

PY=$(uv run which python)       # or: PY=$(pixi run which python)

srun -n 4 $PY my_script.py

Plain srun is enough, you do not need --mpi=pmix. A quick check that the ranks really form one communicator, across more than one node:

srun -N 2 -n 8 $PY -c \
  "from mpi4py import MPI; c=MPI.COMM_WORLD; print(c.Get_rank(), c.Get_size())"

Each rank should report a different rank number and a size of 8. You do not have to load the OpenMPI module inside the job script – the path to the MPI library is baked into the extension module at build time.

That is also the drawback of this approach: your environment is tied to that one MPI installation. If the module is removed or rebuilt, reinstall mpi4py to build it again.

Note

Every rank will print a UCX warning about VM_UNMAP events, mentioning possible performance degradation or data corruption. It is expected, and the recommended OpenMPI settings under MPI Runtime Settings make it go away. Set them in your batch script anyway, they are relevant for performance.

Parallel netCDF and HDF5#

To have several ranks write into one file, h5py and netCDF4 have to be built against Levante’s parallel HDF5 and netCDF-C libraries, using the same MPI as mpi4py. Load the three matching modules – same MPI, same compiler – and build on top of the mpi4py environment from above:

module load openmpi/4.1.2-gcc-11.2.0 \
            hdf5/1.12.1-openmpi-4.1.2-gcc-11.2.0 \
            netcdf-c/4.8.1-openmpi-4.1.2-gcc-11.2.0

The modules do not export a prefix variable, but the libraries can tell you themselves where they are installed:

export NETCDF4_DIR=$(nc-config --prefix)
export HDF5_DIR=$(pkg-config --variable=prefix hdf5)

h5py needs the HDF5_MPI switch, and takes a few minutes to compile. With uv:

CC=mpicc HDF5_MPI=ON uv add --no-binary-package h5py h5py

With pixi, no-binary is a manifest entry, added to the same [pypi-options] block used for mpi4py above:

[pypi-options]
no-binary = ["mpi4py", "h5py"]
CC=mpicc HDF5_MPI=ON pixi add --pypi h5py

netCDF4 needs one thing more. Its build script notices that the netCDF library has parallel functions and then imports mpi4py, which is not visible inside an isolated build environment – so build isolation has to be turned off for this package, and its build requirements provided in the project instead. With uv:

uv add setuptools cython numpy

CC=mpicc USE_NCCONFIG=1 \
  uv add --no-binary-package netcdf4 --no-build-isolation-package netcdf4 netCDF4

With pixi, add no-build-isolation for this one package to the same block:

[pypi-options]
no-binary = ["mpi4py", "h5py", "netcdf4"]
no-build-isolation = ["netcdf4"]
pixi add --pypi setuptools cython numpy
CC=mpicc USE_NCCONFIG=1 pixi add --pypi netCDF4

Check that both really got parallel support, and that they agree on the HDF5 version – if they do not, they found different libraries and will not work together:

uv run python -c "import h5py; print(h5py.get_config().mpi, h5py.version.hdf5_version)"
# or: pixi run python -c "..."
# True 1.12.1

uv run python -c "import netCDF4; print(netCDF4.__has_parallel4_support__, netCDF4.__netcdf4libversion__, netCDF4.__hdf5libversion__)"
# or: pixi run python -c "..."
# 1 4.8.1 1.12.1

Both routes are verified on Levante. If you check this from inside an allocation rather than a login-node shell, run it through srun for the same reason as mpi4py above.

Writing a file collectively#

Open the dataset with parallel=True and pass the communicator. Every rank opens the same file and writes its own part:

# par_write.py
import sys
from mpi4py import MPI
import netCDF4

comm = MPI.COMM_WORLD
ds = netCDF4.Dataset(sys.argv[1], "w", parallel=True, comm=comm, info=MPI.Info())
ds.createDimension("x", comm.Get_size())
v = ds.createVariable("v", "i4", ("x",))
v[comm.Get_rank()] = comm.Get_rank()
ds.close()

Write the output to /work. Parallel I/O goes through MPI-IO, which is meant for the Lustre file system; /home is not a parallel file system. Note that this is the opposite of where the environment itself belongs, see File system, quota and caches.

PY=$(uv run which python)

srun -n 4 $PY par_write.py /work/xx0123/x123456/test_par.nc
ncdump /work/xx0123/x123456/test_par.nc

If every rank contributed, the variable contains one value per rank:

v = 0, 1, 2, 3 ;