Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,9 @@ out = src(
)
```

See [Large problems](https://algorithmiq.github.io/src-method/features/large-problems) for the budgets, the scratch directory and how to read the logged plan.
On a node with several GPUs, `Resources(devices=4)` splits the sweep among four of them, in the same process.

See [Large problems](https://algorithmiq.github.io/src-method/features/large-problems) for the budgets, the scratch directory, several GPUs and how to read the logged plan.

### Logging

Expand Down
18 changes: 16 additions & 2 deletions benches/large/README.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
# Out-of-core benchmark

`bench_large.py` compresses the stack `N . V . M . U` of Pauli-transfer-matrix MPOs
(physical legs of 4) on one GPU, where `M` has a large bond and is read from disk.
(physical legs of 4) on one GPU or several, where `M` has a large bond and is read
from disk.

```bash
# Write the stack to node-local NVMe (205 GB for M at the defaults, complex128).
Expand All @@ -15,9 +16,13 @@ uv run python benches/large/bench_large.py run /local/stack --chi-out 2000 \
uv run python benches/large/bench_large.py generate /local/medium --bond-m 1000
uv run python benches/large/bench_large.py compare /local/medium --chi-out 500 \
--small 40GB --large 80GB --scratch-dir /local/scratch

# Split the sweep among the GPUs of the node: 1, 2, 4 and 8 of them, in turn.
uv run python benches/large/bench_large.py scaling /local/stack --chi-out 2000 \
--devices 1 2 4 8 --scratch-dir /local/scratch --results scaling.jsonl
```

`run` logs the wall time, the size of CuPy's pool at the end (its high-water mark,
`run` logs the wall time, the size of CuPy's pool on each GPU at the end (its high-water mark,
since the pool keeps its blocks), the host peak (`ru_maxrss`), the planned peaks,
the bytes spilled and the tiers. With `--debug` before the command, `src_method` also
logs its plan, the time of each pass and the stall time (`SRC stalls`).
Expand All @@ -31,3 +36,12 @@ memory and node-local NVMe:
budget.
3. The stall time is below 10% of the wall time.
4. `compare` at `D_M = 1000`, `chi_out = 500` reports a distance below `1e-10`.

`scaling` logs, for each GPU count, the wall time, the speed-up and the parallel
efficiency `T_1 / (G T_G)` against the first count, and the relative distance to
its output; `--results` appends the same as JSON lines. `--devices` with
`run` splits a single run. Acceptance of the multi-GPU sweep, at the reference
size on 8xH100 or 8xH200 (and 4xA100 on Leonardo):

1. Every distance is below `1e-10`.
2. The parallel efficiency is at least 80% on all the GPUs of the node.
124 changes: 106 additions & 18 deletions benches/large/bench_large.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,24 @@
"""Out-of-core SRC of ``N . V . M . U`` with a large ``M``, read from disk.

Three commands:
Four commands:

- ``generate``: write the four MPOs site by site as ``.npy`` files, so that ``M``
is never whole in memory.
- ``run``: compress the stack with automatic (or given) budgets and report the
plan, the wall time, the device pool size, the host peak and the stall time.
- ``run``: compress the stack with automatic (or given) budgets, on one or several
GPUs, and report the plan, the wall time, the device pool sizes, the host peak
and the stall time.
- ``compare``: run twice with two GPU budgets and report the relative distance
between the two outputs, which checks that the plan does not change the result.
- ``scaling``: run on 1, 2, 4, ... GPUs of the node and report the wall time, the
parallel efficiency and the distance to the single-GPU output.

The plan, the per-pass times and the stall time are also logged by `src_method`
itself at ``DEBUG``; pass ``--debug`` to see them.
"""

from __future__ import annotations

import json
import logging
import resource
from pathlib import Path # noqa: TC003 (cyclopts reads the annotations at runtime)
Expand Down Expand Up @@ -94,8 +98,12 @@ def _load(directory: Path) -> list[list[np.ndarray]]:


def _run(
directory: Path, chi_out: int, resources: Resources, seed: int
) -> list[np.ndarray]:
directory: Path,
chi_out: int,
resources: Resources,
seed: int,
device: str = "gpu",
) -> tuple[list[np.ndarray], float]:
layers = _load(directory)
plans = []

Expand All @@ -111,29 +119,35 @@ def spy(*args: object, **kwargs: object) -> object:
chi_out=chi_out,
dtype=layers[0][0].dtype,
seed=seed,
device="gpu",
device=device,
resources=resources,
)
seconds = perf_counter() - start
finally:
sweep_module.make_plan = make_plan
import cupy # noqa: PLC0415 (GPU-only benchmark)

(plan,) = plans
tiers = [site.tier for site in plan.sites]
pools = []
if device == "gpu":
import cupy # noqa: PLC0415 (optional dependency)

devices = resources.devices or 1
for gpu in range(devices) if isinstance(devices, int) else devices:
with cupy.cuda.Device(gpu):
pools.append(cupy.get_default_memory_pool().total_bytes())
logger.info(
"Run complete: %.1f s, pool %d B, host peak %d B, planned device peak %d B, "
"Run complete: %.1f s, pools %s B, host peak %d B, planned device peak %d B, "
"planned host peak %d B, disk %d B, tiers %s, sketch batches %s",
seconds,
cupy.get_default_memory_pool().total_bytes(),
pools,
resource.getrusage(resource.RUSAGE_SELF).ru_maxrss * 1024,
plan.device_peak,
plan.host_peak,
plan.disk_bytes,
{tier: tiers.count(tier) for tier in ("device", "host", "disk")},
sorted({site.sketch_batch for site in plan.sites[1:]}),
)
return out
return out, seconds


@app.command
Expand All @@ -144,19 +158,22 @@ def run(
gpu_memory: str | None = None,
host_memory: str | None = None,
scratch_dir: Path | None = None,
devices: int = 1,
seed: int = 0,
) -> None:
"""Compress the stack written by ``generate`` on the GPU.

Args:
directory: The directory given to ``generate``.
chi_out: The output bond dimension.
gpu_memory: GPU budget, e.g. ``36GB``; detected when omitted.
gpu_memory: GPU budget of each device, e.g. ``36GB``; detected when omitted.
host_memory: Host budget; detected when omitted.
scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted.
devices: The number of GPUs to split the sweep among.
seed: Seed of the sketch.
"""
_run(directory, chi_out, Resources(gpu_memory, host_memory, scratch_dir), seed)
resources = Resources(gpu_memory, host_memory, scratch_dir, devices=devices)
_run(directory, chi_out, resources, seed)


def _inner(a: list[np.ndarray], b: list[np.ndarray]) -> complex:
Expand Down Expand Up @@ -187,11 +204,82 @@ def compare(
scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted.
seed: Seed of the sketch, the same for both runs.
"""
first = _run(directory, chi_out, Resources(small, None, scratch_dir), seed)
second = _run(directory, chi_out, Resources(large, None, scratch_dir), seed)
aa, bb, ab = _inner(first, first), _inner(second, second), _inner(first, second)
distance = np.sqrt(max((aa + bb - 2 * ab).real, 0.0) / aa.real)
logger.info("Relative distance between the runs: %.3e", distance)
first, _ = _run(directory, chi_out, Resources(small, None, scratch_dir), seed)
second, _ = _run(directory, chi_out, Resources(large, None, scratch_dir), seed)
logger.info("Relative distance between the runs: %.3e", _distance(first, second))


def _distance(a: list[np.ndarray], b: list[np.ndarray]) -> float:
"""Relative Frobenius distance ``|a - b| / |a|`` of two MPOs."""
aa, bb, ab = _inner(a, a), _inner(b, b), _inner(a, b)
return float(np.sqrt(max((aa + bb - 2 * ab).real, 0.0) / aa.real))


@app.command
def scaling(
directory: Path,
*,
chi_out: int = 2000,
devices: Annotated[tuple[int, ...], cyclopts.Parameter(consume_multiple=True)] = (
1,
2,
4,
8,
),
gpu_memory: str | None = None,
host_memory: str | None = None,
scratch_dir: Path | None = None,
results: Path | None = None,
seed: int = 0,
device: str = "gpu",
) -> None:
"""Run on increasing numbers of GPUs and report the strong scaling.

The first count is the reference: the efficiency of ``G`` GPUs is
``T_ref * G_ref / (T_G * G)``, and every output is compared with the
reference output. Host memory holds two outputs at a time.

Args:
directory: The directory given to ``generate``.
chi_out: The output bond dimension.
devices: The GPU counts, the reference first.
gpu_memory: GPU budget of each device; detected when omitted.
host_memory: Host budget of the node; detected when omitted.
scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted.
results: A JSON-lines file to append one record per run to.
seed: Seed of the sketch, the same for every run.
device: ``gpu``, or ``cpu`` to smoke-test with simulated devices.
"""
reference: list[np.ndarray] | None = None
ref_seconds = ref_devices = 0.0
for count in devices:
resources = Resources(gpu_memory, host_memory, scratch_dir, devices=count)
out, seconds = _run(directory, chi_out, resources, seed, device)
if reference is None:
reference, ref_seconds, ref_devices = out, seconds, count
distance = 0.0
else:
distance = _distance(reference, out)
efficiency = ref_seconds * ref_devices / (seconds * count)
logger.info(
"GPUs %d: %.1f s, speed-up %.2f, efficiency %.0f %%, distance %.3e",
count,
seconds,
ref_seconds / seconds,
100 * efficiency,
distance,
)
if results is not None:
record = {
"gpus": count,
"seconds": seconds,
"efficiency": efficiency,
"distance": distance,
"chi_out": chi_out,
}
with results.open("a") as f:
f.write(json.dumps(record) + "\n")
del out


@app.meta.default
Expand Down
11 changes: 11 additions & 0 deletions docs/content/docs/contributing/testing.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ Always run tests through `pytest`, never as `python tests/test_<name>.py`.
| `test_plan.py` | the planner: budgets, batch sizes and environment tiers |
| `test_store.py` | environment storage on the device, in host memory and on disk |
| `test_sites.py` | reading the cores one site at a time, with prefetching |
| `test_group.py` | the device group: work splitting, collectives, failures |
| `test_logging.py` | the package leaves the host's logging configuration alone |
| `test_gpu_backend.py` | the CuPy path; skipped without CuPy or a GPU |

Expand Down Expand Up @@ -88,6 +89,16 @@ pure function, so `tests/test_plan.py` checks batch sizes and tiers from shapes
alone, and `tests/test_gpu_backend.py` derives GPU budgets from a plan made with
`make_plan` to reach every tier.

## Exercising several devices

`Resources(devices=n)` on the CPU simulates `n` devices with threads that share
the host, through the same code as on GPUs: the split of the sketch columns, the
collectives and the shared reading of the cores. `tests/test_stack.py` compares
such runs with a single device, also with uneven blocks (`chi_out` not a
multiple of `n`, or smaller than it), a cutoff and spilling. Only the peer copies
and device selection are GPU-specific; their tests in `tests/test_gpu_backend.py`
skip unless at least two GPUs are visible.

## Logging in tests

`log_cli` is on, so log records show up live while the tests run. To assert on
Expand Down
4 changes: 4 additions & 0 deletions docs/content/docs/features/gpu.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -27,3 +27,7 @@ result = apply(H, psi, chi_out=256, device="gpu")
The inputs can be NumPy or CuPy arrays; they are moved to the device for the
sweep, and the result always comes back as NumPy arrays on the host. Without
CuPy installed, `device="gpu"` raises an `ImportError`.

To split one sweep among several GPUs of a node, pass
`resources=Resources(devices=n)`; see
[Large problems](/features/large-problems#several-gpus).
62 changes: 59 additions & 3 deletions docs/content/docs/features/large-problems.mdx
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
---
title: Large problems
description: Run stacks that fit neither on the GPU nor in host memory.
description: Run stacks that fit neither on the GPU nor in host memory, on one GPU or several.
---

A single SRC sweep keeps, besides the input and output trains, one sketched
Expand Down Expand Up @@ -75,18 +75,74 @@ reference size of 50 sites, `D_M = 4000` and `chi_out = 2000` in complex128, up
returns, also after an error; a process killed with `SIGKILL` leaves it behind, and
its name identifies the process.

## Several GPUs

On a node with several GPUs, one sweep can be split among them:

```python notest
from src_method import Resources, src

