../

Async & concurrency

Python 3.14's concurrency toolkit: asyncio for I/O, threads, processes, subinterpreters and the free-threaded build, and how to pick between them. The JS side lives in Async & Promises.

Choosing a model

ModelBest forCPU in parallel?Sharing dataCost
asynciothousands of I/O waits (HTTP, sockets, DB)no, one threadplain objects, no locks between awaitscheapest; needs async libraries
threading / ThreadPoolExecutorblocking I/O through sync librariesno with the GIL; yes on 3.14tshared memory, needs locksone OS thread each
multiprocessing / ProcessPoolExecutorCPU-bound pure Pythonyespickled messages, shared_memoryprocess start-up, pickling
concurrent.interpreters / InterpreterPoolExecutorCPU-bound, isolated workyes, one GIL per interpreterqueues of shareable objectslighter than processes; young, extension support varies
Free-threaded build (python3.14t) + threadsCPU-bound with shared datayesshared memory, needs locksa few % slower single-threaded; needs compatible wheels
NumPy, hashlib, zlib, …heavy C workoften: they release the GILarrays/buffersnone: just use threads
Waiting on network/disk?
  many connections, async libs exist  -> asyncio
  sync library (requests, boto3, DB-API) -> thread pool
Burning CPU in Python code?
  -> ProcessPoolExecutor (safe default)
  -> InterpreterPoolExecutor / 3.14t threads (3.14+)
Both?  asyncio + loop.run_in_executor(process_pool, fn)

Event loop & asyncio.run

One thread, one loop. A coroutine runs until it hits an await on something not ready, then the loop runs whatever else is ready. Nothing preempts a coroutine between awaits.

Entry pointUse
asyncio.run(main())the usual entry: new loop, run, cancel leftovers, close
asyncio.run(main(), debug=True)slow-callback warnings, never-awaited tracebacks
asyncio.run(main(), loop_factory=uvloop.new_event_loop)custom loop (3.12+); uvloop for speed
with asyncio.Runner() as r:several r.run(coro) calls on one loop (tests, REPL tools)
asyncio.get_running_loop()the loop, from inside a coroutine or callback
asyncio.get_event_loop()avoid: 3.14 raises RuntimeError when no loop is set
python -m asyncioREPL with top-level await
PYTHONASYNCIODEBUG=1debug mode from the environment
main.py
import asyncio
 
 
async def greet(name: str, delay: float) -> str:
    await asyncio.sleep(delay)  # yields to the loop
    return f"hi {name}"
 
 
async def main() -> None:
    msg = await greet("ada", 0.1)
    print(msg)
 
 
if __name__ == "__main__":
    asyncio.run(main())

Event loop policies (set_event_loop_policy and friends) are deprecated in 3.14 and go in 3.16: pass loop_factory instead.

Coroutines, tasks & futures

ThingWhat it isCreated by
Coroutine functionasync def f()you
Coroutine objectf(): a paused body; runs only when awaited or wrapped in a taskcalling it
Taska scheduled coroutine; starts on the next loop turnasyncio.create_task, TaskGroup.create_task
Futurelow-level "result later" slot; tasks subclass itloop.create_future()
Awaitableanything with __await__: coroutines, tasks, futures
import asyncio
 
 
async def fetch(n: int) -> int:
    await asyncio.sleep(0.1)
    return n * 2
 
 
async def main() -> None:
    # sequential: ~0.2 s, awaiting a coroutine runs it now
    a = await fetch(1)
    b = await fetch(2)
 
    # concurrent: ~0.1 s, tasks start on the next await
    t1 = asyncio.create_task(fetch(3), name="f3")
    t2 = asyncio.create_task(fetch(4))
    c, d = await t1, await t2
    print(a, b, c, d)
 
asyncio.run(main())
Task APIDoes
t.done(), t.cancelled()state checks
t.result(), t.exception()outcome of a finished task (raises if not done)
t.cancel(msg=None)request cancellation (see below)
t.add_done_callback(fn)fn(task) when it finishes
t.get_name(), t.set_name()shows in reprs and python -m asyncio ps
create_task(coro, eager_start=True)3.14: run synchronously until the first real suspension
asyncio.current_task()the task running this code

