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
class Parallel:
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.

Parallel( filter: hojichar.core.composition.Compose, num_jobs: int | None = None, ignore_errors: bool = False, ordered: bool = False, max_in_flight: int | None = None)
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.
def imap_apply( self, docs: Iterator[hojichar.core.models.Document]) -> Iterator[hojichar.core.models.Document]:
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.

def get_total_statistics(self) -> List[hojichar.core.models.Statistics]:
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

def get_total_statistics_map(self) -> List[dict]:
307    def get_total_statistics_map(self) -> List[dict]:
308        return [stat.to_dict() for stat in self.get_total_statistics()]

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