pratik-250620/MultiModal-Coherence-AI
2
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 