Parallelism with concurrent.futures

Reviewed & published by Brayan K

Master Python's concurrent.futures module to build high-performance parallel systems using ThreadPoolExecutor and ProcessPoolExecutor.

Part of the free Python course at LearnCodingFast — hands-on lessons with examples you run in your browser, plus practice exercises and a quick quiz.

What You'll Learn in This Lesson

1. What concurrent.futures Actually Does

It provides two executor types:

ExecutorWorkers AreBest ForExamples
ThreadPoolExecutorThreadsI/O-bound tasks (waiting)API calls, file downloads, DB queries
ProcessPoolExecutorProcessesCPU-bound tasks (thinking)Image processing, ML prep, heavy math

⚙️ 2. The Basic Pattern (Threads)

from concurrent.futures import ThreadPoolExecutor
import time

def download(url):
    # Simulate a network request that takes time
    time.sleep(1)
    return f"Downloaded {url}"

# Create a pool with 4 worker threads
# Think of it as hiring 4 assistants
executor = ThreadPoolExecutor(max_workers=4)

# Submit a task - returns immediately with a Future
# The task runs in the background
future = executor.submit(download, "https://example.com")
print("Task submitted - we can do other work here!")

# When we need the result, call .result()
# This BLOCKS until the task is complete
result = future.result()
print(result)

# Always clean up - release the workers
executor.shutdown()
MethodWhat It DoesReturns
submit(fn, *args)Schedules function to run in backgroundFuture object
future.result()Waits for and returns the resultFunction's return value
shutdown()Releases worker threadsNone

⚡ 3. Running Many Tasks at Once

from concurrent.futures import ThreadPoolExecutor

def fetch(url):
    # simulate work - in real code this would be an HTTP request
    return f"Content from {url}"

urls = ["url1", "url2", "url3", "url4"]

# The 'with' statement ensures automatic cleanup
# No need to call executor.shutdown() manually!
with ThreadPoolExecutor(max_workers=4) as executor:
    # map() is like the built-in map(), but parallel!
    # It automatically:
    # 1. Submits all tasks
    # 2. Collects results in order
    # 3. Returns an iterator of results
    results = list(executor.map(fetch, urls))

print(results)
# Output: ['Content from url1', 'Content from url2', ...]
ApproachWhen to UseReturns
submit()Different functions, custom handlingIndividual Futures
map()Same function, many inputsIterator of results (in order)

4. CPU Parallelism with ProcessPoolExecutor

from concurrent.futures import ProcessPoolExecutor

def heavy_compute(n):
    # CPU-intensive work: sum of squares
    # This keeps the CPU busy (no waiting)
    return sum(i * i for i in range(n))

numbers = [10_000_000, 10_000_000, 10_000_000, 10_000_000]

# ProcessPoolExecutor uses separate processes, not threads
# Each process gets its own Python interpreter and GIL
with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(heavy_compute, numbers))

print(results)
# On a 4-core CPU, this can be ~4x faster than sequential!
ExecutorMemoryGILCPU Usage
ThreadPoolExecutorSharedOne GIL (blocks CPU work)1 core effective
ProcessPoolExecutorSeparateEach has own GILAll cores!

🧩 5. Futures — Understanding the Object

A Future represents a pending operation. It can be in different states:

StateMeaningCheck With
RunningTask is currently being executedfuture.running()
DoneTask completed (success or error)future.done()
CancelledTask was cancelled before runningfuture.cancelled()
from concurrent.futures import ThreadPoolExecutor
import time

def task(n):
    # Simulate work that takes 'n' seconds
    time.sleep(n)
    return n * 2

executor = ThreadPoolExecutor()

# Submit a task that takes 2 seconds
future = executor.submit(task, 2)

# Check state immediately (task is running)
print(f"Done? {future.done()}")      # False
print(f"Running? {future.running()}")  # True

# Wait for result (blocks until complete)
result = future.result()
print(f"Result: {result}")  # 4

# Now it's done
print(f"Done? {future.done()}")  # True

🌀 6. Handling Exceptions in Parallel Tasks

