mhcflurry.parallelism package

Local parallelism helpers.

This package contains the implementation formerly held in mhcflurry.local_parallelism. The old module remains as an import compatibility shim.

class mhcflurry.parallelism.NonDaemonContext[source]

Bases: ForkContext

A multiprocessing context that hands out NonDaemonProcess workers.

Subclasses the current default multiprocessing context so its start method is preserved — we only swap the Process class. The Pool uses self._ctx.Process(...) to create workers and will now get our non-daemonic variant.

Process

alias of NonDaemonProcess

class mhcflurry.parallelism.NonDaemonPool(*args, **kwargs)[source]

Bases: Pool

A multiprocessing.Pool that runs non-daemonic workers.

Pool’s constructor takes a context kwarg — we thread a NonDaemonContext through so each worker is a NonDaemonProcess. Everything else (apply_async, imap, etc.) inherits unchanged.

class mhcflurry.parallelism.NonDaemonProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, Process

A multiprocessing.Process whose daemon flag cannot be set.

Reading .daemon always returns False; writes are no-ops. This lets us instantiate multiprocessing.pool.Pool with a worker class that declines to be a daemon, so the DataLoader inside each worker can spawn its own prefetch children.

class mhcflurry.parallelism.NonDaemonSpawnContext[source]

Bases: SpawnContext

Process

alias of NonDaemonSpawnProcess

class mhcflurry.parallelism.NonDaemonSpawnProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, SpawnProcess

exception mhcflurry.parallelism.WrapException[source]

Bases: Exception

Add traceback info to exception so exceptions raised in worker processes can still show traceback info when re-raised in the parent.

mhcflurry.parallelism.add_local_parallelism_args(parser)[source]

Add local parallelism arguments to the given argparse.ArgumentParser.

Parameters:
parserargparse.ArgumentParser
mhcflurry.parallelism.add_prediction_parallelism_args(parser)[source]

Add prediction-time local parallelism arguments to an argparse parser.

This is the inference subset of add_local_parallelism_args: the worker scheduler, backend selection, and torch forward-kernel knobs, without training-only DataLoader/random-negative options.

mhcflurry.parallelism.apply_dataloader_num_workers_to_work_items(work_items, num_workers, *, log=None)[source]

Inject dataloader_num_workers into every work item’s hyperparameters.

Writes the resolved value into affinity work-item hyperparameters. It is consumed by streaming pretraining; in-memory affinity fitting ignores it. Processing models do not use this affinity DataLoader hyperparameter.

Parameters:
work_itemslist of dict

Each dict has a hyperparameters sub-dict (the canonical shape produced by train_pan_allele_models_command / train_allele_specific_models_command).

num_workersint

The resolved value from resolve_local_parallelism_args. Use 0 to force in-process batching (no prefetch children).

logcallable, optional

Logging hook for the human-readable summary. Defaults to print.

mhcflurry.parallelism.apply_random_negative_pool_epochs_to_work_items(work_items, pool_epochs, *, log=None)[source]

Inject random_negative_pool_epochs into every work item’s hyperparameters.

Parallel to apply_dataloader_num_workers_to_work_items: the orchestrator chooses an auto value once at startup (or honors the CLI int) and writes it into every per-work-item hyperparameter dict. fit() reads it from self.hyperparameters['random_negative_pool_epochs'] when constructing its RandomNegativesPool.

Parameters:
work_itemslist of dict

Each dict has a hyperparameters sub-dict.

pool_epochsint

The resolved value from resolve_local_parallelism_args (>= 1).

logcallable, optional

Logging hook. Defaults to print.

mhcflurry.parallelism.apply_resolved_training_hyperparameters_to_work_items(work_items, args, *, log=None)[source]

Inject resolved per-model training knobs into work item hyperparameters.

resolve_local_parallelism_args owns hardware-dependent CLI resolution. Trainers should call this once after constructing work items so pan-allele, allele-specific, and future affinity trainers all persist the same resolved settings in component-model hyperparameters.

mhcflurry.parallelism.attach_constant_data_to_work_items_if_needed(work_items, constant_data, worker_pool, *, log=None)[source]

Attach constant data only when the Pool cannot inherit it by fork.

mhcflurry.parallelism.auto_dataloader_num_workers(num_fit_workers, vcpus=None, ram_gb=None, hard_cap=None)[source]

Pick per-fit-worker DataLoader child count from box capacity.

Returns an int >= 0. The result is the value that should be plugged into each component model’s dataloader_num_workers hyperparameter. A return of 0 means in-process batching (no prefetch children), which is correct for serial runs and very tight CPU configs.

Parameters:
num_fit_workersint

Total fit() worker processes that will share the box. Equal to num_gpus * max_workers_per_gpu for the canonical GPU run, or num_jobs for CPU-only runs.

vcpusint, optional

Total vCPU count. Default os.cpu_count().

ram_gbfloat, optional

Total system RAM in GB. None skips the RAM cap (CPU-only decision).

hard_capint, optional

Maximum DL children per fit-worker. Default 4 (overridable via MHCFLURRY_AUTO_DATALOADER_HARD_CAP); beyond this, process/queue overhead reduced measured throughput.

Notes

  1. Serial / no GPU work: if num_fit_workers <= 0, return 0. The caller will run fit() in-process; spawning children would buy nothing and cost a process-fork.

  2. CPU budget per fit-worker: cpu_per_fit = vcpus // num_fit_workers. Each DL child needs ~``_AUTO_DATALOADER_CORES_PER_CHILD`` (=2) physical cores to keep up with fancy-indexing + collate without starving the main fit-worker’s OMP/MKL pool.

  3. CPU cap: cpu_cap = cpu_per_fit // 2.

  4. RAM cap (when ram_gb is provided): each DL child holds ~``_AUTO_DATALOADER_RAM_PER_CHILD_GB`` (=0.5) GB of RSS for the torch+mhcflurry imports; the main fit-worker baseline is ~``_AUTO_DATALOADER_RAM_BASELINE_PER_FIT_GB`` (=2.0) GB. ram_cap = max(0, (ram_per_fit_gb - 2.0) / 0.5).

  5. Throughput cap: min(cpu_cap, ram_cap, hard_cap).

  6. Floor: return 0 when any CPU, RAM or explicit hard cap is 0. Otherwise the minimum is 1.

Edge cases

  • num_fit_workers > vcpus → cpu_per_fit = 0, cpu_cap = 0, result 0. The main fit-workers themselves are oversubscribed; adding DL children would make it worse.

  • ram_gb very small (e.g. < 2 GB / fit) → ram_cap = 0, falls back to in-process batching to preserve correctness over throughput.

  • hard_cap env override of 0 → forces in-process for diagnostics.

Cross-checks (see test_orchestrator_helpers.py)

  • 8×A100-80GB Verda (176v / 16 fit / 400G) → 4

  • 8×A100-40GB (176v / 8 fit / 400G) → 4

  • 8×L40S (96v / 16 fit / 200G) → 3

  • Single A100 80G Lambda (30v / 2 fit / 200G) → 4

  • Single A100 80G tight (16v / 2 fit / 64G) → 4

  • Single T4 (8v / 1 fit / 16G) → 4

  • CPU 8-thread (8v / 0 fit) → 0

  • Tight cluster node (32v / 16 fit / 64G) → 1

  • RAM-starved (176v / 16 fit / 32G) → 0

