Parallel

Threading and multiprocessing functions
import asyncio
from contextlib import aclosing
from fastcore.test import *
from nbdev.showdoc import *
from fastcore.nb_imports import *

parallel(f, items) maps f over items using a process pool (threadpool=True for threads instead), with n_workers, optional progress bar, a pause between starts (to stagger e.g. web requests), and return_exceptions to collect errors instead of raising; n_workers=0 runs serially, which makes debugging easy. parallel_async is the asyncio version, using a semaphore to limit concurrency. @threaded makes a function run in a Thread (or Process) whose eventual return value lands in its result attr, and startthread/startproc start one immediately.


source

threaded

def threaded(
    f, # Function to run
    process:bool=False, # Create a Process instead of a Thread?
    daemon:bool=False, # Use daemon mode?
):

Run f in a Thread (or Process if process=True), and returns it

@threaded
def _1():
    time.sleep(0.05)
    print("second")
    return 5

@threaded
def _2():
    time.sleep(0.01)
    print("first")

a = _1()
_2()
time.sleep(0.1)
first
second

After the thread is complete, the return value is stored in the result attr.

a.result
5

Pass daemon=True to make the thread (or process) a daemon, so it won’t prevent the parent from exiting. Useful for background services like webservers, where you don’t want a still-running thread to block process shutdown.

@threaded(daemon=True)
def f(): time.sleep(0.01)

assert f().daemon

source

startthread

def startthread(
    f, daemon:bool=False
):

Like threaded, but start thread immediately

@startthread
def _():
    time.sleep(0.05)
    print("second")

@startthread
def _():
    time.sleep(0.01)
    print("first")

time.sleep(0.1)
first
second
@startthread(daemon=True)
def f(): time.sleep(0.01)

assert f.daemon

source

startproc

def startproc(
    f, daemon:bool=False
):

Like threaded(process=True), but start Process immediately

@startproc
def _():
    time.sleep(0.05)
    print("second")

@startproc
def _():
    time.sleep(0.01)
    print("first")

time.sleep(0.1)

source

parallelable

def parallelable(
    param_name, num_workers, f:NoneType=None
):

Call self as a function.


source

ThreadPoolExecutor

def ThreadPoolExecutor(
    max_workers:int=4, on_exc:builtin_function_or_method=print, pause:int=0, **kwargs
):

Same as Python’s ThreadPoolExecutor, except can pass max_workers==0 for serial execution


source

ProcessPoolExecutor

def ProcessPoolExecutor(
    max_workers:int=4, on_exc:builtin_function_or_method=print, pause:int=0, mp_context:NoneType=None,
    initializer:NoneType=None, initargs:tuple=(), max_tasks_per_child:NoneType=None
):

Same as Python’s ProcessPoolExecutor, except can pass max_workers==0 for serial execution


source

parallel

def parallel(
    f, items, *args, n_workers:int=4, total:NoneType=None, progress:NoneType=None, pause:int=0, method:NoneType=None,
    threadpool:bool=False, timeout:NoneType=None, chunksize:int=1, return_exceptions:bool=False, **kwargs
):

Applies func in parallel to items, using n_workers

inp,exp = range(50),range(1,51)

test_eq(parallel(_add_one, inp, n_workers=2), exp)
test_eq(parallel(_add_one, inp, threadpool=True, n_workers=2), exp)
test_eq(parallel(_add_one, inp, n_workers=1, a=2), range(2,52))
test_eq(parallel(_add_one, inp, n_workers=0), exp)
test_eq(parallel(_add_one, inp, n_workers=0, a=2), range(2,52))

Use the pause parameter to ensure a pause of pause seconds between processes starting. This is in case there are race conditions in starting some process, or to stagger the time each process starts, for example when making many requests to a webserver. Set threadpool=True to use ThreadPoolExecutor instead of ProcessPoolExecutor.

from datetime import datetime
def print_time(i): 
    time.sleep(random.random()/1000)
    print(i, datetime.now())

parallel(print_time, range(5), n_workers=2, pause=0.1);

You can also pass return_exceptions=True to catch any exceptions from parallel workers and return them instead:

def die_sometimes(x):
    if 3<x<6: raise Exception(f"exc: {x}")
    return x*2

parallel(die_sometimes, range(8), return_exceptions=True)
[0, 2, 4, 6, Exception('exc: 4'), Exception('exc: 5'), 12, 14]

parallel_async_gen is the streaming core of parallel_async: it starts every task up front (semaphore-limited to n_workers) and yields (index,result) pairs as tasks complete, so a consumer can act on each result the moment it’s ready. Completion order is not item order, so the index says which item each result came from. When the stream exits – exhausted, errored, or closed early – cancel_on_exit cancels unfinished tasks and waits for them; pass cancel_on_exit=False to let them run to completion in the background instead.


