graphld.multiprocessing_template¶
Parallel processing utilities for LDGM-based algorithms.
ParallelProcessor is a base class for implementing parallel algorithms with LDGMs, wrapping Python's multiprocessing module. It splits work among processes, each of which loads a subset of LD blocks.
For higher-level context, see the Parallel Processing guide.
multiprocessing_template
¶
Base class for parallel processing applications.
SharedData
¶
Wrapper for shared memory data structures.
Attributes:
| Name | Type | Description |
|---|---|---|
_data_dict |
Dictionary mapping keys to shared memory objects |
Initialize shared memory objects.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
sizes
|
Dict[str, Union[int, None]]
|
Dictionary mapping keys to array sizes. If size is None, creates a shared float value. |
required |
Source code in src/graphld/multiprocessing_template.py
__getitem__
¶
Get numpy array view or float value.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
key
|
Union[str, Tuple[str, slice]]
|
Key in data dictionary or tuple of (key, slice) |
required |
Returns:
| Type | Description |
|---|---|
Union[ndarray, float]
|
numpy array view for Array, float for Value |
Source code in src/graphld/multiprocessing_template.py
__setitem__
¶
Set array contents or float value.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
key
|
Union[str, Tuple[str, slice]]
|
Key in data dictionary or tuple of (key, slice) |
required |
value
|
Union[ndarray, float]
|
Array or float to set |
required |
Source code in src/graphld/multiprocessing_template.py
SerialManager
¶
Manager for debugging by running workers in serial.
Initialize worker manager.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
num_processes
|
int
|
Number of worker processes to manage |
required |
Source code in src/graphld/multiprocessing_template.py
start_workers
¶
Signal workers to start processing.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
flag
|
Optional[int]
|
Optional flag value to set (default: 1) |
None
|
Source code in src/graphld/multiprocessing_template.py
add_process
¶
Add a worker process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
target
|
Callable
|
Function to run in process |
required |
args
|
Tuple
|
Arguments to pass to function |
required |
Source code in src/graphld/multiprocessing_template.py
WorkerManager
¶
Manager for coordinating parallel worker processes.
Attributes:
| Name | Type | Description |
|---|---|---|
flags |
List of shared flags for worker control |
|
processes |
List[Process]
|
List of worker processes |
Initialize worker manager.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
num_processes
|
int
|
Number of worker processes to manage |
required |
Source code in src/graphld/multiprocessing_template.py
start_workers
¶
Signal workers to start processing.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
flag
|
Optional[int]
|
Optional flag value to set (default: 1) |
None
|
Source code in src/graphld/multiprocessing_template.py
await_workers
¶
Wait for all workers to finish current task; abort if a worker crashes.
Source code in src/graphld/multiprocessing_template.py
add_process
¶
Add a worker process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
target
|
Callable
|
Function to run in process |
required |
args
|
Tuple
|
Arguments to pass to function |
required |
Source code in src/graphld/multiprocessing_template.py
shutdown
¶
Shutdown all worker processes.
ParallelProcessor
¶
Bases: ABC
Abstract base class for parallel processing applications.
This class provides a framework for parallel processing of LDGM files. Subclasses must implement the following methods: - initialize: Set up shared memory arrays and data structures - supervise: Monitor and control worker processes - process_block: Process a single LDGM block
create_shared_memory
abstractmethod
classmethod
¶
Initialize shared memory and data structures.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
metadata
|
DataFrame
|
polars dataframe continaing LDGM metadata for each LD block |
required |
block_data
|
list
|
List of block-specific data |
required |
**kwargs
|
Any
|
Additional arguments passed from run() |
{}
|
Returns:
| Type | Description |
|---|---|
SharedData
|
SharedData object containing shared memory arrays |
Source code in src/graphld/multiprocessing_template.py
supervise
abstractmethod
classmethod
¶
supervise(manager: Union[WorkerManager, SerialManager], shared_data: SharedData, block_data: list, **kwargs: Any) -> Any
Monitor workers and process results.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
manager
|
Union[WorkerManager, SerialManager]
|
Worker manager for controlling processes |
required |
shared_data
|
SharedData
|
Shared memory data |
required |
**kwargs
|
Any
|
Additional arguments passed from run() |
{}
|
Returns:
| Type | Description |
|---|---|
Any
|
Results of the parallel computation |
Source code in src/graphld/multiprocessing_template.py
process_block
abstractmethod
classmethod
¶
process_block(ldgm: PrecisionOperator, flag: Value, shared_data: SharedData, block_offset: int, block_data: Any = None, worker_params: Any = None) -> None
Process single block.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ldgm
|
PrecisionOperator
|
LDGM object |
required |
flag
|
Value
|
Worker flag |
required |
shared_data
|
SharedData
|
Dictionary-like shared data object |
required |
block_offset
|
int
|
Offset for this block |
required |
block_data
|
Any
|
Optional block-specific data from prepare_block_data |
None
|
worker_params
|
Any
|
Optional parameters passed to each worker process |
None
|
Returns:
| Type | Description |
|---|---|
None
|
None |
Source code in src/graphld/multiprocessing_template.py
prepare_block_data
classmethod
¶
Prepare data specific to each block for processing.
This method should return a list of length equal to the number of blocks, where each element contains any block-specific data needed by process_block. The base implementation returns None for each block.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
metadata
|
DataFrame
|
Metadata DataFrame containing block information |
required |
**kwargs
|
Any
|
Additional arguments passed from run() |
{}
|
Returns:
| Type | Description |
|---|---|
list
|
List of block-specific data, length equal to number of blocks |
Source code in src/graphld/multiprocessing_template.py
worker
classmethod
¶
worker(files: list, block_data: list, flag: Value, shared_data: SharedData, offset: int, worker_params: Any = None) -> None
Worker process that loads LDGMs and processes blocks.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
files
|
list
|
List of LDGM files to process |
required |
block_data
|
list
|
List of block-specific data |
required |
flag
|
Value
|
Shared flag for worker control |
required |
shared_data
|
SharedData
|
Shared memory data |
required |
offset
|
int
|
In shared data, where to start processing |
required |
worker_params
|
Any
|
Optional parameters passed to process_block |
None
|
Source code in src/graphld/multiprocessing_template.py
serial_worker
classmethod
¶
serial_worker(ldgm: Optional[PrecisionOperator], offset: int, file: tuple[Path, Optional[str]], data: list, flag: Value, shared_data: SharedData, worker_params: Any) -> PrecisionOperator
Worker process that loads LDGMs and processes blocks in serial.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ldgm
|
Optional[PrecisionOperator]
|
the LDGM that was previously loaded, or None |
required |
offset
|
int
|
In shared data, where to start processing |
required |
file
|
tuple[Path, Optional[str]]
|
LDGM file and population context to process |
required |
data
|
list
|
List of block-specific data |
required |
flag
|
Value
|
Shared flag for worker control |
required |
shared_data
|
SharedData
|
Shared memory data |
required |
worker_params
|
Any
|
Optional parameters passed to process_block |
required |
Source code in src/graphld/multiprocessing_template.py
run
classmethod
¶
run(ldgm_metadata_path: str, populations: Optional[Union[str, List[str]]] = None, chromosomes: Optional[Union[int, List[int]]] = None, num_processes: Optional[int] = None, worker_params: Any = None, **kwargs: Any) -> Any
Run parallel computation.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ldgm_metadata_path
|
str
|
Path to metadata file |
required |
populations
|
Optional[Union[str, List[str]]]
|
Populations to process; None -> all |
None
|
chromosomes
|
Optional[Union[int, List[int]]]
|
Chromosomes to process; None -> all |
None
|
num_processes
|
Optional[int]
|
Number of processes to use |
None
|
worker_params
|
Any
|
Optional parameters passed to each worker process |
None
|
**kwargs
|
Any
|
Additional arguments |
{}
|
Returns:
| Type | Description |
|---|---|
Any
|
Results of the parallel computation |
Source code in src/graphld/multiprocessing_template.py
432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 | |
run_serial
classmethod
¶
run_serial(ldgm_metadata_path: str, populations: Optional[Union[str, List[str]]] = None, chromosomes: Optional[Union[int, List[int]]] = None, num_processes: Optional[int] = None, worker_params: Any = None, **kwargs: Any) -> Any
Run computation block-by-block in the current process for debugging.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ldgm_metadata_path
|
str
|
Path to metadata file |
required |
populations
|
Optional[Union[str, List[str]]]
|
Populations to process; None -> all |
None
|
chromosomes
|
Optional[Union[int, List[int]]]
|
Chromosomes to process; None -> all |
None
|
num_processes
|
Optional[int]
|
Ignored in serial mode; accepted for API compatibility |
None
|
worker_params
|
Any
|
Optional parameters passed to each worker process |
None
|
**kwargs
|
Any
|
Additional arguments |
{}
|
Returns:
| Type | Description |
|---|---|
Any
|
Results of the computation |