import asyncio from apollo_utils.core.parallel import iter_chunk from inspect import signature from typing import Callable def parallel(chunk_size: int = 1, items_kwarg: str = None) -> Callable: """deco which splits bulk data in chunks, runs wrapped function asynchronously on these chunks, returns all the responses. """ if not items_kwarg: raise ValueError(f"'items_kwarg' should be passed. Got item_kwarg={items_kwarg}") def wrapper(f: Callable): async def wrapped(*args, **kwargs): def set_items(_items, _kwargs): _kwargs[items_kwarg] = _items return _kwargs tasks = [ f(*args, **set_items(chunk, kwargs)) for chunk in iter_chunk(kwargs[items_kwarg], chunk_size=chunk_size) ] results = await asyncio.gather(*tasks) return results wrapped.__signature__ = signature(f) return wrapped return wrapper