TaskGroup vs gather

asyncio.TaskGroup (3.11+) is structured concurrency: every child finishes (or is canceled) before the async with exits. Prefer it for new code.

TaskGroupgather(*aws)gather(..., return_exceptions=True)
A child raisescancels the siblings, raises ExceptionGroupfirst exception propagates; siblings keep runningexception objects returned in the list
Resultstask.result() after the blocklist in argument orderlist of values and exceptions
Outer cancelcancels all childrencancels all childrencancels all children
JS analoguePromise.all + abort on failurePromise.allPromise.allSettled
import asyncio
 
 
async def work(n: int) -> int:
    await asyncio.sleep(0.01 * n)
    if n == 3:
        raise ValueError(f"bad {n}")
    return n
 
 
async def main() -> None:
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(work(n)) for n in (1, 2)]
    print([t.result() for t in tasks])  # [1, 2]
 
    try:
        async with asyncio.TaskGroup() as tg:
            for n in (1, 3, 5):
                tg.create_task(work(n))
    except* ValueError as eg:  # ExceptionGroup
        print("failed:", eg.exceptions)
 
    got = await asyncio.gather(
        work(1), work(3), return_exceptions=True
    )
    print(got)  # [1, ValueError('bad 3')]
 
asyncio.run(main())

Other ways to wait

APIReturns
async for t in asyncio.as_completed(aws)tasks in finish order (async-iterable since 3.13)
await asyncio.wait(tasks, return_when=FIRST_COMPLETED)(done, pending) sets; also FIRST_EXCEPTION, ALL_COMPLETED
await asyncio.wait_for(aw, timeout)the result, or TimeoutError (prefer asyncio.timeout)
await asyncio.shield(aw)protects aw from the caller's cancellation

Timeouts & cancellation

Canceling a task throws asyncio.CancelledError into it at its current await. It subclasses BaseException, so except Exception does not catch it.

import asyncio
 
 
async def slow() -> str:
    try:
        await asyncio.sleep(10)
        return "done"
    except asyncio.CancelledError:
        print("cleaning up")
        raise  # always re-raise
    finally:
        print("finally runs too")
 
 
async def main() -> None:
    try:
        async with asyncio.timeout(0.1):  # 3.11+
            await slow()
    except TimeoutError:  # builtin since 3.11
        print("timed out")
 
    t = asyncio.create_task(slow())
    await asyncio.sleep(0)
    t.cancel("shutting down")
    try:
        await t
    except asyncio.CancelledError:
        print("canceled:", t.cancelled())
 
asyncio.run(main())
APINotes
asyncio.timeout(sec)context manager; None means no limit; cm.reschedule(when) to move it
asyncio.timeout_at(when)absolute deadline in loop.time() units
task.cancel()a request: the task can clean up, and could (wrongly) swallow it
task.uncancel() / task.cancelling()for libraries that suppress one cancellation deliberately
asyncio.shield(aw)outer cancel doesn't reach aw; it keeps running
try / finallythe cleanup path on cancel, timeout and error alike

Queues & synchronization

All asyncio primitives are for tasks on one loop; they are not thread-safe. Across threads use queue.Queue / threading primitives.

PrimitiveUse
asyncio.Queue(maxsize)FIFO; put waits when full, get waits when empty
PriorityQueue, LifoQueuesmallest-first, stack order
q.task_done() + await q.join()wait until every item is processed
q.shutdown(immediate=False)3.13+: put raises QueueShutDown; get raises once drained
asyncio.Lock()async with lock: one task at a time
asyncio.Semaphore(n)at most n tasks inside; BoundedSemaphore errors on extra release
asyncio.Event()set() / await wait() / clear(): one-shot signal
asyncio.Condition()wait_for(predicate) + notify()
asyncio.Barrier(n)3.11+: n tasks wait for each other
import asyncio
 
lock = asyncio.Lock()
balance = 0
 
 
async def deposit(amount: int) -> None:
    global balance
    async with lock:  # read-modify-write across an await
        current = balance
        await asyncio.sleep(0)
        balance = current + amount
 
 