mhcflurry.parallelism.auto_max_workers_per_gpu(num_jobs, num_gpus, backend='auto', per_worker_gb=None)[source]

Pick max_workers_per_gpu based on detected hardware.

Returns an int ≥ 1. Logic:

  • num_gpus == 0 (CPU-only) → 1.

  • Otherwise: take the minimum of two capacity limits:
    • num_jobs // num_gpus — don’t oversubscribe a GPU beyond the jobs that actually exist.

    • complete per-worker working sets that fit in free VRAM after the shared allocator/context reserve.

Free VRAM is read from nvidia-smi (or from MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB). It deliberately avoids torch.cuda so resolving local parallelism before forking does not initialize CUDA in the parent process. Per-worker VRAM upper bound and the optional expert hard cap are overridable via MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_PER_WORKER_GB and MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_HARD_CAP.

Calibrate’s per-worker footprint is dominated by the cached_stages tensor, so the planner passes per_worker_gb explicitly — the affinity calibration workload profile’s device_worker_gb of 24 GB — overriding the 4 GB train default. When given, the explicit hint wins over the env var. (No env var for it: the workload-specific knowledge belongs in the workload profile, not in a global env.)

The result is logged so the chosen value is visible in the worker log alongside the reasoning.

mhcflurry.parallelism.auto_num_jobs(num_gpus, max_workers_per_gpu)[source]

Compute total fit-worker count from GPU plan.

Returns num_gpus * max_workers_per_gpu for GPU runs, 0 when no GPUs are visible (caller decides serial vs explicit CPU pool). Treats "auto" max_workers_per_gpu as not-yet-resolved and raises; callers must resolve it first via auto_max_workers_per_gpu.

mhcflurry.parallelism.auto_random_negative_pool_epochs(num_random_negatives, peptide_max_length, num_workers, ram_gb=None, *, safety_fraction=None, per_pool_epoch_per_worker_bytes=None, hard_cap=None, base_worker_gb=0.0)[source]

Pick random_negative_pool_epochs from box capacity.

The pool sits in the heap of every fit-worker process. Per-pool-epoch memory cost is dominated by:

  • num_random_negatives × peptide_max_length int8 indices,

  • intermediate pandas.Series[str] allocations and encoder buffers observed in practice at ~10–40× the int8 size on 2026-04 measurements (the old pool_epochs=100 run OOM’d a 944 GB box at ~199 GB / worker worth of transient pool cost).

We budget for the pessimistic per_pool_epoch_per_worker_bytes figure so the auto value is safe under transient peaks; tunable via env when a workload is known to behave better than the empirical pessimistic figure.

Returns an int >= 1. 1 means fresh random negatives every epoch. > 1 amortizes random-negative generation + encoding across N epochs.

Parameters:
num_random_negativesint

Number of random-negative peptides per epoch (the planner’s get_total_count()). The size of one pool-epoch in the cycle.

peptide_max_lengthint

Longest peptide the encoding allocates space for. With the fixed-vector encoding lookup, the per-peptide footprint is peptide_max_length int8 bytes.

num_workersint

Total fit() worker processes that will share the box. Each holds its own RN pool, so total RAM cost = num_workers × pool_epochs × per_pool_epoch_per_worker_bytes.

ram_gbfloat, optional

Total system RAM in GB. None returns 1 (safe default — we don’t know the budget so don’t bump pool_epochs above legacy).

safety_fractionfloat, optional

Expert fraction of RAM available to RN pools across all workers. By default the shared host reserve is used and base_worker_gb is subtracted first. Override via MHCFLURRY_AUTO_RN_POOL_SAFETY_FRACTION.

per_pool_epoch_per_worker_bytesfloat, optional

Empirical per-pool-epoch per-worker RAM cost in bytes. Default 1 GB (conservative — captures the int8 indices + transient pandas + encoder buffers seen on the 2026-04 run). Override via MHCFLURRY_AUTO_RN_POOL_PER_EPOCH_PER_WORKER_GB.

hard_capint, optional

Maximum pool epochs. Default 10 (expert-overridable); larger pools add startup/memory cost after generation overhead is already amortized.

base_worker_gbfloat, optional

Baseline host memory per fit worker, subtracted before sizing pools.

Notes

Total available bytes for RN pools across the box is usable host memory after the shared reserve and fit-worker base RSS. An explicit safety_fraction replaces that calculation. Per-worker budget is available / max(num_workers, 1). The number of pool epochs that fit is per_worker_budget / per_pool_epoch_per_worker_bytes, clamped to [1, hard_cap].

Cross-checks

  • 8×A100-80GB Verda shares the same reserve and worker-RSS estimate as the outer process planner, then stops at the throughput cap.

  • Single A100 80G Lambda uses the shared host reserve minus the two workers’ base RSS, then fills the remaining per-worker pool budget.

  • Tight and RAM-starved nodes naturally fall back toward one pooled epoch as the fit-worker base RSS consumes the shared usable budget.

mhcflurry.parallelism.call_wrapped(function, *args, **kwargs)[source]

Run function on args and kwargs and return result, wrapping any exception raised in a WrapException.

Parameters:
functionarbitrary function
Any other arguments provided are passed to the function.
Returns:
object
mhcflurry.parallelism.call_wrapped_kwargs(function, kwargs)[source]

Invoke function on given kwargs and return result, wrapping any exception raised in a WrapException.

Parameters:
functionarbitrary function
kwargsdict
Returns:
object
Result of calling function(**kwargs).
mhcflurry.parallelism.chunk_ranges_for_local_parallelism(num_items, num_jobs=0, chunks_per_worker=4)[source]

Split a row/sequence axis into stable contiguous chunks for local workers.

Parameters:
num_itemsint

Number of input items.

num_jobsint

Number of worker processes. 0 yields one serial chunk.

chunks_per_workerint

Target number of work chunks per worker for load balancing.

Returns:
list of tuple

(chunk_index, start, end) ranges.

mhcflurry.parallelism.configure_cluster_worker_torch_compile_threads()[source]

Auto-size Inductor helper threads inside one cluster worker process.

Cluster parallelism submits each work item as its own process, often on different nodes. We therefore do not try to share a compile cache across the cluster. Each worker process still needs the same local policy: if compile is enabled and TORCHINDUCTOR_COMPILE_THREADS is unset or auto, pick a numeric value on that machine before the first torch.compile call.

If a scheduler packs several mhcflurry work items onto one node, set MHCFLURRY_CLUSTER_WORKERS_PER_NODE so the auto value is divided across those co-resident compiler processes. Otherwise the default assumes one work process owns its scheduler CPU allocation.

mhcflurry.parallelism.detect_free_vram_per_gpu_gb(num_gpus)[source]

Per-GPU free VRAM (GB) as a list, or None if undetectable.

