backend / async concurrency / 15_concurrent_futures.md

Concurrent.futures in Python: A Comprehensive Guide

4 interview angles 9 min read source

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

  1. Executor Classes: Abstract base classes for executing callables
  2. Future Objects: Represent pending computations
  3. Context Managers: Automatic resource management
  4. 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

  1. Use ThreadPoolExecutor for I/O-bound tasks like network requests and file operations
  2. Use ProcessPoolExecutor for CPU-bound tasks that need true parallelism
  3. Always use context managers for automatic resource cleanup
  4. Handle exceptions properly to ensure robust error handling
  5. Use appropriate timeouts to prevent indefinite blocking
  6. Choose the right number of workers based on your workload
  7. Avoid common pitfalls like blocking the main thread
  8. 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.futures give 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 using threading or multiprocessing directly.
  • submit or map?” - map for a uniform function over an iterable with results in order; submit when you need per-task handles, want results as they finish via as_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_executor takes one of these executors, which is how async code offloads blocking or CPU-bound work. asyncio.to_thread is the modern shorthand for the thread case.