async def main() -> None:
    async with asyncio.TaskGroup() as tg:
        for _ in range(100):
            tg.create_task(deposit(1))
    print(balance)  # 100; without the lock: 1
 
asyncio.run(main())

Async iterators & context managers

ConstructProtocol
async for x in it__aiter__() returns an object with async __anext__(); StopAsyncIteration ends it
async def with yieldasync generator, typed AsyncIterator[T] or AsyncGenerator[T, None]
[x async for x in it]async comprehension (inside async def)
async with cmasync __aenter__ / async __aexit__
@contextlib.asynccontextmanagerbuild one from an async generator
contextlib.aclosing(gen)close an async generator deterministically
contextlib.AsyncExitStacka dynamic number of async context managers
anext(it, default), aiter(it)builtins (3.10+)
import asyncio
from collections.abc import AsyncGenerator, AsyncIterator
from contextlib import aclosing, asynccontextmanager
 
 
async def ticks(n: int) -> AsyncGenerator[int]:
    for i in range(n):
        await asyncio.sleep(0.01)
        yield i
 
 
@asynccontextmanager
async def connection(url: str) -> AsyncIterator[str]:
    print("open", url)
    try:
        yield f"conn:{url}"
    finally:
        print("close", url)
 
 
async def main() -> None:
    async with connection("db://local") as conn:
        async with aclosing(ticks(5)) as it:
            evens = [i async for i in it if i % 2 == 0]
        print(conn, evens)  # conn:db://local [0, 2, 4]
 
asyncio.run(main())

Blocking code in async

Anything that doesn't await blocks every task: time.sleep, requests, file reads, heavy math. Push it off the loop.

APIRuns fn inUse
await asyncio.to_thread(fn, *args)the loop's default thread poolblocking I/O, sync SDKs
await loop.run_in_executor(pool, fn, *args)a pool you pass (None = default)CPU work in a ProcessPoolExecutor
asyncio.run_coroutine_threadsafe(coro, loop)the loop, from another threadreturns a concurrent.futures.Future
loop.call_soon_threadsafe(cb, *args)the loop, from another threadwake the loop, set an Event
asyncio.wrap_future(fut)await a concurrent.futures.Future
import asyncio
import time
 
 
def blocking_io(path: str) -> int:
    time.sleep(0.2)  # stands in for a sync library call
    return len(path)
 
 
async def main() -> None:
    sizes = await asyncio.gather(
        asyncio.to_thread(blocking_io, "a.txt"),
        asyncio.to_thread(blocking_io, "bb.txt"),
    )
    print(sizes)  # [5, 6] after ~0.2 s, not 0.4 s
 
asyncio.run(main())

to_thread copies the current contextvars context into the thread. Threads can't be canceled: the await is canceled, the function runs to the end.

concurrent.futures

One API over threads, processes and (3.14) interpreters. Use it from sync code, or from async via run_in_executor.

APINotes
ThreadPoolExecutor(max_workers)default min(32, cpu + 4) threads
ProcessPoolExecutor(max_workers)default os.process_cpu_count(); args and results must pickle
InterpreterPoolExecutor(max_workers)3.14: one subinterpreter per worker
pool.submit(fn, *args)returns a Future
pool.map(fn, it, timeout=, chunksize=, buffersize=)results in input order; buffersize (3.14) bounds work in flight
as_completed(futs, timeout)futures in finish order
wait(futs, return_when=FIRST_EXCEPTION)(done, not_done)
fut.result(timeout)value, re-raises the worker's exception
pool.shutdown(wait=True, cancel_futures=False)with calls it for you
ProcessPoolExecutor(max_tasks_per_child=n)recycle workers (leaky C libs)
terminate_workers(), kill_workers()3.14: stop process workers now
from concurrent.futures import (
    ThreadPoolExecutor,
    as_completed,
)
from urllib.request import urlopen
 
 
def status(url: str) -> tuple[str, int]:
    with urlopen(url, timeout=5) as res:
        return url, res.status
 
 