Env override (MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB) first, then nvidia-smi. Unlike the scalar helpers used by the worker-count math, this preserves the per-GPU values so capacity warnings can flag small or uneven cards.

mhcflurry.parallelism.estimate_worker_context_bytes(value, _seen=None)[source]

Best-effort deep resident size for data copied into spawn workers.

mhcflurry.parallelism.free_vram_per_gpu_from_nvidia_smi_gb(num_gpus)[source]

Return a list of free VRAM (GB) per visible GPU using nvidia-smi.

Returns None if nvidia-smi is unavailable or returns nothing. The per-GPU list (not collapsed to a scalar) lets capacity diagnostics see heterogeneous / partially-occupied cards; the worker-count math separately takes the min via free_vram_from_nvidia_smi_gb.

This is intentionally a subprocess call instead of torch.cuda so the orchestrator can size a fork-based worker pool without initializing CUDA in the parent process.

mhcflurry.parallelism.hoist_torchinductor_compile_threads(args, phase='production')[source]

Auto-size TORCHINDUCTOR_COMPILE_THREADS for local training.

torch.compile (when enabled via MHCFLURRY_TORCH_COMPILE=1) spins up an inductor compile worker pool that defaults to os.cpu_count() threads. With N fit() workers each running their own compile pool, that multiplies into an oversubscribed compile storm. The production phase uses an auto value derived from available cores and the worker count; the warmup phase uses a larger value because only one worker is compiling.

The orchestrator owns “how many workers will exist”, so it owns the env knob too: set once before forking, every worker inherits. Skips the hoist when the user has already pinned the value or when MHCFLURRY_TORCH_COMPILE isn’t on. Cluster workers running on other hosts must size themselves locally; see configure_cluster_worker_torch_compile_threads.

Lives here (not in any one train_*_command module) so processing, allele-specific, and any future train command can call it the same way.

mhcflurry.parallelism.make_worker_pool(processes=None, initializer=None, initializer_kwargs_per_process=None, initializer_shared_kwargs=None, max_tasks_per_worker=None, start_method=None)[source]

Convenience wrapper to create a multiprocessing.Pool.

This function adds support for per-worker initializer arguments, which are not natively supported by the multiprocessing module. The motivation for this feature is to support allocating each worker to a (different) GPU.

IMPLEMENTATION NOTE:

The per-worker initializer arguments are implemented using a SimpleQueue. Each worker reads its arguments from this queue when it starts. When it terminates, it adds its initializer arguments back to the queue, so a future process can initialize itself using these arguments. SimpleQueue is important here: Queue.put uses a feeder thread, so workers can observe a transiently empty queue during startup and duplicate GPU assignments. A worker that puts its arguments back can also hang forever joining that feeder thread during process finalization.

There is one issue with this approach, however. If a worker crashes, it never repopulates the queue of initializer arguments. This will prevent any future worker from re-using those arguments. To deal with this issue we add a second ‘backup queue’. This queue always contains the full set of initializer arguments: whenever a worker reads from it, it always pushes the pop’d args back to the end of the queue immediately. If the primary arg queue is ever empty, then workers will read from this backup queue.

Parameters:
processesint

Number of workers. Default: num CPUs.

initializerfunction, optional

Init function to call in each worker

initializer_kwargs_per_processlist of dict, optional

Arguments to pass to initializer function for each worker. Length of list must equal the number of workers.

initializer_shared_kwargsdict, optional

Arguments passed once to every worker initializer. Unlike work-item arguments, large values here are serialized only once per process.

max_tasks_per_workerint, optional

Restart workers after this many tasks.

start_methodstring, optional

Multiprocessing start method to use for the worker pool.

Returns:
multiprocessing.Pool
mhcflurry.parallelism.non_daemon_context(start_method=None)[source]

Return a multiprocessing context whose workers are non-daemonic.

mhcflurry.parallelism.num_workers_per_gpu_from_args(args)[source]

Return resolved max_workers_per_gpu for model auto-sizing.

Callers must run resolve_local_parallelism_args first so workload- and hardware-aware defaults have already converted the CLI sentinel "auto" into an integer.

mhcflurry.parallelism.refine_local_parallelism_from_spawn_context(args, context_bytes, memory=None)[source]

Cap automatic spawn workers using the loaded per-process context.

Spawn serializes the full worker initializer context into every process, unlike fork’s copy-on-write inheritance. Input files may be compressed, so their on-disk size is not a reliable host-RAM estimate. This refinement is deliberately performed after the dataset/model context has been loaded but before any workers start.

mhcflurry.parallelism.refine_local_parallelism_from_warmup(args, reports)[source]

Tighten an automatic plan using representative warmup measurements.

The analytic estimate remains the floor because a one-batch warmup does not contain a full resident dataset or validation. Measured CUDA reserved memory and process peak RSS can only increase the estimated working set and reduce automatic concurrency. Explicit --num-jobs and --max-workers-per-gpu values are never changed.

mhcflurry.parallelism.refine_local_parallelism_from_worker_context(args, worker_context_data, start_method=None)[source]

Apply spawn-context host sizing before code chooses serial/parallel.

mhcflurry.parallelism.refresh_device_memory_budget(args)[source]

Snapshot and propagate one fixed launch-time entitlement per worker.

mhcflurry.parallelism.resolve_cpu_thread_budget(plan, cpu_count=None)[source]

Resolve native threads per fit worker and environment ownership.

A uniform runtime resize is safe only when all supported environment variables are unset or were written by an earlier mhcflurry auto pass. Any caller-owned value makes the environment authoritative; the numeric auto estimate is still recorded for diagnostics but must not be applied to loaded native or PyTorch pools.

mhcflurry.parallelism.resolve_cpu_threads_per_worker(plan, cpu_count=None)[source]

Set unset BLAS/OpenMP thread env vars from the final worker plan.

mhcflurry.parallelism.resolve_dataloader_num_workers(value, num_fit_workers=None, vcpus=None, ram_gb=None)[source]

Normalize a dataloader_num_workers value to an int.

Accepts "auto" / unset / None (delegates to auto_dataloader_num_workers), or any int-coercible value. Used by both the orchestrator-side resolver and the shell helper that injects dataloader_num_workers into the recipe’s hyperparameters.yaml.

mhcflurry.parallelism.resolve_local_parallelism_args(args, cap_auto_num_jobs=True, per_worker_gb=None, workload_name='generic', workload_hints=None)[source]

Resolve and normalize local parallelism arguments through the planner.

mhcflurry.parallelism.resolve_max_workers_per_gpu(args, per_worker_gb=None, num_gpus=None, backend=None)[source]

Resolve args.max_workers_per_gpu to an int, mutating args.

Accepts the literal string "auto" (the default) or an int. When "auto", calls auto_max_workers_per_gpu with the rest of the args’ parallelism config to pick a value. Idempotent — calling twice on the same args is a no-op the second time.

per_worker_gb lets workload-specific commands (e.g. calibrate, where cached_stages dominates per-worker VRAM at ~15 GB) override the train-default. Falls back to env var / the module default when not given.

