Concurrent.futures in Python: A Comprehensive Guide
Introduction
The concurrent.futures module provides a high-level interface for asynchronously executing callable objects. It was introduced in Python 3.2 and offers a simple, consistent API for both threading and multiprocessing, making it easier to write concurrent code without dealing with the complexities of low-level threading or multiprocessing APIs.
Key Concepts
What is concurrent.futures?
concurrent.futures provides a unified interface for:
- ThreadPoolExecutor: For I/O-bound tasks using threads
- ProcessPoolExecutor: For CPU-bound tasks using processes
- Future objects: Represent the eventual result of an asynchronous operation
Core Components
- Executor Classes: Abstract base classes for executing callables
- Future Objects: Represent pending computations
- Context Managers: Automatic resource management
- Exception Handling: Robust error propagation
ThreadPoolExecutor
Basic Usage
from concurrent.futures import ThreadPoolExecutor
import time
import requests
def fetch_url(url):
"""Simulate fetching data from a URL"""
time.sleep(1) # Simulate network delay
return f"Data from {url}"
def basic_threading_example():
urls = [
"http://example.com/1",
"http://example.com/2",
"http://example.com/3",
"http://example.com/4"
]
with ThreadPoolExecutor(max_workers=4) as executor:
# Submit tasks and get futures
futures = [executor.submit(fetch_url, url) for url in urls]
# Collect results
results = [future.result() for future in futures]
for result in results:
print(result)
# Alternative using map
def map_example():
urls = ["http://example.com/1", "http://example.com/2"]
with ThreadPoolExecutor() as executor:
results = list(executor.map(fetch_url, urls))
print(results)
Advanced ThreadPoolExecutor Features
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task_with_id(task_id):
"""Task that takes time and returns an ID"""
time.sleep(2)
return f"Task {task_id} completed"
def advanced_threading_example():
with ThreadPoolExecutor(max_workers=3) as executor:
# Submit multiple tasks
futures = {
executor.submit(task_with_id, i): i
for i in range(5)
}
# Process results as they complete
for future in as_completed(futures):
task_id = futures[future]
try:
result = future.result(timeout=5)
print(f"Task {task_id}: {result}")
except Exception as e:
print(f"Task {task_id} failed: {e}")
ProcessPoolExecutor
Basic Usage
from concurrent.futures import ProcessPoolExecutor
import time
def cpu_intensive_task(n):
"""CPU-intensive task"""
result = 0
for i in range(n):
result += i ** 2
return result
def basic_multiprocessing_example():
numbers = [1000000, 2000000, 3000000, 4000000]
with ProcessPoolExecutor() as executor:
# Use map for simple parallel execution
results = list(executor.map(cpu_intensive_task, numbers))
for i, result in enumerate(results):
print(f"Task {i}: {result}")
def submit_example():
numbers = [1000000, 2000000, 3000000, 4000000]
with ProcessPoolExecutor() as executor:
# Submit individual tasks
futures = [executor.submit(cpu_intensive_task, n) for n in numbers]
# Collect results
results = [future.result() for future in futures]
print(results)
Future Objects
Understanding Futures
from concurrent.futures import ThreadPoolExecutor, Future
import time
def long_running_task(seconds):
time.sleep(seconds)
return f"Completed after {seconds} seconds"
def future_example():
with ThreadPoolExecutor() as executor:
# Submit a task and get a Future object
future = executor.submit(long_running_task, 3)
print(f"Future object: {future}")
print(f"Done: {future.done()}")
print(f"Running: {future.running()}")
print(f"Cancelled: {future.cancelled()}")
# Wait for result
result = future.result()
print(f"Result: {result}")
print(f"Done: {future.done()}")
Future State Management
from concurrent.futures import ThreadPoolExecutor, Future
import time
def future_states_example():
with ThreadPoolExecutor() as executor:
# Submit a task
future = executor.submit(time.sleep, 2)
# Check different states
print(f"Immediately after submission:")
print(f" Done: {future.done()}")
print(f" Running: {future.running()}")
print(f" Cancelled: {future.cancelled()}")
# Wait a bit
time.sleep(1)
print(f"\nAfter 1 second:")
print(f" Done: {future.done()}")
print(f" Running: {future.running()}")
# Wait for completion
future.result()
print(f"\nAfter completion:")
print(f" Done: {future.done()}")
print(f" Running: {future.running()}")
Exception Handling
Robust Error Handling
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task_with_exception(task_id):
"""Task that may raise an exception"""
time.sleep(1)
if task_id % 3 == 0:
raise ValueError(f"Task {task_id} failed")
return f"Task {task_id} succeeded"
def exception_handling_example():
with ThreadPoolExecutor() as executor:
futures = [executor.submit(task_with_exception, i) for i in range(10)]
successful_results = []
failed_tasks = []
for future in as_completed(futures):
try:
result = future.result()
successful_results.append(result)
except Exception as e:
failed_tasks.append(str(e))
print(f"Successful: {len(successful_results)}")
print(f"Failed: {len(failed_tasks)}")
print(f"Failed tasks: {failed_tasks}")
Timeout Handling
from concurrent.futures import ThreadPoolExecutor, TimeoutError
import time
def slow_task(seconds):
time.sleep(seconds)
return f"Completed after {seconds} seconds"
def timeout_example():
with ThreadPoolExecutor() as executor:
future = executor.submit(slow_task, 5)
try:
result = future.result(timeout=3)
print(result)
except TimeoutError:
print("Task timed out")
future.cancel()
print(f"Cancelled: {future.cancelled()}")
Advanced Patterns
Chaining Futures
from concurrent.futures import ThreadPoolExecutor
import time
def stage1(data):
"""First processing stage"""
time.sleep(1)
return data * 2
def stage2(data):
"""Second processing stage"""
time.sleep(1)
return data + 10
def stage3(data):
"""Third processing stage"""
time.sleep(1)
return data ** 2
def chaining_example():
with ThreadPoolExecutor() as executor:
# Submit first stage
future1 = executor.submit(stage1, 5)
# Chain second stage
future2 = executor.submit(stage2, future1.result())
# Chain third stage
future3 = executor.submit(stage3, future2.result())
final_result = future3.result()
print(f"Final result: {final_result}")
Batch Processing
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def process_batch(batch_id, items):
"""Process a batch of items"""
time.sleep(1)
return f"Batch {batch_id}: Processed {len(items)} items"
def batch_processing_example():
# Create batches of data
all_data = list(range(100))
batch_size = 20
batches = [
all_data[i:i + batch_size]
for i in range(0, len(all_data), batch_size)
]
with ThreadPoolExecutor(max_workers=3) as executor:
# Submit batch processing tasks
futures = {
executor.submit(process_batch, i, batch): i
for i, batch in enumerate(batches)
}
# Process results as they complete
for future in as_completed(futures):
batch_id = futures[future]
result = future.result()
print(result)
Callback Functions
from concurrent.futures import ThreadPoolExecutor, Future
import time
def callback_example():
def task_completed(future):
"""Callback function called when task completes"""
try:
result = future.result()
print(f"Task completed successfully: {result}")
except Exception as e:
print(f"Task failed: {e}")
def long_task(seconds):
time.sleep(seconds)
return f"Task completed after {seconds} seconds"
with ThreadPoolExecutor() as executor:
future = executor.submit(long_task, 2)
future.add_done_callback(task_completed)
# Wait for completion
future.result()
Performance Comparison
ThreadPoolExecutor vs ProcessPoolExecutor
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import time
def cpu_bound_task(n):
"""CPU-intensive task"""
result = 0
for i in range(n):
result += i ** 2
return result
def io_bound_task(n):
"""I/O-bound task"""
time.sleep(0.1)
return n
def performance_comparison():
# CPU-bound comparison
print("CPU-bound task comparison:")
# ThreadPoolExecutor
start_time = time.time()
with ThreadPoolExecutor() as executor:
results = list(executor.map(cpu_bound_task, [100000] * 4))
thread_time = time.time() - start_time
print(f"ThreadPoolExecutor: {thread_time:.2f} seconds")
# ProcessPoolExecutor
start_time = time.time()
with ProcessPoolExecutor() as executor:
results = list(executor.map(cpu_bound_task, [100000] * 4))
process_time = time.time() - start_time
print(f"ProcessPoolExecutor: {process_time:.2f} seconds")
print(f"Speedup: {thread_time / process_time:.2f}x")
# I/O-bound comparison
print("\nI/O-bound task comparison:")
# ThreadPoolExecutor
start_time = time.time()
with ThreadPoolExecutor() as executor:
results = list(executor.map(io_bound_task, range(10)))
thread_time = time.time() - start_time
print(f"ThreadPoolExecutor: {thread_time:.2f} seconds")
# ProcessPoolExecutor
start_time = time.time()
with ProcessPoolExecutor() as executor:
results = list(executor.map(io_bound_task, range(10)))
process_time = time.time() - start_time
print(f"ProcessPoolExecutor: {process_time:.2f} seconds")
Real-World Examples
Web Scraping
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
import time
def fetch_webpage(url):
"""Fetch a webpage"""
try:
response = requests.get(url, timeout=5)
return f"{url}: {len(response.content)} bytes"
except Exception as e:
return f"{url}: Error - {e}"
def web_scraping_example():
urls = [
"https://www.google.com",
"https://www.github.com",
"https://www.stackoverflow.com",
"https://www.python.org"
]
with ThreadPoolExecutor(max_workers=4) as executor:
futures = {executor.submit(fetch_webpage, url): url for url in urls}
for future in as_completed(futures):
url = futures[future]
result = future.result()
print(result)
File Processing
from concurrent.futures import ProcessPoolExecutor
import os
import time
def process_file(filename):
"""Process a single file"""
# Simulate file processing
time.sleep(1)
return f"Processed {filename}"
def file_processing_example():
# Simulate list of files
files = [f"file_{i}.txt" for i in range(10)]
with ProcessPoolExecutor() as executor:
results = list(executor.map(process_file, files))
for result in results:
print(result)
Database Operations
from concurrent.futures import ThreadPoolExecutor
import sqlite3
import time
def database_operation(user_id):
"""Simulate database operation"""
time.sleep(0.5) # Simulate database query
return f"User {user_id} processed"
def database_example():
user_ids = list(range(1, 21))
with ThreadPoolExecutor(max_workers=5) as executor:
futures = [executor.submit(database_operation, user_id)
for user_id in user_ids]
results = [future.result() for future in futures]
for result in results:
print(result)
Best Practices
1. Choose the Right Executor
def choose_executor():
import multiprocessing
# For I/O-bound tasks (network, file operations)
def io_bound_example():
with ThreadPoolExecutor(max_workers=10) as executor:
# Your I/O-bound tasks here
pass
# For CPU-bound tasks (computations, data processing)
def cpu_bound_example():
with ProcessPoolExecutor(max_workers=multiprocessing.cpu_count()) as executor:
# Your CPU-bound tasks here
pass
2. Proper Resource Management
def resource_management():
# Always use context managers
with ThreadPoolExecutor() as executor:
# Your tasks here
pass
# Resources are automatically cleaned up
3. Handle Exceptions Properly
def exception_best_practices():
with ThreadPoolExecutor() as executor:
futures = [executor.submit(risky_task, i) for i in range(10)]
for future in as_completed(futures):
try:
result = future.result()
# Process successful result
except Exception as e:
# Handle exception appropriately
print(f"Task failed: {e}")
4. Use Appropriate Timeouts
def timeout_best_practices():
with ThreadPoolExecutor() as executor:
future = executor.submit(long_running_task, 10)
try:
result = future.result(timeout=5)
except TimeoutError:
future.cancel()
# Handle timeout appropriately
Common Pitfalls
1. Blocking the Main Thread
def avoid_blocking():
# Wrong: This blocks the main thread
with ThreadPoolExecutor() as executor:
future = executor.submit(long_task)
result = future.result() # Blocks here
# Better: Use as_completed or handle futures asynchronously
with ThreadPoolExecutor() as executor:
futures = [executor.submit(task, i) for i in range(10)]
for future in as_completed(futures):
result = future.result()
# Process result
2. Memory Issues with Large Data
def handle_large_data():
# Problem: Passing large data to processes
large_data = [i for i in range(1000000)]
# Solution: Process in chunks
chunk_size = 100000
chunks = [large_data[i:i + chunk_size]
for i in range(0, len(large_data), chunk_size)]
with ProcessPoolExecutor() as executor:
results = list(executor.map(process_chunk, chunks))
3. Pickling Issues
def avoid_pickling_issues():
# Problem: Lambda functions can't be pickled for ProcessPoolExecutor
# with ProcessPoolExecutor() as executor:
# results = executor.map(lambda x: x * 2, range(10)) # Fails
# Solution: Use regular functions
def double(x):
return x * 2
with ProcessPoolExecutor() as executor:
results = list(executor.map(double, range(10)))
Integration with Asyncio
import asyncio
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def sync_io_task():
"""Synchronous I/O task"""
import time
time.sleep(1)
return "I/O task completed"
def sync_cpu_task(n):
"""Synchronous CPU task"""
result = 0
for i in range(n):
result += i ** 2
return result
async def async_with_futures():
"""Use concurrent.futures with asyncio"""
loop = asyncio.get_event_loop()
# Use ThreadPoolExecutor for I/O tasks
with ThreadPoolExecutor() as thread_executor:
io_result = await loop.run_in_executor(thread_executor, sync_io_task)
print(io_result)
# Use ProcessPoolExecutor for CPU tasks
with ProcessPoolExecutor() as process_executor:
cpu_result = await loop.run_in_executor(process_executor, sync_cpu_task, 100000)
print(cpu_result)
# Run the async example
if __name__ == "__main__":
asyncio.run(async_with_futures())
Summary
| Feature | ThreadPoolExecutor | ProcessPoolExecutor |
|---|---|---|
| Use Case | I/O-bound tasks | CPU-bound tasks |
| Parallelism | Limited by GIL | True parallelism |
| Memory | Shared memory | Separate memory |
| Overhead | Low | Higher |
| Best For | Network, file I/O | Computations, data processing |
| API | Same high-level interface | Same high-level interface |
Key Takeaways
- Use ThreadPoolExecutor for I/O-bound tasks like network requests and file operations
- Use ProcessPoolExecutor for CPU-bound tasks that need true parallelism
- Always use context managers for automatic resource cleanup
- Handle exceptions properly to ensure robust error handling
- Use appropriate timeouts to prevent indefinite blocking
- Choose the right number of workers based on your workload
- Avoid common pitfalls like blocking the main thread
- Integrate with asyncio for mixed workloads
The concurrent.futures module provides a powerful, high-level interface for concurrent programming in Python, making it easier to write efficient, scalable applications.
Interview angle
- “What does
concurrent.futuresgive you?” - one API (submit,map,Future) over both threads and processes, so switching execution model is a one-line change. That’s its main value over usingthreadingormultiprocessingdirectly. - “
submitormap?” -mapfor a uniform function over an iterable with results in order;submitwhen you need per-task handles, want results as they finish viaas_completed, or the calls differ. - “How are exceptions handled?” - captured and re-raised when you call
result(). A future whose result is never retrieved swallows its exception silently, which is why fire-and-forget submission hides failures. - “How does this relate to asyncio?” -
loop.run_in_executortakes one of these executors, which is how async code offloads blocking or CPU-bound work.asyncio.to_threadis the modern shorthand for the thread case.