Team Ai
Apppublic

pratik-250620/MultiModal-Coherence-AI

sourceHugging Facemitupdated 8mo agoView on Hugging Face
2likes
parallel_processing.py170 linesDownload Raw Back to utils
1"""2Parallel processing utilities for data pipeline optimization.3 4Supports:5- Parallel data ingestion6- Batch processing with multiprocessing7- Distributed processing support (foundation for Dask/Beam)8"""9 10from __future__ import annotations11 12import multiprocessing as mp13from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor, as_completed14from functools import partial15from typing import Any, Callable, Dict, Iterable, List, Optional, TypeVar16 17T = TypeVar("T")18R = TypeVar("R")19 20 21def parallel_map(22    func: Callable[[T], R],23    items: Iterable[T],24    max_workers: Optional[int] = None,25    use_threads: bool = False,26    chunk_size: int = 1,27    **func_kwargs,28) -> List[R]:29    """30    Parallel map function.31    32    Args:33        func: Function to apply to each item34        items: Iterable of items to process35        max_workers: Maximum number of workers (default: CPU count)36        use_threads: Use threads instead of processes (for I/O-bound tasks)37        chunk_size: Number of items per chunk (for ProcessPoolExecutor)38        **func_kwargs: Additional kwargs to pass to func39    40    Returns:41        List of results in the same order as items42    """43    if max_workers is None:44        max_workers = mp.cpu_count()45    46    items_list = list(items)47    if not items_list:48        return []49    50    # Prepare function with kwargs51    if func_kwargs:52        func = partial(func, **func_kwargs)53    54    executor_class = ThreadPoolExecutor if use_threads else ProcessPoolExecutor55    56    with executor_class(max_workers=max_workers) as executor:57        if use_threads:58            # ThreadPoolExecutor doesn't use chunk_size59            futures = [executor.submit(func, item) for item in items_list]60        else:61            # ProcessPoolExecutor supports chunking62            futures = executor.map(func, items_list, chunksize=chunk_size)63            return list(futures)64        65        # For ThreadPoolExecutor, collect results66        results = []67        for future in as_completed(futures):68            try:69                results.append(future.result())70            except Exception as e:71                # Log error but continue72                print(f"Error in parallel_map: {e}")73                results.append(None)74        75        # Reorder results to match input order (approximate for as_completed)76        # For exact ordering, use executor.map instead77        return results78 79 80def batch_process(81    func: Callable[[List[T]], List[R]],82    items: Iterable[T],83    batch_size: int = 32,84    max_workers: Optional[int] = None,85    **func_kwargs,86) -> List[R]:87    """88    Process items in batches in parallel.89    90    Args:91        func: Function that processes a batch and returns list of results92        items: Iterable of items to process93        batch_size: Number of items per batch94        max_workers: Maximum number of parallel batches95        **func_kwargs: Additional kwargs to pass to func96    97    Returns:98        Flattened list of results99    """100    items_list = list(items)101    batches = [102        items_list[i : i + batch_size] for i in range(0, len(items_list), batch_size)103    ]104    105    if not batches:106        return []107    108    # Prepare function with kwargs109    if func_kwargs:110        func = partial(func, **func_kwargs)111    112    if max_workers is None:113        max_workers = min(len(batches), mp.cpu_count())114    115    results = parallel_map(116        func,117        batches,118        max_workers=max_workers,119        use_threads=False,120    )121    122    # Flatten results123    flattened = []124    for batch_results in results:125        if batch_results:126            flattened.extend(batch_results)127    128    return flattened129 130 131class ParallelProcessor:132    """Wrapper for parallel processing with configuration."""133 134    def __init__(135        self,136        max_workers: Optional[int] = None,137        use_threads: bool = False,138        chunk_size: int = 1,139    ):140        self.max_workers = max_workers or mp.cpu_count()141        self.use_threads = use_threads142        self.chunk_size = chunk_size143 144    def map(self, func: Callable[[T], R], items: Iterable[T], **func_kwargs) -> List[R]:145        """Apply function to items in parallel."""146        return parallel_map(147            func,148            items,149            max_workers=self.max_workers,150            use_threads=self.use_threads,151            chunk_size=self.chunk_size,152            **func_kwargs,153        )154 155    def batch_map(156        self,157        func: Callable[[List[T]], List[R]],158        items: Iterable[T],159        batch_size: int = 32,160        **func_kwargs,161    ) -> List[R]:162        """Process items in batches in parallel."""163        return batch_process(164            func,165            items,166            batch_size=batch_size,167            max_workers=self.max_workers,168            **func_kwargs,169        )170