urls = ["https://example.com", "https://python.org"]
with ThreadPoolExecutor(max_workers=8) as pool:
    futs = [pool.submit(status, u) for u in urls]
    for fut in as_completed(futs):
        try:
            print(fut.result())
        except OSError as err:
            print("failed:", err)

Threads

threading APIUse
Thread(target=fn, args=(...), daemon=True).start(), .join(timeout); daemon threads die with the process
Lock(), RLock()with lock:; RLock is re-entrant for the same thread
Event()set(), wait(timeout), is_set(): stop flags
Condition(), Semaphore(n), Barrier(n)same ideas as the asyncio versions, but blocking
local()per-thread attributes (prefer contextvars)
Timer(sec, fn)run fn once after a delay; .cancel()
queue.Queue(maxsize)thread-safe hand-off; put, get(timeout), task_done, join, shutdown (3.13)
threading.current_thread().namelogging, debugging
import queue
import threading
 
jobs: queue.Queue[int] = queue.Queue()
results: list[int] = []
lock = threading.Lock()
 
 
def worker(stop: threading.Event) -> None:
    while not stop.is_set():
        try:
            n = jobs.get(timeout=0.1)
        except queue.Empty:
            continue
        with lock:
            results.append(n * n)
        jobs.task_done()
 
 
stop = threading.Event()
threads = [
    threading.Thread(target=worker, args=(stop,))
    for _ in range(4)
]
for t in threads:
    t.start()
for n in range(10):
    jobs.put(n)
jobs.join()  # every job done
stop.set()
for t in threads:
    t.join()
print(sorted(results))

The GIL makes single bytecodes atomic, not your read-modify-write sequences: x += 1 from many threads still needs a lock, and on free-threaded builds even more so.

Processes

multiprocessing APIUse
Process(target=fn, args=(...)).start(), .join(), .exitcode
Pool(n)map, imap_unordered, starmap, apply_async; prefer ProcessPoolExecutor
Queue(), Pipe()pickled messages between processes
Manager()proxied dict/list shared across processes (slow)
shared_memory.SharedMemoryraw shared bytes (pair with NumPy)
Value, Arrayshared C scalars/arrays with a lock
get_context("spawn")pick a start method per pool
Start methodDefault onNotes
spawnmacOS, Windowsfresh interpreter; re-imports your module
forkserverLinux and other POSIX (3.14+)forks from a clean server process
forkLinux before 3.14copies the parent; unsafe with threads
from concurrent.futures import ProcessPoolExecutor
 
 
def cpu_heavy(n: int) -> int:
    return sum(i * i for i in range(n))
 
 
if __name__ == "__main__":  # required for spawn/forkserver
    with ProcessPoolExecutor() as pool:
        totals = list(pool.map(cpu_heavy, [10**6] * 8))
    print(totals[0])

Worker functions must be importable top-level functions (no lambdas, no closures), and every argument and result is pickled.

GIL, free threading & subinterpreters

The GIL lets one thread run Python bytecode at a time. Threads still overlap on I/O and in C code that releases it. Two ways around it:

Free-threaded buildSubinterpreters
Status in 3.14officially supported (PEP 779), still optionalnew stdlib module concurrent.interpreters (PEP 734)
Install / runuv python install 3.14t, uv run -p 3.14t app.py, python3.14tnormal python3.14
Parallelismthreads run truly in parallelone GIL per interpreter
Sharingordinary objects (you add the locks)isolated; pass shareable objects through queues
Costsa few % slower single-threaded, more memory; C extensions need cp314t wheelsstart-up and memory per interpreter; many extensions not yet supported
Checksys._is_gil_enabled()
Re-enable GILPYTHON_GIL=1 or -X gil=1
import sys
import sysconfig
 
free_build = sysconfig.get_config_var("Py_GIL_DISABLED")
print("free-threaded build:", bool(free_build))
print("GIL on now:", sys._is_gil_enabled())

Importing an extension module that isn't marked free-threading-safe turns the GIL back on (with a warning). Check sys._is_gil_enabled() after your imports.

