Parallel

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

source

threaded

def threaded(
    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:NoneType=None, 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:NoneType=None, daemon:bool=False
):

Like threaded(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]

source

parallel_async

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

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-17 17:39:11.566435 2026-07-17 17:39:11.569794 0.002859891538949448
3 2026-07-17 17:39:11.569932 2026-07-17 17:39:11.576193 0.005425981904664934
1 2026-07-17 17:39:11.566330 2026-07-17 17:39:11.588110 0.020615564297547868
0 2026-07-17 17:39:11.566115 2026-07-17 17:39:11.596509 0.029235003044882336
4 2026-07-17 17:39:11.576326 2026-07-17 17:39:11.606287 0.02879390679028037
5 2026-07-17 17:39:11.588278 2026-07-17 17:39:11.606438 0.017552992195368936

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