Returns the resolved int (also stored on args.max_workers_per_gpu so subsequent consumers see the int).

mhcflurry.parallelism.resolve_torchinductor_compile_threads_env(num_jobs=1, phase='production')[source]

Ensure TORCHINDUCTOR_COMPILE_THREADS is parseable by PyTorch.

MHCflurry launchers accept TORCHINDUCTOR_COMPILE_THREADS=auto as an orchestrator-owned sentinel so each command can size the Inductor compiler helper pool after it knows its local worker count. PyTorch itself does not accept that sentinel: Inductor parses the env var with int(...).

This helper is for entry points that do not run the full local-parallelism resolver before importing or spawning PyTorch work. It leaves unset and user-pinned integer values alone, resolves orchestrator-owned auto to a numeric value, and fails early on other invalid values with a clearer message than the downstream Inductor traceback.

mhcflurry.parallelism.run_single_worker_resource_probe(args, work_items, work_function, constant_data=None)[source]

Measure one real peak phase for every resource-distinct architecture.

The command’s resource_probe_only path must exercise its true resident data, configured minibatch, and validation phase for one bounded epoch. The single spawned worker records process-level CUDA memory (including the context), allocator peaks, and host RSS. Those observations may only tighten automatic concurrency before the production pool starts.

When torch.compile is enabled the same pass also primes its on-disk cache; resource safety is intentionally independent of compile being enabled.

work_items is not mutated: every task still runs in production.

Returns None when skipped, otherwise the number of unique architectures warmed.

mhcflurry.parallelism.run_single_worker_torch_compile_warmup(args, work_items, work_function, constant_data=None)[source]

Compatibility alias for run_single_worker_resource_probe().

mhcflurry.parallelism.validate_worker_pool_args(num_jobs, num_gpus=0, backend='auto', max_workers_per_gpu=1)[source]

Validate local worker scheduling arguments.

--gpus controls CUDA worker assignment only. It does not select MPS devices and it does not distribute a single model across multiple GPUs.

mhcflurry.parallelism.worker_init(keras_backend=None, backend=None, gpu_device_nums=None, worker_log_dir=None, max_workers_per_gpu=None, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, worker_context_module=None, worker_context_data=None, device_memory_budget_bytes=None)[source]
mhcflurry.parallelism.worker_init_entry_point(init_function, kwargs_per_process=None, slots=None, sequence=None, shared_kwargs=None)[source]
mhcflurry.parallelism.worker_init_kwargs_for_scheduler(num_jobs, num_gpus=0, backend='auto', max_workers_per_gpu=1, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, device_memory_budget_bytes=None)[source]

Build per-worker init kwargs from the local scheduling configuration.

When num_gpus is set, workers are assigned one CUDA GPU each in round robin order. Any additional workers are forced onto CPU by hiding CUDA and setting their backend to cpu.

mhcflurry.parallelism.worker_pool_uses_fork(worker_pool=None)[source]

Return True when local Pool workers inherit parent globals by fork.

mhcflurry.parallelism.worker_pool_with_gpu_assignments(num_jobs, num_gpus=0, backend=None, max_workers_per_gpu=1, max_tasks_per_worker=None, worker_log_dir=None, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, start_method=None, worker_context_module=None, worker_context_data=None, device_memory_budget_bytes=None)[source]

Create a multiprocessing.Pool where each worker uses its own GPU.

Parameters:
num_jobsint

Number of worker processes.

num_gpusint
backendstring
max_workers_per_gpuint
max_tasks_per_workerint
worker_log_dirstring
cpu_threads_per_workerint

Runtime BLAS/OpenMP/PyTorch thread limit applied in each worker.

cpu_threads_per_worker_was_autobool

Whether mhcflurry owns the uniform runtime limit. False preserves caller-provided OMP/MKL/OpenBLAS settings.

start_methodstring

Optional multiprocessing start method, e.g. "spawn" when workers must not inherit PyTorch state from the parent process.

worker_context_modulestring, optional

Module containing a WORKER_CONTEXT dictionary used by worker functions.

worker_context_datadict, optional

Constant data to install once per process during initialization. This avoids serializing the same large payload with every queued task on spawn-based platforms.

device_memory_budget_bytesint, optional

Fixed launch-time device-memory entitlement propagated to each GPU worker for elastic batch sizing. Kept last in the signature so adding the planner-owned value does not change existing positional calls.

Returns:
multiprocessing.Pool
mhcflurry.parallelism.worker_pool_with_gpu_assignments_from_args(args, workload_name='generic', workload_hints=None, start_method=None, worker_context_module=None, worker_context_data=None)[source]

Create a multiprocessing.Pool where each worker uses its own GPU.

Uses commandline arguments. See worker_pool_with_gpu_assignments. Resolves args.max_workers_per_gpu="auto" to an int (mutating args so downstream consumers — e.g. inference batch sizing in calibrate — observe the same value).

Parameters:
argsargparse.Namespace

Parsed local-parallelism options.

workload_namestr

Workload profile used for automatic sizing.

workload_hintsdict, optional

Model/data size hints for that profile.

start_methodstr, optional

Multiprocessing start method.

worker_context_modulestr, optional

Module whose WORKER_CONTEXT receives constant data.

worker_context_datadict, optional

Constant data installed once per worker.

Returns:
multiprocessing.Pool

Submodules

mhcflurry.parallelism.cli_args module

Argument-parser helpers for local parallelism.

mhcflurry.parallelism.cli_args.add_local_parallelism_args(parser)[source]

Add local parallelism arguments to the given argparse.ArgumentParser.

Parameters:
parserargparse.ArgumentParser
mhcflurry.parallelism.cli_args.add_prediction_parallelism_args(parser)[source]

Add prediction-time local parallelism arguments to an argparse parser.

This is the inference subset of add_local_parallelism_args: the worker scheduler, backend selection, and torch forward-kernel knobs, without training-only DataLoader/random-negative options.

mhcflurry.parallelism.planning module

Hardware planning helpers for local parallelism.

mhcflurry.parallelism.planning.cuda_visible_devices_from_env()[source]

Return CUDA_VISIBLE_DEVICES entries, or None when unset.

mhcflurry.parallelism.planning.free_vram_per_gpu_override_gb(num_gpus)[source]

Return env-pinned free VRAM in GB as a per-GPU list, or None.

MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB accepts either one value applied to every GPU or a comma/space-separated list. It exists for tests and unusual launchers where nvidia-smi is hidden but the caller knows the device budget.

mhcflurry.parallelism.planning.free_vram_override_gb(num_gpus)[source]

Minimum env-pinned free VRAM in GB, or None.

mhcflurry.parallelism.planning.detect_num_cuda_devices_no_torch()[source]

Return the number of visible CUDA devices without importing torch.

CUDA_VISIBLE_DEVICES is authoritative when set by a scheduler or container. Otherwise shell out so the orchestrator can size a fork-based worker pool without initializing CUDA in the parent process. Returns 0 when nvidia-smi is unavailable or no GPUs are visible.

mhcflurry.parallelism.planning.free_vram_per_gpu_from_nvidia_smi_gb(num_gpus)[source]