interp.py
from concurrent import interpreters
from concurrent.futures import InterpreterPoolExecutor
 
 
def fib(n: int) -> int:
    return n if n < 2 else fib(n - 1) + fib(n - 2)
 
 
if __name__ == "__main__":
    interp = interpreters.create()
    print(interp.call(fib, 20))  # 6765, in isolation
    interp.close()
 
    with InterpreterPoolExecutor(max_workers=4) as pool:
        print(list(pool.map(fib, [24, 25, 26])))
concurrent.interpretersDoes
create()new isolated interpreter
interp.exec(code)run source text in its __main__
interp.call(fn, *args)call and return the result
interp.call_in_thread(fn, *args)same, in a new thread; returns the Thread
create_queue(maxsize)cross-interpreter queue; put/get
is_shareable(obj)None, bool, int, float, str, bytes, tuples of those, queues, memoryview

Pitfalls

SymptomCauseFix
RuntimeWarning: coroutine ... was never awaitedcalled f() without awaitawait f() or create_task(f())
Everything freezes for secondsblocking call in a coroutine (time.sleep, requests)await asyncio.sleep, an async client, or to_thread
Background task silently vanishesthe loop keeps only a weak reference to taskskeep a reference (set + add_done_callback(discard)) or use TaskGroup
Task exception was never retrievednobody awaited a failed taskawait it, or use TaskGroup
Shutdown hangs or tasks won't stopexcept BaseException / bare except swallowed CancelledErrorre-raise it
gather failure leaves work runningplain gather doesn't cancel siblingsTaskGroup
RuntimeError: asyncio.run() cannot be called from a running event loopJupyter, or nested runawait main() directly
RuntimeError: ... attached to a different loopa Lock/Queue reused across two asyncio.run callscreate primitives inside main()
Pickling error in a process poollambda, closure or local functionmove it to module top level
Infinite process spawningno if __name__ == "__main__": guardadd the guard
import asyncio
from collections.abc import Coroutine
from typing import Any
 
background: set[asyncio.Task[Any]] = set()
 
 
def fire_and_forget(coro: Coroutine[Any, Any, Any]) -> None:
    task = asyncio.create_task(coro)
    background.add(task)  # strong reference
    task.add_done_callback(background.discard)

Debug a live program with python -m asyncio ps PID or pstree PID (3.14): every task, its name and what it awaits. In code, asyncio.print_call_graph() shows the current task's awaiters.

Compared with JS promises

JavaScriptPython
calling an async fn starts itcalling returns an idle coroutine; await or create_task starts it
Promiseasyncio.Task / Future
event loop is implicitasyncio.run(main()) starts one
Promise.all([...])TaskGroup (cancels siblings) or gather
Promise.allSettledgather(..., return_exceptions=True)
Promise.race / anyasyncio.wait(..., FIRST_COMPLETED) / as_completed
AbortController + signaltask.cancel(); no token to thread through
AbortSignal.timeout(ms)async with asyncio.timeout(sec):
for await (const x of it)async for x in it:
setTimeout(fn, ms)loop.call_later(sec, fn); await asyncio.sleep(sec)
queueMicrotask(fn)loop.call_soon(fn)
unhandled rejection event"Task exception was never retrieved" log; loop.set_exception_handler
Worker threadsthreads, processes, subinterpreters

Recipes

Bounded concurrency with a semaphore

When you have many jobs but the server (or your socket limit) allows only a few at once.

import asyncio
from collections.abc import Awaitable, Callable, Iterable
 
 
async def bounded[T, R](
    items: Iterable[T],
    fn: Callable[[T], Awaitable[R]],
    limit: int = 10,
) -> list[R]:
    sem = asyncio.Semaphore(limit)
 
    async def one(item: T) -> R:
        async with sem:
            return await fn(item)
 
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(one(i)) for i in items]
    return [t.result() for t in tasks]

Retry with exponential backoff

When a flaky call is idempotent and worth a few more tries with full jitter.

