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:
ForkContextA multiprocessing context that hands out
NonDaemonProcessworkers.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:
PoolA
multiprocessing.Poolthat runs non-daemonic workers.Pool’s constructor takes a
contextkwarg — we thread aNonDaemonContextthrough so each worker is aNonDaemonProcess. 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,ProcessA
multiprocessing.Processwhosedaemonflag cannot be set.Reading
.daemonalways returns False; writes are no-ops. This lets us instantiatemultiprocessing.pool.Poolwith 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:
ExceptionAdd 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_workersinto 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
hyperparameterssub-dict (the canonical shape produced bytrain_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_epochsinto 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 fromself.hyperparameters['random_negative_pool_epochs']when constructing itsRandomNegativesPool.- Parameters:
- work_itemslist of dict
Each dict has a
hyperparameterssub-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_argsowns 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_workershyperparameter. A return of0means 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_gpufor the canonical GPU run, ornum_jobsfor CPU-only runs.- vcpusint, optional
Total vCPU count. Default
os.cpu_count().- ram_gbfloat, optional
Total system RAM in GB.
Noneskips 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
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.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.CPU cap:
cpu_cap = cpu_per_fit // 2.RAM cap (when
ram_gbis 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).Throughput cap:
min(cpu_cap, ram_cap, hard_cap).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_gbvery small (e.g. < 2 GB / fit) →ram_cap = 0, falls back to in-process batching to preserve correctness over throughput.hard_capenv override of0→ 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_gpubased 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 fromMHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB). It deliberately avoidstorch.cudaso 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 viaMHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_PER_WORKER_GBandMHCFLURRY_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_gbexplicitly — the affinity calibration workload profile’sdevice_worker_gbof 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_gpufor GPU runs,0when no GPUs are visible (caller decides serial vs explicit CPU pool). Treats"auto"max_workers_per_gpuas not-yet-resolved and raises; callers must resolve it first viaauto_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_epochsfrom 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_lengthint8 indices,intermediate
pandas.Series[str]allocations and encoder buffers observed in practice at ~10–40× the int8 size on 2026-04 measurements (the oldpool_epochs=100run 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_bytesfigure 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.
1means fresh random negatives every epoch.> 1amortizes 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_lengthint8 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.
Nonereturns1(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_gbis subtracted first. Override viaMHCFLURRY_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 viaMHCFLURRY_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_fractionreplaces that calculation. Per-worker budget isavailable / max(num_workers, 1). The number of pool epochs that fit isper_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.
0yields 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_THREADSis unset orauto, pick a numeric value on that machine before the firsttorch.compilecall.If a scheduler packs several mhcflurry work items onto one node, set
MHCFLURRY_CLUSTER_WORKERS_PER_NODEso 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
Noneif undetectable.Env override (
MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB) first, thennvidia-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
Noneifnvidia-smiis 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 theminviafree_vram_from_nvidia_smi_gb.This is intentionally a subprocess call instead of
torch.cudaso 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_THREADSfor local training.torch.compile(when enabled viaMHCFLURRY_TORCH_COMPILE=1) spins up an inductor compile worker pool that defaults toos.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_COMPILEisn’t on. Cluster workers running on other hosts must size themselves locally; seeconfigure_cluster_worker_torch_compile_threads.Lives here (not in any one
train_*_commandmodule) 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.SimpleQueueis important here:Queue.putuses 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_gpufor model auto-sizing.Callers must run
resolve_local_parallelism_argsfirst 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-jobsand--max-workers-per-gpuvalues 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_workersvalue to an int.Accepts
"auto"/ unset /None(delegates toauto_dataloader_num_workers), or any int-coercible value. Used by both the orchestrator-side resolver and the shell helper that injectsdataloader_num_workersinto 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_gputo an int, mutatingargs.Accepts the literal string
"auto"(the default) or an int. When"auto", callsauto_max_workers_per_gpuwith 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_gblets 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_gpuso subsequent consumers see the int).
- mhcflurry.parallelism.resolve_torchinductor_compile_threads_env(num_jobs=1, phase='production')[source]
Ensure
TORCHINDUCTOR_COMPILE_THREADSis parseable by PyTorch.MHCflurry launchers accept
TORCHINDUCTOR_COMPILE_THREADS=autoas 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 withint(...).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
autoto 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_onlypath 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_itemsis not mutated: every task still runs in production.Returns
Nonewhen 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.
--gpuscontrols 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_gpusis 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 tocpu.
- 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_CONTEXTdictionary 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. Resolvesargs.max_workers_per_gpu="auto"to an int (mutatingargsso 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
Nonewhen 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_GBaccepts either one value applied to every GPU or a comma/space-separated list. It exists for tests and unusual launchers wherenvidia-smiis 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_DEVICESis 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
Noneifnvidia-smiis 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 theminviafree_vram_from_nvidia_smi_gb.This is intentionally a subprocess call instead of
torch.cudaso 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
Noneif undetectable.Env override (
MHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB) first, thennvidia-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_gpubased 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 fromMHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_FREE_VRAM_GB). It deliberately avoidstorch.cudaso 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 viaMHCFLURRY_AUTO_MAX_WORKERS_PER_GPU_PER_WORKER_GBandMHCFLURRY_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_gbexplicitly — the affinity calibration workload profile’sdevice_worker_gbof 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_workershyperparameter. A return of0means 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_gpufor the canonical GPU run, ornum_jobsfor CPU-only runs.- vcpusint, optional
Total vCPU count. Default
os.cpu_count().- ram_gbfloat, optional
Total system RAM in GB.
Noneskips 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
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.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.CPU cap:
cpu_cap = cpu_per_fit // 2.RAM cap (when
ram_gbis 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).Throughput cap:
min(cpu_cap, ram_cap, hard_cap).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_gbvery small (e.g. < 2 GB / fit) →ram_cap = 0, falls back to in-process batching to preserve correctness over throughput.hard_capenv override of0→ 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_workersvalue to an int.Accepts
"auto"/ unset /None(delegates toauto_dataloader_num_workers), or any int-coercible value. Used by both the orchestrator-side resolver and the shell helper that injectsdataloader_num_workersinto 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_epochsfrom 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_lengthint8 indices,intermediate
pandas.Series[str]allocations and encoder buffers observed in practice at ~10–40× the int8 size on 2026-04 measurements (the oldpool_epochs=100run 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_bytesfigure 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.
1means fresh random negatives every epoch.> 1amortizes 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_lengthint8 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.
Nonereturns1(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_gbis subtracted first. Override viaMHCFLURRY_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 viaMHCFLURRY_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_fractionreplaces that calculation. Per-worker budget isavailable / max(num_workers, 1). The number of pool epochs that fit isper_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.auto_num_jobs(num_gpus, max_workers_per_gpu)[source]
Compute total fit-worker count from GPU plan.
Returns
num_gpus * max_workers_per_gpufor GPU runs,0when no GPUs are visible (caller decides serial vs explicit CPU pool). Treats"auto"max_workers_per_gpuas not-yet-resolved and raises; callers must resolve it first viaauto_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_gputo an int, mutatingargs.Accepts the literal string
"auto"(the default) or an int. When"auto", callsauto_max_workers_per_gpuwith 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_gblets 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_gpuso subsequent consumers see the int).
- mhcflurry.parallelism.planning.num_workers_per_gpu_from_args(args)[source]
Return resolved
max_workers_per_gpufor model auto-sizing.Callers must run
resolve_local_parallelism_argsfirst 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-jobsand--max-workers-per-gpuvalues 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_epochsinto 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 fromself.hyperparameters['random_negative_pool_epochs']when constructing itsRandomNegativesPool.- Parameters:
- work_itemslist of dict
Each dict has a
hyperparameterssub-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_workersinto 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
hyperparameterssub-dict (the canonical shape produced bytrain_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_argsowns 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_THREADSis parseable by PyTorch.MHCflurry launchers accept
TORCHINDUCTOR_COMPILE_THREADS=autoas 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 withint(...).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
autoto 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_THREADSis unset orauto, pick a numeric value on that machine before the firsttorch.compilecall.If a scheduler packs several mhcflurry work items onto one node, set
MHCFLURRY_CLUSTER_WORKERS_PER_NODEso 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_THREADSfor local training.torch.compile(when enabled viaMHCFLURRY_TORCH_COMPILE=1) spins up an inductor compile worker pool that defaults toos.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_COMPILEisn’t on. Cluster workers running on other hosts must size themselves locally; seeconfigure_cluster_worker_torch_compile_threads.Lives here (not in any one
train_*_commandmodule) 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_onlypath 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_itemsis not mutated: every task still runs in production.Returns
Nonewhen 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,ProcessA
multiprocessing.Processwhosedaemonflag cannot be set.Reading
.daemonalways returns False; writes are no-ops. This lets us instantiatemultiprocessing.pool.Poolwith 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:
ForkContextA multiprocessing context that hands out
NonDaemonProcessworkers.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:
PoolA
multiprocessing.Poolthat runs non-daemonic workers.Pool’s constructor takes a
contextkwarg — we thread aNonDaemonContextthrough so each worker is aNonDaemonProcess. 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.
0yields 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. Resolvesargs.max_workers_per_gpu="auto"to an int (mutatingargsso 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_CONTEXTdictionary 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.
--gpuscontrols 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_gpusis 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 tocpu.
- 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.SimpleQueueis important here:Queue.putuses 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:
ExceptionAdd traceback info to exception so exceptions raised in worker processes can still show traceback info when re-raised in the parent.