Return a list of free VRAM (GB) per visible GPU using nvidia-smi.

Returns None if nvidia-smi is unavailable or returns nothing. The per-GPU list (not collapsed to a scalar) lets capacity diagnostics see heterogeneous / partially-occupied cards; the worker-count math separately takes the min via free_vram_from_nvidia_smi_gb.

This is intentionally a subprocess call instead of torch.cuda so the orchestrator can size a fork-based worker pool without initializing CUDA in the parent process.

mhcflurry.parallelism.planning.free_vram_from_nvidia_smi_gb(num_gpus)[source]

Minimum free VRAM (GB) across visible GPUs, or None.

mhcflurry.parallelism.planning.detect_free_vram_per_gpu_gb(num_gpus)[source]

Per-GPU free VRAM (GB) as a list, or None if undetectable.

Env override (MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB) first, then nvidia-smi. Unlike the scalar helpers used by the worker-count math, this preserves the per-GPU values so capacity warnings can flag small or uneven cards.

mhcflurry.parallelism.planning.auto_max_workers_per_gpu(num_jobs, num_gpus, backend='auto', per_worker_gb=None)[source]

Pick max_workers_per_gpu based on detected hardware.

Returns an int ≥ 1. Logic:

  • num_gpus == 0 (CPU-only) → 1.

  • Otherwise: take the minimum of two capacity limits:
    • num_jobs // num_gpus — don’t oversubscribe a GPU beyond the jobs that actually exist.

    • complete per-worker working sets that fit in free VRAM after the shared allocator/context reserve.

Free VRAM is read from nvidia-smi (or from MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB). It deliberately avoids torch.cuda so resolving local parallelism before forking does not initialize CUDA in the parent process. Per-worker VRAM upper bound and the optional expert hard cap are overridable via MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_PER_WORKER_GB and MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_HARD_CAP.

Calibrate’s per-worker footprint is dominated by the cached_stages tensor, so the planner passes per_worker_gb explicitly — the affinity calibration workload profile’s device_worker_gb of 24 GB — overriding the 4 GB train default. When given, the explicit hint wins over the env var. (No env var for it: the workload-specific knowledge belongs in the workload profile, not in a global env.)

The result is logged so the chosen value is visible in the worker log alongside the reasoning.

mhcflurry.parallelism.planning.auto_dataloader_num_workers(num_fit_workers, vcpus=None, ram_gb=None, hard_cap=None)[source]

Pick per-fit-worker DataLoader child count from box capacity.

Returns an int >= 0. The result is the value that should be plugged into each component model’s dataloader_num_workers hyperparameter. A return of 0 means in-process batching (no prefetch children), which is correct for serial runs and very tight CPU configs.

Parameters:
num_fit_workersint

Total fit() worker processes that will share the box. Equal to num_gpus * max_workers_per_gpu for the canonical GPU run, or num_jobs for CPU-only runs.

vcpusint, optional

Total vCPU count. Default os.cpu_count().

ram_gbfloat, optional

Total system RAM in GB. None skips the RAM cap (CPU-only decision).

hard_capint, optional

Maximum DL children per fit-worker. Default 4 (overridable via MHCFLURRY_AUTO_DATALOADER_HARD_CAP); beyond this, process/queue overhead reduced measured throughput.

Notes

  1. Serial / no GPU work: if num_fit_workers <= 0, return 0. The caller will run fit() in-process; spawning children would buy nothing and cost a process-fork.

  2. CPU budget per fit-worker: cpu_per_fit = vcpus // num_fit_workers. Each DL child needs ~``_AUTO_DATALOADER_CORES_PER_CHILD`` (=2) physical cores to keep up with fancy-indexing + collate without starving the main fit-worker’s OMP/MKL pool.

  3. CPU cap: cpu_cap = cpu_per_fit // 2.

  4. RAM cap (when ram_gb is provided): each DL child holds ~``_AUTO_DATALOADER_RAM_PER_CHILD_GB`` (=0.5) GB of RSS for the torch+mhcflurry imports; the main fit-worker baseline is ~``_AUTO_DATALOADER_RAM_BASELINE_PER_FIT_GB`` (=2.0) GB. ram_cap = max(0, (ram_per_fit_gb - 2.0) / 0.5).

  5. Throughput cap: min(cpu_cap, ram_cap, hard_cap).

  6. Floor: return 0 when any CPU, RAM or explicit hard cap is 0. Otherwise the minimum is 1.

Edge cases

  • num_fit_workers > vcpus → cpu_per_fit = 0, cpu_cap = 0, result 0. The main fit-workers themselves are oversubscribed; adding DL children would make it worse.

  • ram_gb very small (e.g. < 2 GB / fit) → ram_cap = 0, falls back to in-process batching to preserve correctness over throughput.

  • hard_cap env override of 0 → forces in-process for diagnostics.

Cross-checks (see test_orchestrator_helpers.py)

  • 8×A100-80GB Verda (176v / 16 fit / 400G) → 4

  • 8×A100-40GB (176v / 8 fit / 400G) → 4

  • 8×L40S (96v / 16 fit / 200G) → 3

  • Single A100 80G Lambda (30v / 2 fit / 200G) → 4

  • Single A100 80G tight (16v / 2 fit / 64G) → 4

  • Single T4 (8v / 1 fit / 16G) → 4

  • CPU 8-thread (8v / 0 fit) → 0

  • Tight cluster node (32v / 16 fit / 64G) → 1

  • RAM-starved (176v / 16 fit / 32G) → 0

mhcflurry.parallelism.planning.resolve_dataloader_num_workers(value, num_fit_workers=None, vcpus=None, ram_gb=None)[source]

Normalize a dataloader_num_workers value to an int.

Accepts "auto" / unset / None (delegates to auto_dataloader_num_workers), or any int-coercible value. Used by both the orchestrator-side resolver and the shell helper that injects dataloader_num_workers into the recipe’s hyperparameters.yaml.

mhcflurry.parallelism.planning.auto_random_negative_pool_epochs(num_random_negatives, peptide_max_length, num_workers, ram_gb=None, *, safety_fraction=None, per_pool_epoch_per_worker_bytes=None, hard_cap=None, base_worker_gb=0.0)[source]

Pick random_negative_pool_epochs from box capacity.

The pool sits in the heap of every fit-worker process. Per-pool-epoch memory cost is dominated by:

  • num_random_negatives × peptide_max_length int8 indices,

  • intermediate pandas.Series[str] allocations and encoder buffers observed in practice at ~10–40× the int8 size on 2026-04 measurements (the old pool_epochs=100 run OOM’d a 944 GB box at ~199 GB / worker worth of transient pool cost).

We budget for the pessimistic per_pool_epoch_per_worker_bytes figure so the auto value is safe under transient peaks; tunable via env when a workload is known to behave better than the empirical pessimistic figure.

Returns an int >= 1. 1 means fresh random negatives every epoch. > 1 amortizes random-negative generation + encoding across N epochs.

Parameters:
num_random_negativesint

Number of random-negative peptides per epoch (the planner’s get_total_count()). The size of one pool-epoch in the cycle.

