hojichar.core.parallel
1from __future__ import annotations 2 3import functools 4import logging 5import multiprocessing 6import os 7import signal 8import threading 9from copy import copy 10from multiprocessing.pool import Pool 11from typing import Iterator, List 12 13import hojichar 14from hojichar.core import inspection 15from hojichar.core.models import Statistics 16 17logger = logging.getLogger(__name__) 18 19 20_START_METHOD_ENV_VAR = "HOJICHAR_MP_START_METHOD" 21 22PARALLEL_BASE_FILTER: hojichar.Compose 23WORKER_PARAM_IGNORE_ERRORS: bool 24 25 26def _get_parallel_context() -> multiprocessing.context.BaseContext: 27 """Return the multiprocessing context used by :class:`Parallel`. 28 29 If no global start method has been configured, ``fork`` is selected when it 30 is available so worker processes can inherit filters which cannot be pickled. 31 The environment variable takes precedence over the global setting; 32 ``default`` delegates the choice back to Python's global/default context. 33 """ 34 configured_method = os.getenv(_START_METHOD_ENV_VAR) 35 if configured_method is None: 36 start_method = multiprocessing.get_start_method(allow_none=True) 37 if start_method is None and "fork" in multiprocessing.get_all_start_methods(): 38 start_method = "fork" 39 elif configured_method == "default": 40 start_method = None 41 else: 42 start_method = configured_method 43 44 try: 45 return multiprocessing.get_context(start_method) 46 except ValueError as error: 47 supported_methods = ["default", *multiprocessing.get_all_start_methods()] 48 raise ValueError( 49 f"Invalid {_START_METHOD_ENV_VAR}={configured_method!r}. " 50 f"Choose one of: {', '.join(supported_methods)}." 51 ) from error 52 53 54def _init_worker(filter: hojichar.Compose, ignore_errors: bool) -> None: 55 signal.signal(signal.SIGINT, signal.SIG_IGN) 56 global PARALLEL_BASE_FILTER, WORKER_PARAM_IGNORE_ERRORS 57 PARALLEL_BASE_FILTER = hojichar.Compose(copy(filter.filters)) # TODO random state treatment 58 WORKER_PARAM_IGNORE_ERRORS = ignore_errors 59 60 61def _worker( 62 doc: hojichar.Document, 63) -> tuple[hojichar.Document, int, List[Statistics], str | None]: 64 global PARALLEL_BASE_FILTER, WORKER_PARAM_IGNORE_ERRORS 65 ignore_errors = WORKER_PARAM_IGNORE_ERRORS 66 error_message = None 67 try: 68 result = PARALLEL_BASE_FILTER.apply(doc) 69 except Exception as e: 70 if ignore_errors: 71 logger.error(e) 72 error_message = str(e) 73 result = hojichar.Document("", is_rejected=True) 74 else: 75 raise e # If we're not ignoring errors, let this one propagate 76 return result, os.getpid(), PARALLEL_BASE_FILTER.get_total_statistics(), error_message 77 78 79class _InFlightGate: 80 """Bounds documents drawn from the input but not yet returned to the caller. 81 82 ``feed`` wraps the input iterator and acquires one permit *before* each 83 document is drawn (acquiring afterwards would hold one extra pre-fetched 84 document while waiting). ``imap_apply`` releases the permit only after the 85 corresponding result has been handed back to the caller, so at most 86 ``max_in_flight`` documents exist anywhere between the input iterator and 87 the caller at any moment. 88 89 ``feed`` runs inside Pool's task-handler thread. The polling acquire 90 observes :meth:`stop` before and after each successful acquire (handing 91 the permit back if stopped), and abnormal exits stop the gate before 92 returning their permit, so the feeder never draws from the source after 93 consumption ended and pool shutdown never blocks on the gate. 94 """ 95 96 def __init__(self, max_in_flight: int) -> None: 97 self._semaphore = threading.Semaphore(max_in_flight) 98 self._stop_feeding = threading.Event() 99 100 def feed(self, docs: Iterator[hojichar.Document]) -> Iterator[hojichar.Document]: 101 iterator = iter(docs) 102 while self._acquire(): 103 try: 104 doc = next(iterator) 105 except StopIteration: 106 self.release() 107 return 108 yield doc 109 110 def release(self) -> None: 111 self._semaphore.release() 112 113 def stop(self) -> None: 114 """Unblock the feeder; called when consumption ends and at shutdown.""" 115 self._stop_feeding.set() 116 117 def _acquire(self) -> bool: 118 while not self._stop_feeding.is_set(): 119 if self._semaphore.acquire(timeout=0.1): 120 if self._stop_feeding.is_set(): 121 # Stopped while blocked on (or right after) the acquire: 122 # hand the permit back instead of drawing one more 123 # document from a possibly blocking source. 124 self._semaphore.release() 125 return False 126 return True 127 return False 128 129 130class Parallel: 131 """ 132 The Parallel class provides a way to apply a hojichar.Compose filter 133 to an iterator of documents in a parallel manner using a specified 134 number of worker processes. This class should be used as a context 135 manager with a 'with' statement. 136 137 When no global start method has been configured, Parallel uses ``fork`` on 138 platforms which support it so unpicklable filters can be inherited by 139 workers. Set ``HOJICHAR_MP_START_METHOD`` to explicitly choose ``fork``, 140 ``spawn``, ``forkserver``, or ``default``. Non-fork contexts require the 141 Compose object and its filters to be picklable. 142 143 Example: 144 145 doc_iter = (hojichar.Document(d) for d in open("my_text.txt")) 146 with Parallel(my_filter, num_jobs=8) as pfilter: 147 for doc in pfilter.imap_apply(doc_iter): 148 pass # Process the filtered document as needed. 149 """ 150 151 def __init__( 152 self, 153 filter: hojichar.Compose, 154 num_jobs: int | None = None, 155 ignore_errors: bool = False, 156 ordered: bool = False, 157 max_in_flight: int | None = None, 158 ): 159 """ 160 Initializes a new instance of the Parallel class. 161 162 Args: 163 filter (hojichar.Compose): A composed filter object that specifies the 164 processing operations to apply to each document in parallel. 165 A copy of the filter is made within a 'with' statement. When the 'with' 166 block terminates,the statistical information obtained through `filter.statistics` 167 or`filter.statistics_obj` is replaced with the total value of the statistical 168 information processed within the 'with' block. 169 170 num_jobs (int | None, optional): The number of worker processes to use. 171 If None, then the number returned by os.cpu_count() is used. Defaults to None. 172 ignore_errors (bool, optional): If set to True, any exceptions thrown during 173 the processing of a document will be caught and logged, but will not 174 stop the processing of further documents. If set to False, the first 175 exception thrown will terminate the entire parallel processing operation. 176 Defaults to False. 177 ordered (bool, optional): If set to True, processed documents are yielded in 178 the same order as the input documents. If set to False, documents are 179 yielded as soon as their processing completes. Defaults to False. 180 max_in_flight (int | None, optional): Strict upper bound on the number of 181 documents drawn from the input iterator whose results have not yet been 182 handed back to the caller. A yielded document keeps its permit until the 183 caller requests the next one, so with `max_in_flight=1` the pool holds a 184 single document end to end (no pipelining). Without a bound, Pool's 185 task-handler thread drains the input as fast as the worker pipe accepts, 186 so a producer that outruns the filters can buffer a large number of 187 documents. If None, no explicit bound is applied. Defaults to None. 188 """ 189 if max_in_flight is not None and max_in_flight < 1: 190 raise ValueError("max_in_flight must be at least 1") 191 self.filter = filter 192 self.num_jobs = num_jobs 193 self.ignore_errors = ignore_errors 194 self.ordered = ordered 195 self.max_in_flight = max_in_flight 196 197 self._pool: Pool | None = None 198 self._pid_stats: dict[int, List[Statistics]] | None = None 199 self._gates: list[_InFlightGate] = [] 200 201 def __enter__(self) -> Parallel: 202 context = _get_parallel_context() 203 self._pool = context.Pool( 204 processes=self.num_jobs, 205 initializer=_init_worker, 206 initargs=(self.filter, self.ignore_errors), 207 ) 208 self._pid_stats = dict() 209 self._gates = [] 210 return self 211 212 def imap_apply(self, docs: Iterator[hojichar.Document]) -> Iterator[hojichar.Document]: 213 """ 214 Takes an iterator of Documents and applies the Compose filter to 215 each Document in a parallel manner. This is a generator method 216 that yields processed Documents. 217 218 Args: 219 docs (Iterator[hojichar.Document]): An iterator of Documents to be processed. 220 221 Raises: 222 RuntimeError: If the Parallel instance is not properly initialized. This 223 generally happens when the method is called outside of a 'with' statement. 224 Exception: If any exceptions are raised within the worker processes. 225 226 Yields: 227 Iterator[hojichar.Document]: An iterator that yields processed Documents. 228 """ 229 if self._pool is None or self._pid_stats is None: 230 raise RuntimeError( 231 "Parallel instance not properly initialized. Use within a 'with' statement." 232 ) 233 gate: _InFlightGate | None = None 234 if self.max_in_flight is not None: 235 gate = _InFlightGate(self.max_in_flight) 236 self._gates.append(gate) 237 docs = gate.feed(docs) 238 try: 239 results = ( 240 self._pool.imap(_worker, docs) 241 if self.ordered 242 else self._pool.imap_unordered(_worker, docs) 243 ) 244 for doc, pid, stat, err_msg in results: 245 self._pid_stats[pid] = stat 246 if err_msg is not None: 247 logger.error(f"Error in worker {pid}: {err_msg}") 248 # The permit is returned only once the caller has taken the 249 # document: releasing before the yield would let the feeder 250 # momentarily draw max_in_flight + 1 documents. On an abnormal 251 # exit (close/throw at the yield point) the gate is stopped 252 # *before* the permit is returned, so a feeder woken by the 253 # release always observes the stop and cannot draw again. 254 try: 255 yield doc 256 except BaseException: 257 if gate is not None: 258 gate.stop() 259 gate.release() 260 raise 261 else: 262 if gate is not None: 263 gate.release() 264 except Exception: 265 self.__exit__(None, None, None) 266 raise 267 finally: 268 if gate is not None: 269 gate.stop() 270 271 def __exit__(self, exc_type, exc_value, traceback) -> None: # type: ignore 272 # Feeders blocked on a max_in_flight gate must exit before the pool is 273 # joined, or shutdown would wait on them forever. 274 for gate in self._gates: 275 gate.stop() 276 if self._pool: 277 self._pool.terminate() 278 self._pool.join() 279 if self._pid_stats: 280 total_stats = functools.reduce( 281 lambda x, y: Statistics.add_list_of_stats(x, y), self._pid_stats.values() 282 ) 283 self.filter._statistics.update(Statistics.get_filter("Total", total_stats)) 284 for stat in total_stats: 285 for filt in self.filter.filters: 286 if stat.name == filt.name: 287 filt._statistics.update(stat) 288 break 289 290 def get_total_statistics(self) -> List[Statistics]: 291 """ 292 Returns a statistics object of the total statistical 293 values processed within the Parallel block. 294 295 Returns: 296 StatsContainer: Statistics object 297 """ 298 if self._pid_stats: 299 total_stats = functools.reduce( 300 lambda x, y: Statistics.add_list_of_stats(x, y), self._pid_stats.values() 301 ) 302 return total_stats 303 else: 304 return [] 305 306 def get_total_statistics_map(self) -> List[dict]: 307 return [stat.to_dict() for stat in self.get_total_statistics()] 308 309 @property 310 def statistics_obj(self) -> inspection.StatsContainer: 311 """ 312 Returns the statistics object of the Parallel instance. 313 This is a StatsContainer object which contains the statistics 314 of the Parallel instance and sub filters. 315 316 Returns: 317 StatsContainer: Statistics object 318 """ 319 return inspection.statistics_obj_adapter(self.get_total_statistics()) # type: ignore
131class Parallel: 132 """ 133 The Parallel class provides a way to apply a hojichar.Compose filter 134 to an iterator of documents in a parallel manner using a specified 135 number of worker processes. This class should be used as a context 136 manager with a 'with' statement. 137 138 When no global start method has been configured, Parallel uses ``fork`` on 139 platforms which support it so unpicklable filters can be inherited by 140 workers. Set ``HOJICHAR_MP_START_METHOD`` to explicitly choose ``fork``, 141 ``spawn``, ``forkserver``, or ``default``. Non-fork contexts require the 142 Compose object and its filters to be picklable. 143 144 Example: 145 146 doc_iter = (hojichar.Document(d) for d in open("my_text.txt")) 147 with Parallel(my_filter, num_jobs=8) as pfilter: 148 for doc in pfilter.imap_apply(doc_iter): 149 pass # Process the filtered document as needed. 150 """ 151 152 def __init__( 153 self, 154 filter: hojichar.Compose, 155 num_jobs: int | None = None, 156 ignore_errors: bool = False, 157 ordered: bool = False, 158 max_in_flight: int | None = None, 159 ): 160 """ 161 Initializes a new instance of the Parallel class. 162 163 Args: 164 filter (hojichar.Compose): A composed filter object that specifies the 165 processing operations to apply to each document in parallel. 166 A copy of the filter is made within a 'with' statement. When the 'with' 167 block terminates,the statistical information obtained through `filter.statistics` 168 or`filter.statistics_obj` is replaced with the total value of the statistical 169 information processed within the 'with' block. 170 171 num_jobs (int | None, optional): The number of worker processes to use. 172 If None, then the number returned by os.cpu_count() is used. Defaults to None. 173 ignore_errors (bool, optional): If set to True, any exceptions thrown during 174 the processing of a document will be caught and logged, but will not 175 stop the processing of further documents. If set to False, the first 176 exception thrown will terminate the entire parallel processing operation. 177 Defaults to False. 178 ordered (bool, optional): If set to True, processed documents are yielded in 179 the same order as the input documents. If set to False, documents are 180 yielded as soon as their processing completes. Defaults to False. 181 max_in_flight (int | None, optional): Strict upper bound on the number of 182 documents drawn from the input iterator whose results have not yet been 183 handed back to the caller. A yielded document keeps its permit until the 184 caller requests the next one, so with `max_in_flight=1` the pool holds a 185 single document end to end (no pipelining). Without a bound, Pool's 186 task-handler thread drains the input as fast as the worker pipe accepts, 187 so a producer that outruns the filters can buffer a large number of 188 documents. If None, no explicit bound is applied. Defaults to None. 189 """ 190 if max_in_flight is not None and max_in_flight < 1: 191 raise ValueError("max_in_flight must be at least 1") 192 self.filter = filter 193 self.num_jobs = num_jobs 194 self.ignore_errors = ignore_errors 195 self.ordered = ordered 196 self.max_in_flight = max_in_flight 197 198 self._pool: Pool | None = None 199 self._pid_stats: dict[int, List[Statistics]] | None = None 200 self._gates: list[_InFlightGate] = [] 201 202 def __enter__(self) -> Parallel: 203 context = _get_parallel_context() 204 self._pool = context.Pool( 205 processes=self.num_jobs, 206 initializer=_init_worker, 207 initargs=(self.filter, self.ignore_errors), 208 ) 209 self._pid_stats = dict() 210 self._gates = [] 211 return self 212 213 def imap_apply(self, docs: Iterator[hojichar.Document]) -> Iterator[hojichar.Document]: 214 """ 215 Takes an iterator of Documents and applies the Compose filter to 216 each Document in a parallel manner. This is a generator method 217 that yields processed Documents. 218 219 Args: 220 docs (Iterator[hojichar.Document]): An iterator of Documents to be processed. 221 222 Raises: 223 RuntimeError: If the Parallel instance is not properly initialized. This 224 generally happens when the method is called outside of a 'with' statement. 225 Exception: If any exceptions are raised within the worker processes. 226 227 Yields: 228 Iterator[hojichar.Document]: An iterator that yields processed Documents. 229 """ 230 if self._pool is None or self._pid_stats is None: 231 raise RuntimeError( 232 "Parallel instance not properly initialized. Use within a 'with' statement." 233 ) 234 gate: _InFlightGate | None = None 235 if self.max_in_flight is not None: 236 gate = _InFlightGate(self.max_in_flight) 237 self._gates.append(gate) 238 docs = gate.feed(docs) 239 try: 240 results = ( 241 self._pool.imap(_worker, docs) 242 if self.ordered 243 else self._pool.imap_unordered(_worker, docs) 244 ) 245 for doc, pid, stat, err_msg in results: 246 self._pid_stats[pid] = stat 247 if err_msg is not None: 248 logger.error(f"Error in worker {pid}: {err_msg}") 249 # The permit is returned only once the caller has taken the 250 # document: releasing before the yield would let the feeder 251 # momentarily draw max_in_flight + 1 documents. On an abnormal 252 # exit (close/throw at the yield point) the gate is stopped 253 # *before* the permit is returned, so a feeder woken by the 254 # release always observes the stop and cannot draw again. 255 try: 256 yield doc 257 except BaseException: 258 if gate is not None: 259 gate.stop() 260 gate.release() 261 raise 262 else: 263 if gate is not None: 264 gate.release() 265 except Exception: 266 self.__exit__(None, None, None) 267 raise 268 finally: 269 if gate is not None: 270 gate.stop() 271 272 def __exit__(self, exc_type, exc_value, traceback) -> None: # type: ignore 273 # Feeders blocked on a max_in_flight gate must exit before the pool is 274 # joined, or shutdown would wait on them forever. 275 for gate in self._gates: 276 gate.stop() 277 if self._pool: 278 self._pool.terminate() 279 self._pool.join() 280 if self._pid_stats: 281 total_stats = functools.reduce( 282 lambda x, y: Statistics.add_list_of_stats(x, y), self._pid_stats.values() 283 ) 284 self.filter._statistics.update(Statistics.get_filter("Total", total_stats)) 285 for stat in total_stats: 286 for filt in self.filter.filters: 287 if stat.name == filt.name: 288 filt._statistics.update(stat) 289 break 290 291 def get_total_statistics(self) -> List[Statistics]: 292 """ 293 Returns a statistics object of the total statistical 294 values processed within the Parallel block. 295 296 Returns: 297 StatsContainer: Statistics object 298 """ 299 if self._pid_stats: 300 total_stats = functools.reduce( 301 lambda x, y: Statistics.add_list_of_stats(x, y), self._pid_stats.values() 302 ) 303 return total_stats 304 else: 305 return [] 306 307 def get_total_statistics_map(self) -> List[dict]: 308 return [stat.to_dict() for stat in self.get_total_statistics()] 309 310 @property 311 def statistics_obj(self) -> inspection.StatsContainer: 312 """ 313 Returns the statistics object of the Parallel instance. 314 This is a StatsContainer object which contains the statistics 315 of the Parallel instance and sub filters. 316 317 Returns: 318 StatsContainer: Statistics object 319 """ 320 return inspection.statistics_obj_adapter(self.get_total_statistics()) # type: ignore
The Parallel class provides a way to apply a hojichar.Compose filter to an iterator of documents in a parallel manner using a specified number of worker processes. This class should be used as a context manager with a 'with' statement.
When no global start method has been configured, Parallel uses fork on
platforms which support it so unpicklable filters can be inherited by
workers. Set HOJICHAR_MP_START_METHOD to explicitly choose fork,
spawn, forkserver, or default. Non-fork contexts require the
Compose object and its filters to be picklable.
Example:
doc_iter = (hojichar.Document(d) for d in open("my_text.txt")) with Parallel(my_filter, num_jobs=8) as pfilter: for doc in pfilter.imap_apply(doc_iter): pass # Process the filtered document as needed.
152 def __init__( 153 self, 154 filter: hojichar.Compose, 155 num_jobs: int | None = None, 156 ignore_errors: bool = False, 157 ordered: bool = False, 158 max_in_flight: int | None = None, 159 ): 160 """ 161 Initializes a new instance of the Parallel class. 162 163 Args: 164 filter (hojichar.Compose): A composed filter object that specifies the 165 processing operations to apply to each document in parallel. 166 A copy of the filter is made within a 'with' statement. When the 'with' 167 block terminates,the statistical information obtained through `filter.statistics` 168 or`filter.statistics_obj` is replaced with the total value of the statistical 169 information processed within the 'with' block. 170 171 num_jobs (int | None, optional): The number of worker processes to use. 172 If None, then the number returned by os.cpu_count() is used. Defaults to None. 173 ignore_errors (bool, optional): If set to True, any exceptions thrown during 174 the processing of a document will be caught and logged, but will not 175 stop the processing of further documents. If set to False, the first 176 exception thrown will terminate the entire parallel processing operation. 177 Defaults to False. 178 ordered (bool, optional): If set to True, processed documents are yielded in 179 the same order as the input documents. If set to False, documents are 180 yielded as soon as their processing completes. Defaults to False. 181 max_in_flight (int | None, optional): Strict upper bound on the number of 182 documents drawn from the input iterator whose results have not yet been 183 handed back to the caller. A yielded document keeps its permit until the 184 caller requests the next one, so with `max_in_flight=1` the pool holds a 185 single document end to end (no pipelining). Without a bound, Pool's 186 task-handler thread drains the input as fast as the worker pipe accepts, 187 so a producer that outruns the filters can buffer a large number of 188 documents. If None, no explicit bound is applied. Defaults to None. 189 """ 190 if max_in_flight is not None and max_in_flight < 1: 191 raise ValueError("max_in_flight must be at least 1") 192 self.filter = filter 193 self.num_jobs = num_jobs 194 self.ignore_errors = ignore_errors 195 self.ordered = ordered 196 self.max_in_flight = max_in_flight 197 198 self._pool: Pool | None = None 199 self._pid_stats: dict[int, List[Statistics]] | None = None 200 self._gates: list[_InFlightGate] = []
Initializes a new instance of the Parallel class.
Args:
filter (hojichar.Compose): A composed filter object that specifies the
processing operations to apply to each document in parallel.
A copy of the filter is made within a 'with' statement. When the 'with'
block terminates,the statistical information obtained through filter.statistics
orfilter.statistics_obj is replaced with the total value of the statistical
information processed within the 'with' block.
num_jobs (int | None, optional): The number of worker processes to use.
If None, then the number returned by os.cpu_count() is used. Defaults to None.
ignore_errors (bool, optional): If set to True, any exceptions thrown during
the processing of a document will be caught and logged, but will not
stop the processing of further documents. If set to False, the first
exception thrown will terminate the entire parallel processing operation.
Defaults to False.
ordered (bool, optional): If set to True, processed documents are yielded in
the same order as the input documents. If set to False, documents are
yielded as soon as their processing completes. Defaults to False.
max_in_flight (int | None, optional): Strict upper bound on the number of
documents drawn from the input iterator whose results have not yet been
handed back to the caller. A yielded document keeps its permit until the
caller requests the next one, so with `max_in_flight=1` the pool holds a
single document end to end (no pipelining). Without a bound, Pool's
task-handler thread drains the input as fast as the worker pipe accepts,
so a producer that outruns the filters can buffer a large number of
documents. If None, no explicit bound is applied. Defaults to None.
213 def imap_apply(self, docs: Iterator[hojichar.Document]) -> Iterator[hojichar.Document]: 214 """ 215 Takes an iterator of Documents and applies the Compose filter to 216 each Document in a parallel manner. This is a generator method 217 that yields processed Documents. 218 219 Args: 220 docs (Iterator[hojichar.Document]): An iterator of Documents to be processed. 221 222 Raises: 223 RuntimeError: If the Parallel instance is not properly initialized. This 224 generally happens when the method is called outside of a 'with' statement. 225 Exception: If any exceptions are raised within the worker processes. 226 227 Yields: 228 Iterator[hojichar.Document]: An iterator that yields processed Documents. 229 """ 230 if self._pool is None or self._pid_stats is None: 231 raise RuntimeError( 232 "Parallel instance not properly initialized. Use within a 'with' statement." 233 ) 234 gate: _InFlightGate | None = None 235 if self.max_in_flight is not None: 236 gate = _InFlightGate(self.max_in_flight) 237 self._gates.append(gate) 238 docs = gate.feed(docs) 239 try: 240 results = ( 241 self._pool.imap(_worker, docs) 242 if self.ordered 243 else self._pool.imap_unordered(_worker, docs) 244 ) 245 for doc, pid, stat, err_msg in results: 246 self._pid_stats[pid] = stat 247 if err_msg is not None: 248 logger.error(f"Error in worker {pid}: {err_msg}") 249 # The permit is returned only once the caller has taken the 250 # document: releasing before the yield would let the feeder 251 # momentarily draw max_in_flight + 1 documents. On an abnormal 252 # exit (close/throw at the yield point) the gate is stopped 253 # *before* the permit is returned, so a feeder woken by the 254 # release always observes the stop and cannot draw again. 255 try: 256 yield doc 257 except BaseException: 258 if gate is not None: 259 gate.stop() 260 gate.release() 261 raise 262 else: 263 if gate is not None: 264 gate.release() 265 except Exception: 266 self.__exit__(None, None, None) 267 raise 268 finally: 269 if gate is not None: 270 gate.stop()
Takes an iterator of Documents and applies the Compose filter to each Document in a parallel manner. This is a generator method that yields processed Documents.
Args: docs (Iterator[hojichar.Document]): An iterator of Documents to be processed.
Raises: RuntimeError: If the Parallel instance is not properly initialized. This generally happens when the method is called outside of a 'with' statement. Exception: If any exceptions are raised within the worker processes.
Yields: Iterator[hojichar.Document]: An iterator that yields processed Documents.
291 def get_total_statistics(self) -> List[Statistics]: 292 """ 293 Returns a statistics object of the total statistical 294 values processed within the Parallel block. 295 296 Returns: 297 StatsContainer: Statistics object 298 """ 299 if self._pid_stats: 300 total_stats = functools.reduce( 301 lambda x, y: Statistics.add_list_of_stats(x, y), self._pid_stats.values() 302 ) 303 return total_stats 304 else: 305 return []
Returns a statistics object of the total statistical values processed within the Parallel block.
Returns: StatsContainer: Statistics object
Returns the statistics object of the Parallel instance. This is a StatsContainer object which contains the statistics of the Parallel instance and sub filters.
Returns: StatsContainer: Statistics object