If a function raises an exception inside a worker, .result() will re-raise it in the main thread:

from concurrent.futures import ThreadPoolExecutor

def failing_task():
    # This exception happens in a worker thread
    raise ValueError("Something went wrong")

executor = ThreadPoolExecutor()
future = executor.submit(failing_task)

# The exception is "stored" in the Future
# When we call .result(), it's re-raised here
try:
    future.result()
except ValueError as e:
    print(f"Caught: {e}")
# Output: Caught: Something went wrong
MethodOn SuccessOn Error
future.result()Returns the valueRaises the exception
future.exception()Returns NoneReturns the exception object

7. Real-World Example — Parallel Web Requests

from concurrent.futures import ThreadPoolExecutor
import requests

def fetch_url(url):
    response = requests.get(url)
    return len(response.content)

urls = [
    "https://api.github.com/users/github",
    "https://api.github.com/users/microsoft",
    "https://api.github.com/users/google"
]

with ThreadPoolExecutor(max_workers=10) as executor:
    sizes = list(executor.map(fetch_url, urls))

print(f"Total bytes: {sum(sizes)}")
ApproachTime for 100 URLs (0.3s each)Speedup
Sequential (one by one)30 seconds1x (baseline)
ThreadPool (10 workers)~3 seconds~10x faster!

⚡ 8. Real-World Example — CPU Parallel Data Processing

from concurrent.futures import ProcessPoolExecutor
import hashlib

def hash_password(password):
    return hashlib.pbkdf2_hmac('sha256', 
                                 password.encode(), 
                                 b'salt', 
                                 100000).hex()

passwords = ["pass123", "secret", "mypass"] * 1000

with ProcessPoolExecutor() as executor:
    hashed = list(executor.map(hash_password, passwords))

print(f"Hashed {len(hashed)} passwords")

🔄 9. Mixing Concurrency & Parallelism

For the best performance, systems combine:

Pipeline StageBest ToolWhy
Download dataThreadPoolExecutorI/O-bound: waiting for servers
Parse/transform dataProcessPoolExecutorCPU-bound: heavy computation
Coordinate/scheduleasyncioLightweight: manage task flow

This is how modern Python backends (FastAPI, aiohttp) achieve massive throughput.

📦 10. Choosing the Right Executor

ScenarioBest Choice
Many API callsThreadPoolExecutor
Downloading filesThreadPoolExecutor
Reading thousands of filesThreadPoolExecutor
Image processingProcessPoolExecutor
ML preprocessingProcessPoolExecutor
Large math loopsProcessPoolExecutor
ETL pipelinesBoth (mixed)

If your task is waiting, use threads.

If your task is thinking, use processes.

🧩 11-20. Behind the Scenes & Advanced Patterns

11. How Executors Work

Task queue, worker threads/processes, IPC mechanisms

Pickling costs, lambda limitations, using top-level functions

13. Batching Large Jobs

Grouping tasks for efficiency, reducing pickling overhead

Optimizing executor.map() with chunksize parameter

15. Managing Shared State

multiprocessing.Manager, Queue, shared_memory

Preventing long tasks from blocking the pool

Never call .result() inside worker tasks

Running futures without blocking main thread

Pinning threads to CPU cores for consistent latency

20. Hybrid Pipeline Architecture

AsyncIO → ThreadPool → ProcessPool → AsyncIO pattern

🧪 21. Full Real-World Example — Data ETL Pipeline

import asyncio
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

io_pool = ThreadPoolExecutor(20)
cpu_pool = ProcessPoolExecutor(8)

def read_file(path):
    with open(path) as f:
        return f.read()

def transform(text):
    return text.upper()

async def main(files):
    loop = asyncio.get_event_loop()

    # Stage 1: Read files (I/O bound)
    contents = await asyncio.gather(*[
        loop.run_in_executor(io_pool, read_file, f)
        for f in files
    ])

    # Stage 2: Transform (CPU bound)
    results = await asyncio.gather(*[
        loop.run_in_executor(cpu_pool, transform, c)
        for c in contents
    ])

    print("Pipeline complete.")

