Skip to content

voiage.parallel.distributed.distributed_reduce

distributed_reduce([positional or keyword] items: Sequence[Item] | Iterable[Item] = None, [positional or keyword] worker_func: Callable[[Item], ChunkResult] = None, [positional or keyword] reducer: Callable[[Sequence[ChunkResult]], ChunkResult] = None, [keyword-only] config: ClusterExecutionConfig | None = None, [keyword-only] n_workers: int | None = None, [keyword-only] use_processes: bool | None = None, [keyword-only] executor_factory: Callable[[int, bool], Executor] | None = None) -> ChunkResult

Map distributed work and reduce it deterministically.

Parameters:

  • items Sequence[Item] | Iterable[Item]
  • worker_func Callable[[Item], ChunkResult]
  • reducer Callable[[Sequence[ChunkResult]], ChunkResult]
  • config ClusterExecutionConfig | None (default: None)
  • n_workers int | None (default: None)
  • use_processes bool | None (default: None)
  • executor_factory Callable[[int, bool], Executor] | None (default: None)

Returns: ChunkResult