import asyncio
import random
from collections.abc import Awaitable, Callable
 
 
async def retry[T](
    fn: Callable[[], Awaitable[T]],
    *,
    attempts: int = 4,
    base: float = 0.2,
    cap: float = 5.0,
    retry_on: tuple[type[Exception], ...] = (OSError,),
) -> T:
    for attempt in range(attempts):
        try:
            return await fn()
        except retry_on:
            if attempt == attempts - 1:
                raise
            ceiling = min(cap, base * 2**attempt)
            await asyncio.sleep(random.uniform(0, ceiling))
    raise AssertionError("unreachable")

Timeout per call, keep the rest

When each call gets its own deadline and a slow one shouldn't sink the batch.

import asyncio
from collections.abc import Awaitable
 
 
async def with_timeout[T](
    aw: Awaitable[T], sec: float
) -> T | None:
    try:
        async with asyncio.timeout(sec):
            return await aw
    except TimeoutError:
        return None
 
 
async def main() -> None:
    delays = [0.05, 2.0, 0.1]
    results = await asyncio.gather(
        *(with_timeout(asyncio.sleep(d, d), 0.5)
          for d in delays)
    )
    print(results)  # [0.05, None, 0.1]
 
asyncio.run(main())

Producer / consumer with a queue

When work arrives as a stream and a fixed set of workers should drain it with back-pressure.

import asyncio
 
 
async def producer(q: asyncio.Queue[int]) -> None:
    for n in range(20):
        await q.put(n)  # waits while the queue is full
    q.shutdown()  # 3.13+: consumers stop once drained
 
 
async def consumer(name: str, q: asyncio.Queue[int]) -> None:
    while True:
        try:
            n = await q.get()
        except asyncio.QueueShutDown:
            return
        await asyncio.sleep(0.01)  # process n
        print(name, n)
 
 
async def main() -> None:
    q: asyncio.Queue[int] = asyncio.Queue(maxsize=5)
    async with asyncio.TaskGroup() as tg:
        tg.create_task(producer(q))
        for i in range(3):
            tg.create_task(consumer(f"w{i}", q))
 
asyncio.run(main())

Graceful shutdown on SIGINT / SIGTERM

When a long-running service should finish in-flight work and clean up on Ctrl+C or docker stop (Unix).

import asyncio
import signal
 
 
async def serve(stop: asyncio.Event) -> None:
    while not stop.is_set():
        await asyncio.sleep(0.5)  # handle one unit of work
    print("draining...")
 
 
async def main() -> None:
    stop = asyncio.Event()
    loop = asyncio.get_running_loop()
    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, stop.set)
    try:
        await serve(stop)
    finally:
        print("closing connections")
 
if __name__ == "__main__":
    asyncio.run(main())

Without handlers, asyncio.run turns the first Ctrl+C into a cancel of main() and raises KeyboardInterrupt once it has unwound.

CPU work in a process pool from async code

When a coroutine must crunch numbers without freezing the loop.

import asyncio
from concurrent.futures import ProcessPoolExecutor
 
 
def crunch(n: int) -> int:  # top level: must pickle
    return sum(i * i for i in range(n))
 
 
async def main() -> None:
    loop = asyncio.get_running_loop()
    with ProcessPoolExecutor() as pool:
        results = await asyncio.gather(*(
            loop.run_in_executor(pool, crunch, n)
            for n in (10**6, 2 * 10**6, 3 * 10**6)
        ))
    print(results)
 
if __name__ == "__main__":
    asyncio.run(main())

Async HTTP fan-out with httpx

When you need many HTTP calls at once over one pooled client (uv add httpx).

import asyncio
 
import httpx
 
 
async def fetch_all(urls: list[str]) -> dict[str, int]:
    limits = httpx.Limits(max_connections=10)
    timeout = httpx.Timeout(10.0)
    async with httpx.AsyncClient(
        limits=limits, timeout=timeout
    ) as client:
        async with asyncio.TaskGroup() as tg:
            tasks = {
                u: tg.create_task(client.get(u))
                for u in urls
            }
    return {u: t.result().status_code
            for u, t in tasks.items()}
 
 
urls = ["https://example.com", "https://www.python.org"]
print(asyncio.run(fetch_all(urls)))

One failed request cancels the rest; wrap client.get in retry or with_timeout above to tolerate failures.

References