asyncio.run(main(["a.txt", "b.txt", "c.txt"]))

This is TRUE professional pipeline architecture.

🧬 22. Zero-Copy Shared Memory (Python 3.8+)

Normal multiprocessing copies all data via pickle.

Approach100MB Array to 4 WorkersMemory Used
Normal (pickle)~2 seconds400MB (4 copies)
Shared memory~0.001 seconds100MB (1 shared)
from multiprocessing import shared_memory
import numpy as np

# Create shared array
data = np.arange(10_000_000, dtype=np.int32)
shm = shared_memory.SharedMemory(create=True, size=data.nbytes)

# Write to shared block
shared_arr = np.ndarray(data.shape, dtype=data.dtype, buffer=shm.buf)
shared_arr[:] = data[:]

# Now multiple workers can read the same data:
# ✔ Zero copy
# ✔ Zero serialization
# ✔ Lightning-fast ML preprocessing

🚀 23-32. Expert-Level Parallel Patterns

Professional parallel computing patterns:

23. ML Tensor Preprocessing

Parallel normalization using shared memory

24. Fan-Out / Fan-In Architecture

Split work, collect results — backbone of scalable systems

as_completed() for responsive systems

Timeouts, retries, error-tolerant submission

Separate pools for IO, CPU, and orchestration

Queue-based throttling to prevent overload

Parallel map, sequential reduce — Hadoop/Spark ancestor

30. Mini Distributed Engine Design

Task graph, scheduler, future tracking — Ray/Dask concepts

31. Building Distributed Executors

ProcessPool + sockets for multi-machine tasks

32. Master Hybrid Pipeline

AsyncIO orchestration + ThreadPool I/O + ProcessPool CPU

🚀 32. Master Hybrid Pipeline (Complete Example)

The ultimate architecture combining AsyncIO + Threads + Processes:

import asyncio
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

async def master_pipeline(files):
    loop = asyncio.get_event_loop()
    io_pool = ThreadPoolExecutor(32)
    cpu_pool = ProcessPoolExecutor(8)

    # Stage 1 — Async orchestrates
    file_contents = await asyncio.gather(*[
        loop.run_in_executor(io_pool, open_file, f)
        for f in files
    ])

    # Stage 2 — CPU processing
    processed = await asyncio.gather(*[
        loop.run_in_executor(cpu_pool, transform, data)
        for data in file_contents
    ])

    # Stage 3 — Async merge & upload
    await upload_results(processed)

    print("Pipeline complete.")
# 🎯 YOUR TURN — replace each ___ using the hint beside it.

from concurrent.futures import ThreadPoolExecutor

def word_count(text):
    return len(text.split())

pages = ["one two three", "four five", "six seven eight nine", "ten"]

# 1) "with" shuts the pool down and waits for every task, even on an error.
with ThreadPoolExecutor(max_workers=3) as pool:
    # 2) map keeps the INPUT order, whatever order the threads finish in.
    #    That is why this exercise can promise an exact answer at all.
    counts = list(pool.___(word_count, pages))    # 👉 replace ___ with map

print("Counts:", counts)
print("Total words:", sum(counts))

with ThreadPoolExecutor(max_workers=2) as pool:
    # 3) submit() hands back a Future straight away and does not block.
    futures = [pool.___(word_count, p) for p in pages]   # 👉 replace ___ with submit
    # Sorted, because completion order is a race and must never be claimed.
    done = sorted(f.___() for f in futures)              # 👉 replace ___ with result

print("Sorted results:", done)

def risky(n):
    if n == 0:
        raise ValueError("cannot divide by zero")
    return 10 // n

with ThreadPoolExecutor() as pool:
    future = pool.submit(risky, 0)
    # 4) The exception is stored in the Future and re-raised here, not in
    #    the worker — so this is where you catch it.
    ___:                                          # 👉 replace ___ with try
        future.result()
    except ValueError as e:
        print("Raised in the worker, caught here:", e)

# ✅ Expected output:
# Counts: [3, 2, 4, 1]
# Total words: 10
# Sorted results: [1, 2, 3, 4]
# Raised in the worker, caught here: cannot divide by zero

