concurrent.futures, Executors & Future Objects

The concurrent.futures module provides a high-level, unified asynchronous execution framework built around the Executor interface. By decoupling task submission from task execution models, developers can switch between thread pools (ThreadPoolExecutor) for I/O-bound work and process pools (ProcessPoolExecutor) for CPU-bound work without modifying task code.

This chapter details the Executor interface, ThreadPoolExecutor vs ProcessPoolExecutor, Future object state machines (PENDING $\rightarrow$ RUNNING $\rightarrow$ FINISHED), exception handling, and as_completed() iterator processing.


1. The Executor Abstraction Interface

The concurrent.futures module defines two concrete Executor implementations sharing an identical API:

The Executor Abstraction Architecture:

[ Task Submission: executor.submit(fn, *args) ]
                       |
                       v Returns Future Object
              [ Future Object (State: PENDING) ]
                       |
     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
     β”‚                                   β”‚
     v                                   v
[ ThreadPoolExecutor ]          [ ProcessPoolExecutor ]
(Concurrently handles I/O)       (Parallelizes CPU on cores)
  • ThreadPoolExecutor(max_workers=N): Ideal for network requests, file I/O, and database queries.
  • ProcessPoolExecutor(max_workers=N): Ideal for CPU-intensive data transformations, image processing, and cryptography.

2. Future Object State Machine & API

A Future represents an eventual result of an asynchronous computation.

Future Object State Machine:

[ PENDING ] ───> (Worker starts) ───> [ RUNNING ] ───> (Task finishes) ───> [ FINISHED ]
     |                                                                           |
     └─── (cancel() invoked) ───────> [ CANCELLED ]                             β”œβ”€β”€ (Normal)  --> result()
                                                                                └── (Error)   --> exception()

Key Future Methods:

  • future.result(timeout=None): Blocks until the computation completes and returns the result. If the task raised an exception, result() re-raises the exception at the caller site!
  • future.exception(timeout=None): Returns the exception raised by the task (or None if successful).
  • future.done(): Returns True if the task is finished or cancelled.
  • future.add_done_callback(fn): Registers a callback function executed automatically when the task finishes.

3. Asynchronous Iterator Processing (as_completed & map)

1. as_completed(futures): Yields Results as They Finish

Processes futures in order of completion, regardless of submission order:

from concurrent.futures import ThreadPoolExecutor, as_completed
import requests

URLS = ["https://httpbin.org/delay/2", "https://httpbin.org/delay/1", "https://httpbin.org/delay/3"]

def fetch(url: str) -> int:
    return requests.get(url, timeout=5).status_code

with ThreadPoolExecutor(max_workers=3) as executor:
    # Submit tasks; returns dict mapping Future -> URL
    future_to_url = {executor.submit(fetch, url): url for url in URLS}

    # Yields futures as soon as they COMPLETE (fastest response yields first!)
    for future in as_completed(future_to_url):
        url = future_to_url[future]
        try:
            status = future.result()
            print(f"[{url}] Status: {status}")
        except Exception as exc:
            print(f"[{url}] Generated exception: {exc}")

2. executor.map(fn, *iterables): Preserves Submission Order

Executes fn over iterables concurrently, yielding results strictly in original submission order.


4. Production Context Manager Behavior (with)

Always use with Executor() as executor: context manager blocks. Exiting the with block automatically calls executor.shutdown(wait=True), joining worker threads/processes and guaranteeing clean resource teardown.

Display Options
Appearance
Text Size
100%