peptide_max_lengthint

Longest peptide the encoding allocates space for. With the fixed-vector encoding lookup, the per-peptide footprint is peptide_max_length int8 bytes.

num_workersint

Total fit() worker processes that will share the box. Each holds its own RN pool, so total RAM cost = num_workers × pool_epochs × per_pool_epoch_per_worker_bytes.

ram_gbfloat, optional

Total system RAM in GB. None returns 1 (safe default — we don’t know the budget so don’t bump pool_epochs above legacy).

safety_fractionfloat, optional

Expert fraction of RAM available to RN pools across all workers. By default the shared host reserve is used and base_worker_gb is subtracted first. Override via MHCFLURRY_AUTO_RN_POOL_SAFETY_FRACTION.

per_pool_epoch_per_worker_bytesfloat, optional

Empirical per-pool-epoch per-worker RAM cost in bytes. Default 1 GB (conservative — captures the int8 indices + transient pandas + encoder buffers seen on the 2026-04 run). Override via MHCFLURRY_AUTO_RN_POOL_PER_EPOCH_PER_WORKER_GB.

hard_capint, optional

Maximum pool epochs. Default 10 (expert-overridable); larger pools add startup/memory cost after generation overhead is already amortized.

base_worker_gbfloat, optional

Baseline host memory per fit worker, subtracted before sizing pools.

Notes

Total available bytes for RN pools across the box is usable host memory after the shared reserve and fit-worker base RSS. An explicit safety_fraction replaces that calculation. Per-worker budget is available / max(num_workers, 1). The number of pool epochs that fit is per_worker_budget / per_pool_epoch_per_worker_bytes, clamped to [1, hard_cap].

Cross-checks

  • 8×A100-80GB Verda shares the same reserve and worker-RSS estimate as the outer process planner, then stops at the throughput cap.

  • Single A100 80G Lambda uses the shared host reserve minus the two workers’ base RSS, then fills the remaining per-worker pool budget.

  • Tight and RAM-starved nodes naturally fall back toward one pooled epoch as the fit-worker base RSS consumes the shared usable budget.

mhcflurry.parallelism.planning.resolved_int(value, name)[source]
mhcflurry.parallelism.planning.auto_num_jobs(num_gpus, max_workers_per_gpu)[source]

Compute total fit-worker count from GPU plan.

Returns num_gpus * max_workers_per_gpu for GPU runs, 0 when no GPUs are visible (caller decides serial vs explicit CPU pool). Treats "auto" max_workers_per_gpu as not-yet-resolved and raises; callers must resolve it first via auto_max_workers_per_gpu.

mhcflurry.parallelism.planning.resolve_max_workers_per_gpu(args, per_worker_gb=None, num_gpus=None, backend=None)[source]

Resolve args.max_workers_per_gpu to an int, mutating args.

Accepts the literal string "auto" (the default) or an int. When "auto", calls auto_max_workers_per_gpu with the rest of the args’ parallelism config to pick a value. Idempotent — calling twice on the same args is a no-op the second time.

per_worker_gb lets workload-specific commands (e.g. calibrate, where cached_stages dominates per-worker VRAM at ~15 GB) override the train-default. Falls back to env var / the module default when not given.

Returns the resolved int (also stored on args.max_workers_per_gpu so subsequent consumers see the int).

mhcflurry.parallelism.planning.num_workers_per_gpu_from_args(args)[source]

Return resolved max_workers_per_gpu for model auto-sizing.

Callers must run resolve_local_parallelism_args first so workload- and hardware-aware defaults have already converted the CLI sentinel "auto" into an integer.

mhcflurry.parallelism.planning.refresh_device_memory_budget(args)[source]

Snapshot and propagate one fixed launch-time entitlement per worker.

mhcflurry.parallelism.planning.resolve_local_parallelism_args(args, cap_auto_num_jobs=True, per_worker_gb=None, workload_name='generic', workload_hints=None)[source]

Resolve and normalize local parallelism arguments through the planner.

mhcflurry.parallelism.planning.resolve_cpu_threads_per_worker(plan, cpu_count=None)[source]

Set unset BLAS/OpenMP thread env vars from the final worker plan.

mhcflurry.parallelism.planning.resolve_cpu_thread_budget(plan, cpu_count=None)[source]

Resolve native threads per fit worker and environment ownership.

A uniform runtime resize is safe only when all supported environment variables are unset or were written by an earlier mhcflurry auto pass. Any caller-owned value makes the environment authoritative; the numeric auto estimate is still recorded for diagnostics but must not be applied to loaded native or PyTorch pools.

mhcflurry.parallelism.planning.refine_local_parallelism_from_warmup(args, reports)[source]

Tighten an automatic plan using representative warmup measurements.

The analytic estimate remains the floor because a one-batch warmup does not contain a full resident dataset or validation. Measured CUDA reserved memory and process peak RSS can only increase the estimated working set and reduce automatic concurrency. Explicit --num-jobs and --max-workers-per-gpu values are never changed.

mhcflurry.parallelism.planning.refine_local_parallelism_from_spawn_context(args, context_bytes, memory=None)[source]

Cap automatic spawn workers using the loaded per-process context.

Spawn serializes the full worker initializer context into every process, unlike fork’s copy-on-write inheritance. Input files may be compressed, so their on-disk size is not a reliable host-RAM estimate. This refinement is deliberately performed after the dataset/model context has been loaded but before any workers start.

mhcflurry.parallelism.planning.apply_random_negative_pool_epochs_to_work_items(work_items, pool_epochs, *, log=None)[source]

Inject random_negative_pool_epochs into every work item’s hyperparameters.

Parallel to apply_dataloader_num_workers_to_work_items: the orchestrator chooses an auto value once at startup (or honors the CLI int) and writes it into every per-work-item hyperparameter dict. fit() reads it from self.hyperparameters['random_negative_pool_epochs'] when constructing its RandomNegativesPool.

Parameters:
work_itemslist of dict

Each dict has a hyperparameters sub-dict.

pool_epochsint

The resolved value from resolve_local_parallelism_args (>= 1).

logcallable, optional

Logging hook. Defaults to print.

mhcflurry.parallelism.planning.apply_dataloader_num_workers_to_work_items(work_items, num_workers, *, log=None)[source]

Inject dataloader_num_workers into every work item’s hyperparameters.

Writes the resolved value into affinity work-item hyperparameters. It is consumed by streaming pretraining; in-memory affinity fitting ignores it. Processing models do not use this affinity DataLoader hyperparameter.

Parameters:
work_itemslist of dict

Each dict has a hyperparameters sub-dict (the canonical shape produced by train_pan_allele_models_command / train_allele_specific_models_command).

num_workersint

The resolved value from resolve_local_parallelism_args. Use 0 to force in-process batching (no prefetch children).

logcallable, optional

Logging hook for the human-readable summary. Defaults to print.

mhcflurry.parallelism.planning.apply_resolved_training_hyperparameters_to_work_items(work_items, args, *, log=None)[source]

Inject resolved per-model training knobs into work item hyperparameters.

