Skip to content

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

SharedData(sizes: Dict[str, Union[int, None]])

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
def __init__(self, sizes: Dict[str, Union[int, None]]):
    """Initialize shared memory objects.

    Args:
        sizes: Dictionary mapping keys to array sizes.
              If size is None, creates a shared float value.
    """
    self._data_dict = {}
    for key, size in sizes.items():
        if size is None:
            self._data_dict[key] = Value('d', 0.0)
        else:
            self._data_dict[key] = Array('d', size)

__getitem__

__getitem__(key: Union[str, Tuple[str, slice]]) -> Union[np.ndarray, float]

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
def __getitem__(
    self, key: Union[str, Tuple[str, slice]]
) -> Union[np.ndarray, float]:
    """Get numpy array view or float value.

    Args:
        key: Key in data dictionary or tuple of (key, slice)

    Returns:
        numpy array view for Array, float for Value
    """
    if isinstance(key, tuple):
        array_key, slice_obj = key
        return self[array_key][slice_obj]

    data = self._data_dict[key]
    if hasattr(data, 'get_obj'):  # Array has get_obj, Value doesn't
        return np.frombuffer(data.get_obj(), dtype=np.float64)
    return data.value

__setitem__

__setitem__(key: Union[str, Tuple[str, slice]], value: Union[ndarray, float])

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
def __setitem__(
    self,
    key: Union[str, Tuple[str, slice]],
    value: Union[np.ndarray, float]
):
    """Set array contents or float value.

    Args:
        key: Key in data dictionary or tuple of (key, slice)
        value: Array or float to set
    """
    if isinstance(key, tuple):
        array_key, slice_obj = key
        data = self._data_dict[array_key]
        if hasattr(data, 'get_obj'):
            data[slice_obj] = value
        else:
            raise ValueError("Slice assignment only supported for Array types")
        return

    data = self._data_dict[key]
    if hasattr(data, 'get_obj'):  # Array has get_obj, Value doesn't
        np.copyto(np.frombuffer(data.get_obj(), dtype=np.float64), value)
    else:
        data.value = float(value)

SerialManager

SerialManager(num_processes: int)

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
def __init__(self, num_processes: int):
    """Initialize worker manager.

    Args:
        num_processes: Number of worker processes to manage
    """
    self.flags = [Value('i', 0) for _ in range(num_processes)]
    self.functions: List[Callable] = []
    self.arguments: List[tuple] = []
    self.states: List[Any] = []

start_workers

start_workers(flag: Optional[int] = None) -> None

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
def start_workers(self, flag: Optional[int] = None) -> None:
    """Signal workers to start processing.

    Args:
        flag: Optional flag value to set (default: 1)
    """
    for f in self.flags:
        f.value = flag or 1

    offset = 0
    for i in range(len(self.flags)):
        func, args, state = self.functions[i], self.arguments[i], self.states[i]
        self.states[i] = func(state, offset, *args)
        offset += self.states[i].shape[0]

    assert all([f.value == 0 for f in self.flags])

add_process

add_process(target: Callable, args: Tuple) -> None

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
def add_process(self, target: Callable, args: Tuple) -> None:
    """Add a worker process.

    Args:
        target: Function to run in process
        args: Arguments to pass to function
    """
    self.functions.append(target)
    self.arguments.append(args)
    self.states.append(None)

shutdown

shutdown() -> None

Shutdown all worker processes.

Source code in src/graphld/multiprocessing_template.py
def shutdown(self) -> None:
    """Shutdown all worker processes."""
    pass

WorkerManager

WorkerManager(num_processes: int)

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
def __init__(self, num_processes: int):
    """Initialize worker manager.

    Args:
        num_processes: Number of worker processes to manage
    """
    self.flags = [Value('i', 0) for _ in range(num_processes)]
    self.processes: List[Process] = []

start_workers

start_workers(flag: Optional[int] = None) -> None

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
def start_workers(self, flag: Optional[int] = None) -> None:
    """Signal workers to start processing.

    Args:
        flag: Optional flag value to set (default: 1)
    """
    if flag is None:
        flag = 1
    for f in self.flags:
        f.value = flag

await_workers

await_workers() -> None

Wait for all workers to finish current task; abort if a worker crashes.