out = src(
N, V, M, U,
chi_out=2000,
dtype=np.complex128,
device="gpu",
resources=Resources(devices=4, scratch_dir="/local/scratch"),
)
```

`devices` is a count, for the first `n` visible GPUs, or a list of device ids;
the first one gathers the sketches and runs the QR. Everything runs in the calling
process, one thread per GPU, so nothing changes for the caller: the call returns
the output cores as with one GPU, errors are raised as usual, and the inputs are
read once and shared.

The sketch index is split: each GPU owns a contiguous block of the `chi_out`
sketch columns and computes only those columns of every environment and of every
sketch, so the environments, the dominant memory, are split among the GPUs too.
At every site of the right-to-left pass, the first GPU gathers the sketch, runs
the QR and gives each GPU a block of rows of the output core; each GPU projects
its rows of the next projected environment, and those rows are gathered on every
GPU. The random draws are those of one GPU, so the result is that of one GPU up to
rounding.

What is not split:

- every GPU holds a whole site: the cores, the projected environment and the
output core of the site, so the largest site that fits is that of one GPU;
- every core is read from the source once, but every GPU needs a copy. When the
GPUs reach each other directly (NVLink, or PCIe peer-to-peer), each uploads a
slice of the core and gathers the rest from the others; otherwise each uploads
the whole core;
- the QR runs on the first GPU while the others wait.

These costs grow more slowly with the problem than the contractions do, so large
problems, the ones this is for, scale close to linearly, while small ones may run
slower on several GPUs than on one.

The budgets apply as follows: `gpu_memory` is the budget of each GPU (detected:
the smallest free memory of the GPUs), while `host_memory` and `scratch_dir` are
shared by all of them. Every GPU writes its environments to its own scratch
directory.

On the CPU, `devices` simulates the GPUs with threads that share the host; the
result is the same, and the mode exists to test the multi-GPU sweep without GPUs.

## Reading the plan

With `DEBUG` logging on for the `src_method` logger (see
[Logging](/features/logging)), every sweep logs its plan as `SRC plan: ...`:

- `prefetch`: whether the next site is read ahead;
- `device peak`, `host peak`, `disk`: the planned peaks, in bytes;
- `devices`: the device ids the sweep is split among;
- `device peak`, `host peak`, `disk`: the planned peaks, in bytes; the device
peak is that of each device, the others are totals;
- `scratch`: the scratch directory, when anything spills to disk;
- `tiers`: where the environment of each site lives (`device`, `host` or `disk`);
- `batches`: for each site, the batch of the environment, sketch and projection
steps.

It then logs the time of each pass and, at the end, `SRC stalls`: the seconds
the sweep waited for input cores and for environments. If they are a large
fraction of the run, the disk is too slow for the compute.
fraction of the run, the disk is too slow for the compute. With several devices,
these lines start with the device they refer to.
Loading
Loading