Skip to content
Merged
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
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ jobs:
test:
name: test (${{ matrix.os }}, py${{ matrix.python }})
runs-on: ${{ matrix.os }}
timeout-minutes: 20 # the suite runs in ~2 min; fail fast instead of hanging 6 h
strategy:
fail-fast: false
matrix:
Expand Down
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@ First public release.
- Batch sinks are correct across re-runs: a file sink is keyed by `(key, method)`, a
directory sink refuses a mismatched `chunk_size` on resume, and a CSV sink asked to
store raw spectra fails loudly rather than dropping the columns.
- The batch process pool uses the `spawn` start method on every platform, so a CPU/GPU
pool no longer deadlocks on Linux (the default `fork` copies parent native thread pools
/ CUDA contexts into the workers).

### Validated

Expand Down
9 changes: 8 additions & 1 deletion src/cuperiod/batch/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
from __future__ import annotations

import json
import multiprocessing
import warnings
from collections.abc import Mapping, Sequence
from concurrent.futures import ProcessPoolExecutor, as_completed
Expand Down Expand Up @@ -401,8 +402,14 @@ def _run_pool(
initializer: Any = None,
initargs: tuple[Any, ...] = (),
) -> None:
# Always use "spawn". Linux's default "fork" copies the parent's already-built
# native thread pools (numba / OpenBLAS / OpenMP, plus any CUDA context for the
# GPU pool) into the child and deadlocks the workers. Spawn starts fresh,
# thread-pinned workers (the Windows/macOS default) — see pin_worker_threads.
ctx = multiprocessing.get_context("spawn")
with ProcessPoolExecutor(
max_workers=max_workers, initializer=initializer, initargs=initargs
max_workers=max_workers, mp_context=ctx,
initializer=initializer, initargs=initargs,
) as pool:
futures = {
pool.submit(_process_chunk, chunks[idx], cfg): idx for idx in pending
Expand Down
Loading