source

parallel_async_gen

def parallel_async_gen(
    f, items, *args, n_workers:int=16, pause:int=0, timeout:NoneType=None, return_exceptions:bool=False,
    cancel_on_exit:bool=True, **kwargs
):

Yield (index,result) pairs as f applied to each of items completes, in completion order

Results arrive in completion order – here the later items sleep less, so they come back first, and sorting by index recovers item order:

async def wait_then(i):
    await asyncio.sleep((3-i)/20)
    return i*2

res = [(i,r) async for i,r in parallel_async_gen(wait_then, range(3))]
test_eq(sorted(res), [(0,0),(1,2),(2,4)])
res

Stopping early cancels the work you no longer want. A break only runs the stream’s cleanup deterministically inside contextlib.aclosing, so wrap the generator whenever you might not exhaust it:

done = []
async def track(i):
    await asyncio.sleep(i/50)
    done.append(i)
    return i

async with aclosing(parallel_async_gen(track, range(10), n_workers=2)) as stream:
    async for i,r in stream:
        if len(done)==2: break
await asyncio.sleep(0.2)
test_eq(done, [0,1])
done

parallel_async_dict collects the whole stream into a dict. Insertion order preserves completion order, so it keeps both facts the generator produces: res[i] addresses a result by item, and list(res) is the order tasks finished. @delegates means its signature carries the generator’s options without restating them.


source

parallel_async_dict

async def parallel_async_dict(
    f, items, *args, n_workers:int=16, pause:int=0, timeout:NoneType=None, return_exceptions:bool=False,
    cancel_on_exit:bool=True
):

Apply f to items in parallel, returning {index: result} in completion order

res = await parallel_async_dict(wait_then, range(3))
test_eq(res, {0:0, 1:2, 2:4})
test_eq(list(res), [2,1,0])
res

source

parallel_async

async def parallel_async(
    f, items, *args, cancel_on_error:bool=False, n_workers:int=16, pause:int=0, timeout:NoneType=None,
    return_exceptions:bool=False
):

Applies f to items in parallel using asyncio and a semaphore to limit concurrency.

async def print_time_async(i): 
    start =datetime.now()
    wait = random.random()/30
    await asyncio.sleep(wait)
    print(i, start, datetime.now(), wait)
    if i==5: raise Exception(f"exc {i}")
    return i

res = await parallel_async(print_time_async, range(6), n_workers=3, return_exceptions=True)
test_eq(res[:5], [0, 1, 2, 3, 4])
test_eq(type(res[5]), Exception)
2 2026-07-29 10:15:51.983068 2026-07-29 10:15:51.992934 0.008671375369168333
1 2026-07-29 10:15:51.982961 2026-07-29 10:15:51.993273 0.009861243238330394
3 2026-07-29 10:15:51.993146 2026-07-29 10:15:52.000805 0.006693202182398147
0 2026-07-29 10:15:51.982728 2026-07-29 10:15:52.001106 0.018118661669984538
5 2026-07-29 10:15:52.000973 2026-07-29 10:15:52.024793 0.02267549253732972
4 2026-07-29 10:15:51.993397 2026-07-29 10:15:52.026390 0.03271171331399413

Adding pause ensures a gap between starts:

await parallel_async(print_time_async, range(6), n_workers=3, pause=0.1, return_exceptions=True);

With cancel_on_error=True, the first failure cancels all remaining tasks, waits for them to finish cancelling, then re-raises that first exception. (Earlier versions raised an ExceptionGroup here; now you catch the original exception directly.)

async def maybe_fail(i:int):
    "Double i unless it's 3, in which case fail"
    await asyncio.sleep(random.random()/50)
    if i==3: raise ValueError(f"bad: {i}")
    return i*2
with expect_fail(ValueError, contains='bad: 3'): await parallel_async(maybe_fail, range(6), n_workers=3, cancel_on_error=True)

With return_exceptions=False, an exception is raised on error:

with expect_fail(ValueError): await parallel_async(maybe_fail, range(6), n_workers=3)

source

bg_task

def bg_task(
    coro
):

Like asyncio.create_task but logs exceptions for fire-and-forget tasks

async def _ok(): return 42
async def _fail(): raise ValueError("this error will be printed")

t1 = bg_task(_ok())
t2 = bg_task(_fail())
await asyncio.sleep(0.01)
test_eq(t1.result(), 42)
Traceback (most recent call last):
  File "<ipython-input-1-48a55f4f8ca9>", line 2, in _fail
    async def _fail(): raise ValueError("this error will be printed")
                       ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
ValueError: this error will be printed