Source code in src/graphld/multiprocessing_template.py
def await_workers(self) -> None:
    """Wait for all workers to finish current task; abort if a worker crashes."""
    while True:
        # Kill all workers if any throws an error
        if any(flag.value == -1 for flag in self.flags):
            for f in self.flags:
                f.value = -1
            for p in self.processes:
                p.join(timeout=0.5)
            raise RuntimeError("Worker process crashed. Aborting.")

        # Normal completion
        if all(flag.value < 1 for flag in self.flags):
            break

        time.sleep(0.01)

add_process

add_process(target: Callable, args: Tuple) -> None

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
def add_process(self, target: Callable, args: Tuple) -> None:
    """Add a worker process.

    Args:
        target: Function to run in process
        args: Arguments to pass to function
    """
    process = Process(target=target, args=args)
    process.start()
    self.processes.append(process)

shutdown

shutdown() -> None

Shutdown all worker processes.

Source code in src/graphld/multiprocessing_template.py
def shutdown(self) -> None:
    """Shutdown all worker processes."""
    # Signal shutdown
    for flag in self.flags:
        flag.value = -1

    # Wait for processes to finish
    for process in self.processes:
        process.join()

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

create_shared_memory(metadata: DataFrame, block_data: list, **kwargs: Any) -> SharedData

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
@classmethod
@abstractmethod
def create_shared_memory(
    cls, metadata: pl.DataFrame, block_data: list, **kwargs: Any
) -> 'SharedData':
    """Initialize shared memory and data structures.

    Args:
        metadata: polars dataframe continaing LDGM metadata for each LD block
        block_data: List of block-specific data
        **kwargs: Additional arguments passed from run()

    Returns:
        SharedData object containing shared memory arrays
    """
    pass

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
@classmethod
@abstractmethod
def supervise(cls, manager: Union[WorkerManager, SerialManager],
            shared_data: SharedData,
            block_data: list, **kwargs: Any) -> Any:
    """Monitor workers and process results.

    Args:
        manager: Worker manager for controlling processes
        shared_data: Shared memory data
        **kwargs: Additional arguments passed from run()

    Returns:
        Results of the parallel computation
    """
    pass

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
@classmethod
@abstractmethod
def process_block(cls, ldgm: PrecisionOperator,
                flag: Value,
                shared_data: SharedData,
                block_offset: int,
                block_data: Any = None,
                worker_params: Any = None) -> None:
    """Process single block.

    Args:
        ldgm: LDGM object
        flag: Worker flag
        shared_data: Dictionary-like shared data object
        block_offset: Offset for this block
        block_data: Optional block-specific data from prepare_block_data
        worker_params: Optional parameters passed to each worker process

    Returns:
        None
    """
    pass

prepare_block_data classmethod

