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
- • How concurrent.futures abstracts threads and processes into a unified API
- • When to use ThreadPoolExecutor vs ProcessPoolExecutor
- • How Futures represent pending work — and how to collect results
- • Running parallel tasks with map() and submit()
- • Handling errors, cancellations, and timeouts in parallel jobs
- • Real patterns: parallel image processing, batch API calls, data pipelines
1. What concurrent.futures Actually Does
It provides two executor types:
| Executor | Workers Are | Best For | Examples |
|---|---|---|---|
| ThreadPoolExecutor | Threads | I/O-bound tasks (waiting) | API calls, file downloads, DB queries |
| ProcessPoolExecutor | Processes | CPU-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()| Method | What It Does | Returns |
|---|---|---|
| submit(fn, *args) | Schedules function to run in background | Future object |
| future.result() | Waits for and returns the result | Function's return value |
| shutdown() | Releases worker threads | None |
⚡ 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', ...]| Approach | When to Use | Returns |
|---|---|---|
| submit() | Different functions, custom handling | Individual Futures |
| map() | Same function, many inputs | Iterator 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!| Executor | Memory | GIL | CPU Usage |
|---|---|---|---|
| ThreadPoolExecutor | Shared | One GIL (blocks CPU work) | 1 core effective |
| ProcessPoolExecutor | Separate | Each has own GIL | All cores! |
🧩 5. Futures — Understanding the Object
A Future represents a pending operation. It can be in different states:
| State | Meaning | Check With |
|---|---|---|
| Running | Task is currently being executed | future.running() |
| Done | Task completed (success or error) | future.done() |
| Cancelled | Task was cancelled before running | future.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| Method | On Success | On Error |
|---|---|---|
| future.result() | Returns the value | Raises the exception |
| future.exception() | Returns None | Returns 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)}")| Approach | Time for 100 URLs (0.3s each) | Speedup |
|---|---|---|
| Sequential (one by one) | 30 seconds | 1x (baseline) |
| ThreadPool (10 workers) | ~3 seconds | ~10x faster! |
- API pipelines
- ETL ingestion
⚡ 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")- AI data prep
- big dataset cleaning
- batch computation services
🔄 9. Mixing Concurrency & Parallelism
For the best performance, systems combine:
- I/O parallelism (threads)
- CPU parallelism (processes)
- Task scheduling (asyncio)
| Pipeline Stage | Best Tool | Why |
|---|---|---|
| Download data | ThreadPoolExecutor | I/O-bound: waiting for servers |
| Parse/transform data | ProcessPoolExecutor | CPU-bound: heavy computation |
| Coordinate/schedule | asyncio | Lightweight: manage task flow |
- asyncio → orchestrates
- ThreadPoolExecutor → handles blocking file/network
- ProcessPoolExecutor → handles heavy CPU tasks
This is how modern Python backends (FastAPI, aiohttp) achieve massive throughput.
📦 10. Choosing the Right Executor
| Scenario | Best Choice |
|---|---|
| Many API calls | ThreadPoolExecutor |
| Downloading files | ThreadPoolExecutor |
| Reading thousands of files | ThreadPoolExecutor |
| Image processing | ProcessPoolExecutor |
| ML preprocessing | ProcessPoolExecutor |
| Large math loops | ProcessPoolExecutor |
| ETL pipelines | Both (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.
| Approach | 100MB Array to 4 Workers | Memory Used |
|---|---|---|
| Normal (pickle) | ~2 seconds | 400MB (4 copies) |
| Shared memory | ~0.001 seconds | 100MB (1 shared) |
- ❌ Large NumPy arrays (100MB+) = slow
- ❌ ML tensors = slow
- ❌ Video frames = slow
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- PyTorch multiprocessing
- TensorFlow data pipelines
- High-performance ETL
🚀 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 zeroThis is the same architecture used by:
- ✔ Netflix data pipelines
- ✔ TikTok's ML recommendation ingest
- ✔ YouTube's batch processing
- ✔ OpenAI's internal preprocessors
🎉 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
| Syntax | What 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
- Previous: Concurrency in Python: Threads vs Processes
- Next: Profiling & Optimising Python Performance — Find bottlenecks with cProfile, timeit, and memory profilers
- Quick reference: Python cheat sheet