resolve_local_parallelism_args owns hardware-dependent CLI resolution. Trainers should call this once after constructing work items so pan-allele, allele-specific, and future affinity trainers all persist the same resolved settings in component-model hyperparameters.

mhcflurry.parallelism.torch_compile module

Torch compile environment and cache-warmup helpers.

mhcflurry.parallelism.torch_compile.resolve_torchinductor_compile_threads_env(num_jobs=1, phase='production')[source]

Ensure TORCHINDUCTOR_COMPILE_THREADS is parseable by PyTorch.

MHCflurry launchers accept TORCHINDUCTOR_COMPILE_THREADS=auto as an orchestrator-owned sentinel so each command can size the Inductor compiler helper pool after it knows its local worker count. PyTorch itself does not accept that sentinel: Inductor parses the env var with int(...).

This helper is for entry points that do not run the full local-parallelism resolver before importing or spawning PyTorch work. It leaves unset and user-pinned integer values alone, resolves orchestrator-owned auto to a numeric value, and fails early on other invalid values with a clearer message than the downstream Inductor traceback.

mhcflurry.parallelism.torch_compile.configure_cluster_worker_torch_compile_threads()[source]

Auto-size Inductor helper threads inside one cluster worker process.

Cluster parallelism submits each work item as its own process, often on different nodes. We therefore do not try to share a compile cache across the cluster. Each worker process still needs the same local policy: if compile is enabled and TORCHINDUCTOR_COMPILE_THREADS is unset or auto, pick a numeric value on that machine before the first torch.compile call.

If a scheduler packs several mhcflurry work items onto one node, set MHCFLURRY_CLUSTER_WORKERS_PER_NODE so the auto value is divided across those co-resident compiler processes. Otherwise the default assumes one work process owns its scheduler CPU allocation.

mhcflurry.parallelism.torch_compile.hoist_torchinductor_compile_threads(args, phase='production')[source]

Auto-size TORCHINDUCTOR_COMPILE_THREADS for local training.

torch.compile (when enabled via MHCFLURRY_TORCH_COMPILE=1) spins up an inductor compile worker pool that defaults to os.cpu_count() threads. With N fit() workers each running their own compile pool, that multiplies into an oversubscribed compile storm. The production phase uses an auto value derived from available cores and the worker count; the warmup phase uses a larger value because only one worker is compiling.

The orchestrator owns “how many workers will exist”, so it owns the env knob too: set once before forking, every worker inherits. Skips the hoist when the user has already pinned the value or when MHCFLURRY_TORCH_COMPILE isn’t on. Cluster workers running on other hosts must size themselves locally; see configure_cluster_worker_torch_compile_threads.

Lives here (not in any one train_*_command module) so processing, allele-specific, and any future train command can call it the same way.

mhcflurry.parallelism.torch_compile.run_single_worker_resource_probe(args, work_items, work_function, constant_data=None)[source]

Measure one real peak phase for every resource-distinct architecture.

The command’s resource_probe_only path must exercise its true resident data, configured minibatch, and validation phase for one bounded epoch. The single spawned worker records process-level CUDA memory (including the context), allocator peaks, and host RSS. Those observations may only tighten automatic concurrency before the production pool starts.

When torch.compile is enabled the same pass also primes its on-disk cache; resource safety is intentionally independent of compile being enabled.

work_items is not mutated: every task still runs in production.

Returns None when skipped, otherwise the number of unique architectures warmed.

mhcflurry.parallelism.torch_compile.run_single_worker_torch_compile_warmup(args, work_items, work_function, constant_data=None)[source]

Compatibility alias for run_single_worker_resource_probe().

mhcflurry.parallelism.worker_pool module

Multiprocessing worker-pool scheduling helpers.

class mhcflurry.parallelism.worker_pool.NonDaemonProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, Process

A multiprocessing.Process whose daemon flag cannot be set.

Reading .daemon always returns False; writes are no-ops. This lets us instantiate multiprocessing.pool.Pool with a worker class that declines to be a daemon, so the DataLoader inside each worker can spawn its own prefetch children.

class mhcflurry.parallelism.worker_pool.NonDaemonContext[source]

Bases: ForkContext

A multiprocessing context that hands out NonDaemonProcess workers.

Subclasses the current default multiprocessing context so its start method is preserved — we only swap the Process class. The Pool uses self._ctx.Process(...) to create workers and will now get our non-daemonic variant.

Process

alias of NonDaemonProcess

class mhcflurry.parallelism.worker_pool.NonDaemonSpawnProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, SpawnProcess

class mhcflurry.parallelism.worker_pool.NonDaemonSpawnContext[source]

Bases: SpawnContext

Process

alias of NonDaemonSpawnProcess

class mhcflurry.parallelism.worker_pool.NonDaemonForkProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, ForkProcess

class mhcflurry.parallelism.worker_pool.NonDaemonForkContext[source]

Bases: ForkContext

Process

alias of NonDaemonForkProcess

class mhcflurry.parallelism.worker_pool.NonDaemonForkServerProcess(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)[source]

Bases: _NonDaemonProcessMixin, ForkServerProcess

class mhcflurry.parallelism.worker_pool.NonDaemonForkServerContext[source]

Bases: ForkServerContext

Process

alias of NonDaemonForkServerProcess

mhcflurry.parallelism.worker_pool.non_daemon_context(start_method=None)[source]

Return a multiprocessing context whose workers are non-daemonic.

class mhcflurry.parallelism.worker_pool.NonDaemonPool(*args, **kwargs)[source]

Bases: Pool

A multiprocessing.Pool that runs non-daemonic workers.

Pool’s constructor takes a context kwarg — we thread a NonDaemonContext through so each worker is a NonDaemonProcess. Everything else (apply_async, imap, etc.) inherits unchanged.

mhcflurry.parallelism.worker_pool.chunk_ranges_for_local_parallelism(num_items, num_jobs=0, chunks_per_worker=4)[source]

Split a row/sequence axis into stable contiguous chunks for local workers.

Parameters:
num_itemsint

Number of input items.

num_jobsint

Number of worker processes. 0 yields one serial chunk.

chunks_per_workerint

Target number of work chunks per worker for load balancing.

Returns:
list of tuple

(chunk_index, start, end) ranges.

mhcflurry.parallelism.worker_pool.worker_pool_with_gpu_assignments_from_args(args, workload_name='generic', workload_hints=None, start_method=None, worker_context_module=None, worker_context_data=None)[source]

Create a multiprocessing.Pool where each worker uses its own GPU.

Uses commandline arguments. See worker_pool_with_gpu_assignments. Resolves args.max_workers_per_gpu="auto" to an int (mutating args so downstream consumers — e.g. inference batch sizing in calibrate — observe the same value).

Parameters:
argsargparse.Namespace

Parsed local-parallelism options.

workload_namestr

Workload profile used for automatic sizing.

workload_hintsdict, optional

Model/data size hints for that profile.

start_methodstr, optional

Multiprocessing start method.

worker_context_modulestr, optional

Module whose WORKER_CONTEXT receives constant data.