prepare_block_data(metadata: DataFrame, **kwargs: Any) -> list

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
@classmethod
def prepare_block_data(cls, metadata: pl.DataFrame, **kwargs: Any) -> list:
    """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.

    Args:
        metadata: Metadata DataFrame containing block information
        **kwargs: Additional arguments passed from run()

    Returns:
        List of block-specific data, length equal to number of blocks
    """
    return [None] * len(metadata)

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
@classmethod
def worker(cls,
           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.

    Args:
        files: List of LDGM files to process
        block_data: List of block-specific data
        flag: Shared flag for worker control
        shared_data: Shared memory data
        offset: In shared data, where to start processing
        worker_params: Optional parameters passed to process_block
    """
    try:
        # Load LDGMs once
        ldgms = []
        for file, population in files:
            ldgm = load_ldgm(str(file), population=population)
            ldgm.factor()
            ldgms.append(ldgm)

        while True:
            # Wait for signal to start new iteration
            while flag.value == 0:
                time.sleep(0.01)

            if flag.value == -1:  # shutdown signal
                break

            # Process all blocks and collect solutions
            block_offset = offset
            starting_flag = flag.value
            for ldgm, data in zip(ldgms, block_data, strict=False):
                try:
                    cls.process_block(ldgm, flag, shared_data, block_offset, data, worker_params)
                    # Check that process_block didn't modify the flag during normal execution
                    assert flag.value == starting_flag, "process_block should not change flag"
                except Exception:
                    # Ensure the flag is flipped so the supervisor doesn't spin forever
                    flag.value = -1
                    raise
                block_offset += ldgm.shape[0]
            # Signal completion
            flag.value = 0

    except Exception as e:
        print(f"Error in worker: {e}")
        flag.value = -1

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
@classmethod
def serial_worker(cls,
           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.

    Args:
        ldgm: the LDGM that was previously loaded, or None
        offset: In shared data, where to start processing
        file: LDGM file and population context to process
        data: List of block-specific data
        flag: Shared flag for worker control
        shared_data: Shared memory data
        worker_params: Optional parameters passed to process_block
    """
    if flag.value <= 0:
        raise ValueError("Serial worker should never be started with flag <= 0")

    # Load LDGMs once
    if ldgm is None:
        edgelist_file, population = file
        ldgm = load_ldgm(str(edgelist_file), population=population)
        ldgm.factor()

    # Process all blocks and collect solutions
    cls.process_block(ldgm, flag, shared_data, offset, data, worker_params)

    # Signal completion
    flag.value = 0

    return ldgm

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
@classmethod
def run(cls,
        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.

    Args:
        ldgm_metadata_path: Path to metadata file
        populations: Populations to process; None -> all
        chromosomes: Chromosomes to process; None -> all
        num_processes: Number of processes to use
        worker_params: Optional parameters passed to each worker process
        **kwargs: Additional arguments

    Returns:
        Results of the parallel computation
    """


    # Read metadata first
    metadata = read_ldgm_metadata(
        ldgm_metadata_path,
        populations=populations,
        chromosomes=chromosomes
    )

    # Get list of files from metadata
    ldgm_directory = Path(ldgm_metadata_path).parent
    edgelist_files = [
        _ldgm_block_context(ldgm_directory, block)
        for block in metadata.iter_rows(named=True)
    ]
    if not edgelist_files:
        raise FileNotFoundError("No edgelist files found in metadata")

    if num_processes is None:
        num_processes = min(len(edgelist_files), cpu_count())

    # Split files among processes
    process_block_ranges, process_offsets = cls._split_blocks(metadata, num_processes)
    process_files = [edgelist_files[start:end] for start, end in process_block_ranges]

    # Data to be sent to each block individually
    block_data = cls.prepare_block_data(metadata, **kwargs)
    process_block_data = [block_data[start:end] for start, end in process_block_ranges]

    # Data shared among all blocks
    shared_data = cls.create_shared_memory(metadata, block_data, **kwargs)

    # Create worker manager
    manager = WorkerManager(num_processes)

    # Start workers
    for i in range(num_processes):
        manager.add_process(
            target=cls.worker,
            args=(
                process_files[i],
                process_block_data[i],
                manager.flags[i],
                shared_data,
                process_offsets[i],
                worker_params,
            )
        )

    # Run supervisor process
    results = cls.supervise(manager, shared_data, block_data, **kwargs)

    # Cleanup
    manager.shutdown()

    return results

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

Source code in src/graphld/multiprocessing_template.py
@classmethod
def run_serial(cls,
        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.

    Args:
        ldgm_metadata_path: Path to metadata file
        populations: Populations to process; None -> all
        chromosomes: Chromosomes to process; None -> all
        num_processes: Ignored in serial mode; accepted for API compatibility
        worker_params: Optional parameters passed to each worker process
        **kwargs: Additional arguments

    Returns:
        Results of the computation
    """

    # Read metadata first
    metadata = read_ldgm_metadata(
        ldgm_metadata_path,
        populations=populations,
        chromosomes=chromosomes
    )

    # Get list of files from metadata
    ldgm_directory = Path(ldgm_metadata_path).parent
    edgelist_files = [
        _ldgm_block_context(ldgm_directory, block)
        for block in metadata.iter_rows(named=True)
    ]
    if not edgelist_files:
        raise FileNotFoundError("No edgelist files found in metadata")

    # Data to be sent to each block individually
    block_data = cls.prepare_block_data(metadata, **kwargs)
    num_blocks = len(block_data)

    # Data shared among all blocks
    shared_data = cls.create_shared_memory(metadata, block_data, **kwargs)

    manager = SerialManager(num_blocks)

    # Start workers
    for i in range(num_blocks):
        manager.add_process(
            target=cls.serial_worker,
            args=(
                edgelist_files[i],
                block_data[i],
                manager.flags[i],
                shared_data,
                worker_params,
            )
        )

    # Run supervisor process
    results = cls.supervise(manager, shared_data, block_data, **kwargs)

    # Cleanup
    manager.shutdown()

    return results