This is the same architecture used by:

🎉 Conclusion

You now understand ULTRA-ADVANCED concurrency and parallelism:

✔ Futures and parallel execution

✔ Exception handling in parallel tasks

✔ Real-world I/O & CPU examples

✔ Batching and chunksize tuning

This module is the foundation of high-performance Python systems — from ML pipelines to scalable backend services.

📋 Quick Reference — Parallelism

SyntaxWhat it does
ThreadPoolExecutor(max_workers=4)Pool for I/O-bound tasks
ProcessPoolExecutor(max_workers=4)Pool for CPU-bound tasks
executor.submit(fn, arg)Submit one task, returns Future
executor.map(fn, items)Map function over iterable
as_completed(futures)Iterate futures as they finish

🎉 Great work! You've completed this lesson.

You can now use ThreadPoolExecutor and ProcessPoolExecutor to parallelize real work efficiently and safely.

Practice quiz

Which concurrent.futures executor is best for I/O-bound work like API calls and downloads?

  • ProcessPoolExecutor
  • asyncio.Executor
  • ThreadPoolExecutor
  • Both are equally bad for I/O

Answer: ThreadPoolExecutor. ThreadPoolExecutor suits I/O-bound tasks where workers spend most time waiting.

Which executor is best for CPU-bound work like image processing or heavy math?

  • ProcessPoolExecutor
  • ThreadPoolExecutor
  • A single thread
  • asyncio alone

Answer: ProcessPoolExecutor. ProcessPoolExecutor uses separate processes, each with its own GIL, to use multiple cores for CPU work.

What does executor.submit(fn, arg) return?

  • The function's result immediately
  • None
  • A list of results
  • A Future object representing the pending work

Answer: A Future object representing the pending work. submit() schedules the call and immediately returns a Future you can query later.

What does calling future.result() do?

  • Cancels the task
  • Blocks until the task finishes, then returns its value (or re-raises its exception)
  • Returns instantly even if not done
  • Starts the task

Answer: Blocks until the task finishes, then returns its value (or re-raises its exception). .result() waits for completion and returns the value, or re-raises any exception from the worker.

What ordering does executor.map(fn, items) guarantee for its results?

  • Results in the same order as the input items
  • Results in completion order (fastest first)
  • Random order
  • Reverse order

Answer: Results in the same order as the input items. map() preserves input order regardless of which task finishes first.

If a function raises an exception inside a worker, when is it surfaced?

  • Immediately in the worker, crashing the program
  • It is silently ignored
  • When you call future.result(), which re-raises it in the calling thread
  • Only when the executor shuts down

Answer: When you call future.result(), which re-raises it in the calling thread. The exception is stored in the Future and re-raised when you call .result().

What is a Future?

  • A scheduled date for the task
  • An object representing work that has started but may not be finished yet
  • A type of thread
  • The result value itself

Answer: An object representing work that has started but may not be finished yet. A Future is a placeholder for a pending result that you can check, wait on, or cancel.

Why can't ThreadPoolExecutor speed up pure-Python CPU-bound loops much?

  • Threads are slower than processes always
  • Threads can't share memory
  • ThreadPoolExecutor has a 1-worker limit
  • The GIL lets only one thread run Python bytecode at a time

Answer: The GIL lets only one thread run Python bytecode at a time. Because of the GIL, threads can't execute Python bytecode truly in parallel, so CPU-bound code stays on one core.

Using a ProcessPoolExecutor as 'with ... as executor:' provides what benefit?

  • It disables the GIL
  • Automatic cleanup/shutdown of workers when the block exits
  • It makes tasks run sequentially
  • It removes the need for results

Answer: Automatic cleanup/shutdown of workers when the block exits. The context manager automatically shuts down the executor, so you don't call shutdown() manually.

Which method returns the stored exception (or None) without raising it?

  • future.result()
  • future.running()
  • future.exception()
  • future.cancel()

Answer: future.exception(). future.exception() returns the exception object if the task failed, or None on success.

Continue this course