worker_context_datadict, optional

Constant data installed once per worker.

Returns:
multiprocessing.Pool
mhcflurry.parallelism.worker_pool.refine_local_parallelism_from_worker_context(args, worker_context_data, start_method=None)[source]

Apply spawn-context host sizing before code chooses serial/parallel.

mhcflurry.parallelism.worker_pool.estimate_worker_context_bytes(value, _seen=None)[source]

Best-effort deep resident size for data copied into spawn workers.

mhcflurry.parallelism.worker_pool.worker_pool_uses_fork(worker_pool=None)[source]

Return True when local Pool workers inherit parent globals by fork.

mhcflurry.parallelism.worker_pool.attach_constant_data_to_work_items_if_needed(work_items, constant_data, worker_pool, *, log=None)[source]

Attach constant data only when the Pool cannot inherit it by fork.

mhcflurry.parallelism.worker_pool.worker_pool_with_gpu_assignments(num_jobs, num_gpus=0, backend=None, max_workers_per_gpu=1, max_tasks_per_worker=None, worker_log_dir=None, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, start_method=None, worker_context_module=None, worker_context_data=None, device_memory_budget_bytes=None)[source]

Create a multiprocessing.Pool where each worker uses its own GPU.

Parameters:
num_jobsint

Number of worker processes.

num_gpusint
backendstring
max_workers_per_gpuint
max_tasks_per_workerint
worker_log_dirstring
cpu_threads_per_workerint

Runtime BLAS/OpenMP/PyTorch thread limit applied in each worker.

cpu_threads_per_worker_was_autobool

Whether mhcflurry owns the uniform runtime limit. False preserves caller-provided OMP/MKL/OpenBLAS settings.

start_methodstring

Optional multiprocessing start method, e.g. "spawn" when workers must not inherit PyTorch state from the parent process.

worker_context_modulestring, optional

Module containing a WORKER_CONTEXT dictionary used by worker functions.

worker_context_datadict, optional

Constant data to install once per process during initialization. This avoids serializing the same large payload with every queued task on spawn-based platforms.

device_memory_budget_bytesint, optional

Fixed launch-time device-memory entitlement propagated to each GPU worker for elastic batch sizing. Kept last in the signature so adding the planner-owned value does not change existing positional calls.

Returns:
multiprocessing.Pool
mhcflurry.parallelism.worker_pool.validate_worker_pool_args(num_jobs, num_gpus=0, backend='auto', max_workers_per_gpu=1)[source]

Validate local worker scheduling arguments.

--gpus controls CUDA worker assignment only. It does not select MPS devices and it does not distribute a single model across multiple GPUs.

mhcflurry.parallelism.worker_pool.worker_init_kwargs_for_scheduler(num_jobs, num_gpus=0, backend='auto', max_workers_per_gpu=1, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, device_memory_budget_bytes=None)[source]

Build per-worker init kwargs from the local scheduling configuration.

When num_gpus is set, workers are assigned one CUDA GPU each in round robin order. Any additional workers are forced onto CPU by hiding CUDA and setting their backend to cpu.

mhcflurry.parallelism.worker_pool.make_worker_pool(processes=None, initializer=None, initializer_kwargs_per_process=None, initializer_shared_kwargs=None, max_tasks_per_worker=None, start_method=None)[source]

Convenience wrapper to create a multiprocessing.Pool.

This function adds support for per-worker initializer arguments, which are not natively supported by the multiprocessing module. The motivation for this feature is to support allocating each worker to a (different) GPU.

IMPLEMENTATION NOTE:

The per-worker initializer arguments are implemented using a SimpleQueue. Each worker reads its arguments from this queue when it starts. When it terminates, it adds its initializer arguments back to the queue, so a future process can initialize itself using these arguments. SimpleQueue is important here: Queue.put uses a feeder thread, so workers can observe a transiently empty queue during startup and duplicate GPU assignments. A worker that puts its arguments back can also hang forever joining that feeder thread during process finalization.

There is one issue with this approach, however. If a worker crashes, it never repopulates the queue of initializer arguments. This will prevent any future worker from re-using those arguments. To deal with this issue we add a second ‘backup queue’. This queue always contains the full set of initializer arguments: whenever a worker reads from it, it always pushes the pop’d args back to the end of the queue immediately. If the primary arg queue is ever empty, then workers will read from this backup queue.

Parameters:
processesint

Number of workers. Default: num CPUs.

initializerfunction, optional

Init function to call in each worker

initializer_kwargs_per_processlist of dict, optional

Arguments to pass to initializer function for each worker. Length of list must equal the number of workers.

initializer_shared_kwargsdict, optional

Arguments passed once to every worker initializer. Unlike work-item arguments, large values here are serialized only once per process.

max_tasks_per_workerint, optional

Restart workers after this many tasks.

start_methodstring, optional

Multiprocessing start method to use for the worker pool.

Returns:
multiprocessing.Pool

mhcflurry.parallelism.worker_runtime module

Worker initialization and error-wrapping helpers.

mhcflurry.parallelism.worker_runtime.install_worker_context(module_name, data)[source]

Install one process-wide constant-data mapping for worker functions.

mhcflurry.parallelism.worker_runtime.configure_worker_cpu_threads(num_threads, auto_owned=True)[source]

Apply a fit worker’s CPU-thread budget to loaded native runtimes.

Forked workers inherit already-initialized BLAS/OpenMP pools, so setting only environment variables in the parent is insufficient. Keep the environment synchronized for libraries loaded later and use threadpoolctl to resize the native pools already present in this process. Caller-owned OMP/MKL/OpenBLAS settings are authoritative and must not be replaced by a uniform auto-derived limit.

mhcflurry.parallelism.worker_runtime.release_initializer_slot(slots, slot, generation)[source]

Return a worker’s assignment without writing a shutdown pipe.

mhcflurry.parallelism.worker_runtime.worker_init_entry_point(init_function, kwargs_per_process=None, slots=None, sequence=None, shared_kwargs=None)[source]
mhcflurry.parallelism.worker_runtime.worker_init(keras_backend=None, backend=None, gpu_device_nums=None, worker_log_dir=None, max_workers_per_gpu=None, cpu_threads_per_worker=None, cpu_threads_per_worker_was_auto=True, worker_context_module=None, worker_context_data=None, device_memory_budget_bytes=None)[source]
exception mhcflurry.parallelism.worker_runtime.WrapException[source]

Bases: Exception

Add traceback info to exception so exceptions raised in worker processes can still show traceback info when re-raised in the parent.

mhcflurry.parallelism.worker_runtime.call_wrapped(function, *args, **kwargs)[source]

Run function on args and kwargs and return result, wrapping any exception raised in a WrapException.

Parameters:
functionarbitrary function
Any other arguments provided are passed to the function.
Returns:
object
mhcflurry.parallelism.worker_runtime.call_wrapped_kwargs(function, kwargs)[source]

Invoke function on given kwargs and return result, wrapping any exception raised in a WrapException.

Parameters:
functionarbitrary function
kwargsdict
Returns:
object
Result of calling function(**kwargs).