Python asyncio in Practice: Event Loops, Async Context Managers, Concurrency Limits and Debugging
Key takeaways
Async code that looks concurrent often is not, because one blocking call stalls the whole event loop. The post shows how to offload blocking work with run_in_executor, close aiohttp sessions properly, handle errors across many tasks, and fit asyncio into a FastAPI service.
Understanding Python Async Programming
Python’s asyncio module provides powerful tools for writing concurrent code using async/await syntax. This approach is particularly effective for I/O-bound applications where traditional threading might introduce complexity or performance overhead.
The core concept revolves around cooperative multitasking—functions voluntarily yield control when waiting for operations to complete, allowing other tasks to run. This eliminates many of the synchronization issues inherent in preemptive threading models.
It helps to understand why Python needed this model in the first place. CPython’s Global Interpreter Lock (GIL) means that only one thread executes Python bytecode at a time, regardless of how many OS threads you spin up. For CPU-bound work, this makes classic threading nearly useless as a way to speed things up — you get concurrency, not parallelism. But most real-world backend code isn’t actually CPU-bound; it’s waiting. It’s waiting on a database round trip, on a downstream HTTP call, on a disk read. During that wait, the GIL is released and the interpreter is free to do something else. Threading lets the OS scheduler decide what “something else” is, with all the locking and race-condition risk that implies. asyncio instead lets you decide, explicitly, at every await point — a coroutine only ever yields control where you write await, never in the middle of a line of pure Python. That single property is what makes async code easier to reason about than threaded code: you don’t need locks around shared in-memory state, because nothing preempts you between two statements that don’t contain an await.
The trade-off is that this cooperative model only pays off when your bottleneck really is I/O. If you accidentally call a blocking function — a synchronous requests.get(), a heavy time.sleep(), or a tight numeric loop — inside a coroutine without yielding, you don’t get a crash; you get a event loop that silently stalls every other coroutine scheduled on it until that call returns. This is the single most common asyncio bug in production code, and it rarely shows up in local testing because a small workload rarely stalls long enough to notice. It shows up under load, as a mysterious latency spike affecting requests that have nothing to do with the slow one.
Understanding when and how to use async programming will dramatically improve your application’s performance and resource utilization, especially in web applications, data processing pipelines, and API integrations. The rest of this guide works through the concrete patterns — event loops, context managers, file I/O, error handling, and framework integration — with an eye toward exactly where those gotchas tend to appear.
Core Concepts: Event Loops and Coroutines
The Event Loop Foundation
The event loop is the heart of asyncio, managing and executing coroutines, handling callbacks, and coordinating I/O operations. Calling an async def function does not run its body immediately — it returns a coroutine object, an inert value describing work that hasn’t started yet. Nothing happens until that object is awaited, wrapped in asyncio.create_task(), or passed to asyncio.run(). This is a frequent source of confusion for developers coming from JavaScript’s promises, which begin executing as soon as they’re created: in Python, basic_coroutine("x", 1) on its own does nothing but allocate an object, and forgetting the await produces a RuntimeWarning: coroutine was never awaited rather than any visible effect — a warning that’s easy to miss in noisy logs but that always signals a real bug.
asyncio.run() is the recommended entry point for a top-level script: it creates a fresh event loop, runs the given coroutine to completion, and then closes the loop, cleaning up any lingering async generators along the way. You should call it exactly once per process — nesting asyncio.run() calls, or calling it from inside a coroutine that’s already running on a loop, raises a RuntimeError. The example below contrasts sequential await calls, which block one after another exactly like synchronous code, with asyncio.gather(), which schedules both coroutines on the loop and lets their asyncio.sleep() waits overlap:
import asyncio
import time
async def basic_coroutine(name, duration):
"""Simple coroutine demonstrating async behavior"""
print(f"{name} starting")
await asyncio.sleep(duration) # Non-blocking sleep
print(f"{name} finished after {duration}s")
return f"Result from {name}"
async def main():
"""Demonstrate concurrent execution"""
# Sequential execution (slow)
start = time.time()
result1 = await basic_coroutine("Task 1", 2)
result2 = await basic_coroutine("Task 2", 1)
sequential_time = time.time() - start
# Concurrent execution (fast)
start = time.time()
results = await asyncio.gather(
basic_coroutine("Concurrent 1", 2),
basic_coroutine("Concurrent 2", 1)
)
concurrent_time = time.time() - start
print(f"Sequential: {sequential_time:.2f}s")
print(f"Concurrent: {concurrent_time:.2f}s")
return results
# Running the event loop
if __name__ == "__main__":
results = asyncio.run(main())
Notice that the sequential block above awaits basic_coroutine directly, while the concurrent block passes bare coroutine calls into gather(). Both are valid, but they behave very differently: await coro() runs that coroutine to completion before the next line executes, while gather(coro1(), coro2()) schedules both onto the event loop so their internal await points interleave. If either sleep were replaced with real network I/O, the concurrent version would take roughly as long as the slowest request rather than the sum of all requests — the entire performance case for asyncio in one example.
Advanced Coroutine Patterns
Real-world applications require more sophisticated patterns for managing coroutine lifecycles. The example above works fine for two known coroutines, but production code typically needs to fan out across dozens or hundreds of tasks, cap how many run at once so you don’t overwhelm a downstream service or exhaust file descriptors, enforce a deadline on the whole batch, and still get partial results back when some of those tasks fail. asyncio.Semaphore handles the concurrency cap, asyncio.wait_for handles the deadline, and gather(..., return_exceptions=True) turns individual task failures into values instead of letting one bad response abort the entire batch. The AsyncTaskManager class below combines all three:
import asyncio
from typing import List, Any, Optional
import aiohttp
import logging
class AsyncTaskManager:
"""Manages concurrent tasks with error handling and timeouts"""
def __init__(self, max_concurrent: int = 10):
self.max_concurrent = max_concurrent
self.semaphore = asyncio.Semaphore(max_concurrent)
async def execute_with_semaphore(self, coro):
"""Execute coroutine with concurrency control"""
async with self.semaphore:
return await coro
async def gather_with_timeout(
self,
tasks: List[asyncio.Task],
timeout: float = 30.0
) -> List[Any]:
"""Gather tasks with timeout and error handling"""
try:
return await asyncio.wait_for(
asyncio.gather(*tasks, return_exceptions=True),
timeout=timeout
)
except asyncio.TimeoutError:
# Cancel remaining tasks
for task in tasks:
if not task.done():
task.cancel()
raise
async def process_urls_batch(self, urls: List[str]) -> List[Optional[dict]]:
"""Process URLs concurrently with proper error handling"""
async with aiohttp.ClientSession() as session:
tasks = [
asyncio.create_task(
self.execute_with_semaphore(
self.fetch_url(session, url)
)
)
for url in urls
]
results = await self.gather_with_timeout(tasks)
# Process results and exceptions
processed_results = []
for i, result in enumerate(results):
if isinstance(result, Exception):
logging.error(f"Error processing {urls[i]}: {result}")
processed_results.append(None)
else:
processed_results.append(result)
return processed_results
async def fetch_url(self, session: aiohttp.ClientSession, url: str) -> dict:
"""Fetch URL with retries and proper error handling"""
max_retries = 3
for attempt in range(max_retries):
try:
async with session.get(url, timeout=10) as response:
if response.status == 200:
data = await response.json()
return {'url': url, 'data': data, 'status': response.status}
else:
raise aiohttp.ClientResponseError(
request_info=response.request_info,
history=response.history,
status=response.status
)
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
if attempt == max_retries - 1:
raise
await asyncio.sleep(2 ** attempt) # Exponential backoff
A few details in this class matter more than they might look. First, execute_with_semaphore receives an already-constructed coroutine object (self.fetch_url(session, url)), not a callable — building the coroutine outside the semaphore’s async with block is intentional and harmless, since constructing it doesn’t start execution; only the await inside execute_with_semaphore does. Second, wrapping each call in asyncio.create_task() before handing it to gather_with_timeout is what actually gives you concurrency: create_task schedules the coroutine onto the event loop immediately, independent of when it’s awaited, whereas passing bare coroutines to gather schedules them only at the moment gather runs. For this code the distinction is mostly stylistic, but it becomes important the moment you need to start work before you’re ready to wait for it — for example, kicking off a slow request and doing other setup in the meantime.
Third, look closely at what happens on a timeout: asyncio.wait_for raises asyncio.TimeoutError, and the except block explicitly cancels every task that isn’t already done before re-raising. This step is easy to skip and easy to regret skipping — a coroutine wrapped in create_task keeps running on the event loop even after the code that created it has moved on or raised an exception. Left uncancelled, those tasks continue consuming the semaphore slots and the HTTP session’s connection pool, and you’ll eventually see Task was destroyed but it is pending! warnings in the logs when the garbage collector finally reaps them — a sign that cleanup was missing somewhere upstream. The retry logic inside fetch_url layers exponential backoff on top of this: each failed attempt waits 2 ** attempt seconds before retrying, which is a reasonable default for transient network errors but should generally be paired with a maximum total retry budget in production, since three retries against a completely dead host can still add up to a noticeable multiple of your per-request timeout.
Async Context Managers and Resource Management
Proper Resource Cleanup
Async context managers ensure resources are properly cleaned up even when exceptions occur. The async with statement is to __aenter__/__aexit__ what with is to __enter__/__exit__ — the difference is that both entry and exit hooks can themselves await, which matters because closing a database connection or an HTTP session is frequently an I/O operation in its own right, not just a synchronous pointer decrement. Writing @asynccontextmanager on a generator function, as shown below, is usually the fastest way to get this behavior without hand-writing a class with __aenter__/__aexit__ methods: everything before yield runs on entry, everything after runs on exit, and the surrounding try/finally guarantees the cleanup code runs even if the code inside the async with block raises.
This matters more for async resources than for synchronous ones because the failure mode is quieter. A leaked synchronous file handle or socket is usually caught quickly by the OS or by ResourceWarning. A leaked aiohttp.ClientSession, by contrast, keeps its connector — and every open TCP connection it’s holding — alive for as long as the event loop runs, and aiohttp only emits a warning about it when the object is eventually garbage collected, which can be long after the leak actually happened. The pattern below wraps both a mock database connection and an aiohttp session pool in finally blocks specifically to make sure disconnect() and session.close() run unconditionally:
import asyncio
import aiofiles
from contextlib import asynccontextmanager
import aiohttp
from typing import AsyncGenerator
class DatabaseConnection:
"""Mock database connection with async operations"""
def __init__(self, connection_string: str):
self.connection_string = connection_string
self.connected = False
async def connect(self):
"""Simulate async connection establishment"""
await asyncio.sleep(0.1) # Simulate connection time
self.connected = True
print(f"Connected to {self.connection_string}")
async def disconnect(self):
"""Simulate async disconnection"""
await asyncio.sleep(0.05)
self.connected = False
print(f"Disconnected from {self.connection_string}")
async def execute_query(self, query: str):
"""Simulate async query execution"""
if not self.connected:
raise RuntimeError("Not connected to database")
await asyncio.sleep(0.1) # Simulate query time
return f"Result for: {query}"
@asynccontextmanager
async def database_transaction(connection_string: str) -> AsyncGenerator[DatabaseConnection, None]:
"""Async context manager for database operations"""
db = DatabaseConnection(connection_string)
try:
await db.connect()
yield db
except Exception as e:
print(f"Error in transaction: {e}")
raise
finally:
await db.disconnect()
@asynccontextmanager
async def http_session_pool() -> AsyncGenerator[aiohttp.ClientSession, None]:
"""Managed HTTP session with proper cleanup"""
connector = aiohttp.TCPConnector(limit=100, limit_per_host=10)
session = aiohttp.ClientSession(connector=connector)
try:
yield session
finally:
await session.close()
await connector.close()
async def process_data_with_resources():
"""Demonstrate proper resource management"""
# Database operations with automatic cleanup
async with database_transaction("postgresql://localhost/mydb") as db:
result = await db.execute_query("SELECT * FROM users")
print(f"Database result: {result}")
# HTTP operations with session management
async with http_session_pool() as session:
tasks = [
fetch_and_process(session, f"https://api.example.com/data/{i}")
for i in range(5)
]
results = await asyncio.gather(*tasks, return_exceptions=True)
return results
async def fetch_and_process(session: aiohttp.ClientSession, url: str):
"""Process individual HTTP requests"""
try:
async with session.get(url) as response:
data = await response.json()
# Process data here
return {'url': url, 'processed': True, 'items': len(data)}
except Exception as e:
return {'url': url, 'error': str(e)}
One subtlety worth calling out in http_session_pool: it explicitly creates a TCPConnector with limit=100 and limit_per_host=10 rather than letting aiohttp fall back to its defaults. Without an explicit cap, a burst of concurrent requests to the same slow host can open far more simultaneous connections than the server (or your own outbound network policy) is comfortable with, which tends to show up as connection resets or throttling under load rather than as an obvious error in your own code. Setting limit_per_host explicitly is a small habit that avoids a class of bug that’s genuinely hard to diagnose after the fact, because the symptoms appear on the server side of the connection, not the client.
File I/O and Data Processing Patterns
Async File Operations
Combining file I/O with async patterns enables efficient data processing, but it’s worth being precise about what “async” means here, because file I/O is a case where the abstraction leaks. Most operating systems don’t offer the same kind of true asynchronous, non-blocking file API that they offer for sockets — aiofiles, the library used below, actually delegates reads and writes to a background thread pool and exposes the result through an awaitable interface. That’s a meaningful implementation detail: it means aiofiles gives you the ergonomics of async code (no blocked event loop, code that composes with gather()) without literally avoiding OS-level blocking, because something still has to block, just off the main thread. For a handful of files, this is a fine trade. For truly enormous read/write volumes, a thread pool has its own overhead and thread-count ceiling, and it’s worth benchmarking against, say, a synchronous approach running in a ProcessPoolExecutor before assuming aiofiles is automatically the fastest option.
import asyncio
import aiofiles
import json
from pathlib import Path
from typing import List, Dict, Any
import csv
class AsyncFileProcessor:
"""Handles various async file operations efficiently"""
async def read_json_files(self, file_paths: List[Path]) -> List[Dict[str, Any]]:
"""Read multiple JSON files concurrently"""
tasks = [self.read_json_file(path) for path in file_paths]
results = await asyncio.gather(*tasks, return_exceptions=True)
# Filter out exceptions and return valid data
valid_results = [
result for result in results
if not isinstance(result, Exception)
]
return valid_results
async def read_json_file(self, file_path: Path) -> Dict[str, Any]:
"""Read a single JSON file asynchronously"""
try:
async with aiofiles.open(file_path, 'r') as file:
content = await file.read()
return json.loads(content)
except (FileNotFoundError, json.JSONDecodeError) as e:
raise ValueError(f"Error reading {file_path}: {e}")
async def write_processed_data(
self,
data: List[Dict[str, Any]],
output_path: Path
) -> None:
"""Write processed data to file asynchronously"""
processed_content = json.dumps(data, indent=2)
async with aiofiles.open(output_path, 'w') as file:
await file.write(processed_content)
async def process_large_csv(
self,
csv_path: Path,
processor_func,
batch_size: int = 1000
) -> List[Any]:
"""Process large CSV files in batches"""
results = []
async with aiofiles.open(csv_path, 'r') as file:
# Read header
header_line = await file.readline()
headers = header_line.strip().split(',')
batch = []
async for line in file:
if line.strip():
row_data = dict(zip(headers, line.strip().split(',')))
batch.append(row_data)
if len(batch) >= batch_size:
# Process batch asynchronously
batch_results = await self.process_batch(batch, processor_func)
results.extend(batch_results)
batch = []
# Process remaining items
if batch:
batch_results = await self.process_batch(batch, processor_func)
results.extend(batch_results)
return results
async def process_batch(self, batch: List[Dict], processor_func) -> List[Any]:
"""Process a batch of data concurrently"""
tasks = [processor_func(item) for item in batch]
return await asyncio.gather(*tasks, return_exceptions=True)
# Example usage
async def enrich_user_data(user_data: Dict[str, Any]) -> Dict[str, Any]:
"""Simulate enriching user data with external API calls"""
# Simulate API call delay
await asyncio.sleep(0.1)
# Add computed fields
user_data['processed_at'] = asyncio.get_event_loop().time()
user_data['score'] = hash(user_data.get('email', '')) % 100
return user_data
async def main_file_processing():
"""Demonstrate file processing patterns"""
processor = AsyncFileProcessor()
# Process multiple configuration files
config_files = [
Path('config1.json'),
Path('config2.json'),
Path('config3.json')
]
configs = await processor.read_json_files(config_files)
# Process large dataset
if Path('users.csv').exists():
processed_users = await processor.process_large_csv(
Path('users.csv'),
enrich_user_data,
batch_size=500
)
# Write results
await processor.write_processed_data(
processed_users,
Path('processed_users.json')
)
return configs, len(processed_users) if 'processed_users' in locals() else 0
Performance Optimization Strategies
Choosing the Right Concurrency Model
Different scenarios require different approaches to maximize performance, and picking the wrong one is a common source of disappointing benchmarks. The three tools below — pure asyncio, a ThreadPoolExecutor, and a ProcessPoolExecutor — look interchangeable because run_in_executor lets you await all three the same way, but they solve different problems and combining them incorrectly either wastes resources or doesn’t speed anything up at all.
Pure asyncio (no executor) is correct exactly when the work is I/O-bound and the library you’re calling has native async support, like aiohttp for HTTP or asyncpg for Postgres — the coroutine yields control at the await point and the event loop runs something else while the network does its thing. The moment you need to call a library that’s synchronous under the hood — requests, most ORMs’ default drivers, psycopg2 — running it directly inside a coroutine blocks the entire event loop, defeating the point of using asyncio at all. run_in_executor with a ThreadPoolExecutor is the fix: it hands the blocking call to a worker thread and awaits the result, which works because threads do release the GIL during I/O waits (the syscall itself isn’t holding the Python interpreter lock), so the event loop stays responsive while a thread blocks on requests.get().
A ProcessPoolExecutor solves a different problem entirely: genuine CPU-bound work, like the cpu_intensive_task below, where the bottleneck is Python bytecode execution, not waiting. Because of the GIL, running that function in a thread pool wouldn’t help — only one thread can execute Python code at a time regardless of how many threads exist, so a CPU-bound task in a thread pool just serializes behind whatever else is running. A process pool sidesteps the GIL entirely by using separate OS processes, each with its own interpreter and its own GIL, at the cost of the serialization overhead needed to pickle arguments and results across the process boundary — which is why it’s worth it for a chunk of 2,500 numbers but would be counterproductive for a function call that itself takes microseconds.
import asyncio
import threading
import multiprocessing
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import time
import requests # Synchronous HTTP library
from typing import Callable, List, Any
class PerformanceOptimizer:
"""Demonstrates different concurrency strategies"""
def __init__(self):
self.thread_pool = ThreadPoolExecutor(max_workers=10)
self.process_pool = ProcessPoolExecutor(max_workers=4)
async def io_bound_async(self, urls: List[str]) -> List[dict]:
"""Optimal for I/O-bound tasks: pure async"""
async with aiohttp.ClientSession() as session:
tasks = [self.fetch_async(session, url) for url in urls]
return await asyncio.gather(*tasks, return_exceptions=True)
async def fetch_async(self, session: aiohttp.ClientSession, url: str) -> dict:
"""Async HTTP request"""
try:
async with session.get(url, timeout=10) as response:
return {
'url': url,
'status': response.status,
'length': len(await response.text())
}
except Exception as e:
return {'url': url, 'error': str(e)}
async def io_bound_mixed(self, urls: List[str]) -> List[dict]:
"""Mix async with thread pool for blocking libraries"""
loop = asyncio.get_event_loop()
tasks = [
loop.run_in_executor(
self.thread_pool,
self.fetch_sync,
url
)
for url in urls
]
return await asyncio.gather(*tasks, return_exceptions=True)
def fetch_sync(self, url: str) -> dict:
"""Synchronous HTTP request (for demonstration)"""
try:
response = requests.get(url, timeout=10)
return {
'url': url,
'status': response.status_code,
'length': len(response.text)
}
except Exception as e:
return {'url': url, 'error': str(e)}
async def cpu_bound_async(self, data: List[int]) -> List[int]:
"""CPU-bound tasks using process pool"""
loop = asyncio.get_event_loop()
# Split work into chunks for parallel processing
chunk_size = len(data) // multiprocessing.cpu_count()
chunks = [
data[i:i + chunk_size]
for i in range(0, len(data), chunk_size)
]
# Process chunks in parallel
tasks = [
loop.run_in_executor(
self.process_pool,
self.cpu_intensive_task,
chunk
)
for chunk in chunks
]
results = await asyncio.gather(*tasks)
# Flatten results
flattened = []
for chunk_result in results:
flattened.extend(chunk_result)
return flattened
@staticmethod
def cpu_intensive_task(numbers: List[int]) -> List[int]:
"""Simulate CPU-intensive processing"""
return [n * n + n // 2 for n in numbers if n % 2 == 0]
async def benchmark_approaches(self):
"""Compare different concurrency approaches"""
urls = [f"https://httpbin.org/delay/{i % 3}" for i in range(10)]
test_data = list(range(10000))
results = {}
# I/O-bound: Pure async
start = time.time()
await self.io_bound_async(urls)
results['io_async'] = time.time() - start
# I/O-bound: Mixed with threads
start = time.time()
await self.io_bound_mixed(urls)
results['io_mixed'] = time.time() - start
# CPU-bound: Process pool
start = time.time()
await self.cpu_bound_async(test_data)
results['cpu_async'] = time.time() - start
return results
# Memory and resource optimization
class AsyncResourceManager:
"""Manages resources efficiently in async applications"""
def __init__(self):
self.active_tasks = set()
self.completed_tasks = []
async def managed_task_execution(self, task_generators: List[Callable]):
"""Execute tasks with proper cleanup and monitoring"""
for generator in task_generators:
task = asyncio.create_task(generator())
self.active_tasks.add(task)
# Add completion callback for cleanup
task.add_done_callback(self.task_completed)
# Wait for all tasks with timeout
try:
await asyncio.wait(
self.active_tasks,
timeout=60.0,
return_when=asyncio.ALL_COMPLETED
)
except asyncio.TimeoutError:
# Cancel remaining tasks
for task in self.active_tasks:
if not task.done():
task.cancel()
return len(self.completed_tasks)
def task_completed(self, task: asyncio.Task):
"""Callback for task completion"""
self.active_tasks.discard(task)
if task.exception():
logging.error(f"Task failed: {task.exception()}")
else:
self.completed_tasks.append(task.result())
Error Handling and Debugging Patterns
Comprehensive Error Management
Robust async applications require sophisticated error handling strategies, because errors in concurrent code have a way of getting lost or duplicated if you don’t handle them deliberately. The AsyncErrorHandler class below bundles three patterns that are easy to write individually but easy to forget to combine: retry-with-backoff for transient failures, a context manager for structured operation logging, and a gather wrapper that separates successes from failures instead of letting one exception propagate and cancel every sibling task.
That last point is worth dwelling on, because it’s the detail that trips people up most often when moving from sequential to concurrent error handling. asyncio.gather() without return_exceptions=True raises the first exception it encounters and, critically, does not wait for or cancel the other still-running tasks by default — they keep executing in the background, detached from the code that launched them, which is a subtler version of the same leaked-task problem discussed earlier. Passing return_exceptions=True, as safe_gather does below, changes the contract entirely: every task runs to completion (or failure) and the exception is returned as a value in the results list rather than raised. That shifts responsibility onto your code to check each result for isinstance(result, Exception), but in exchange you get a batch operation that degrades gracefully — three failed requests out of fifty don’t discard the other forty-seven successful ones.
import asyncio
import logging
from contextlib import asynccontextmanager
from typing import Optional, Callable, Any
import functools
class AsyncErrorHandler:
"""Comprehensive error handling for async operations"""
def __init__(self):
self.setup_logging()
self.error_counts = {}
self.circuit_breakers = {}
def setup_logging(self):
"""Configure logging for async operations"""
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
self.logger = logging.getLogger(__name__)
def retry_async(self, max_retries: int = 3, delay: float = 1.0):
"""Decorator for async function retries with exponential backoff"""
def decorator(func):
@functools.wraps(func)
async def wrapper(*args, **kwargs):
for attempt in range(max_retries):
try:
return await func(*args, **kwargs)
except Exception as e:
self.error_counts[func.__name__] = self.error_counts.get(func.__name__, 0) + 1
if attempt == max_retries - 1:
self.logger.error(f"Final attempt failed for {func.__name__}: {e}")
raise
backoff_delay = delay * (2 ** attempt)
self.logger.warning(
f"Attempt {attempt + 1} failed for {func.__name__}: {e}. "
f"Retrying in {backoff_delay}s"
)
await asyncio.sleep(backoff_delay)
return None
return wrapper
return decorator
@asynccontextmanager
async def error_context(self, operation_name: str):
"""Context manager for operation-level error handling"""
start_time = asyncio.get_event_loop().time()
try:
self.logger.info(f"Starting operation: {operation_name}")
yield
except Exception as e:
duration = asyncio.get_event_loop().time() - start_time
self.logger.error(
f"Operation {operation_name} failed after {duration:.2f}s: {e}"
)
raise
else:
duration = asyncio.get_event_loop().time() - start_time
self.logger.info(
f"Operation {operation_name} completed successfully in {duration:.2f}s"
)
async def safe_gather(self, *coroutines, return_exceptions: bool = True):
"""Gather with detailed error reporting"""
results = await asyncio.gather(*coroutines, return_exceptions=return_exceptions)
successful = []
failed = []
for i, result in enumerate(results):
if isinstance(result, Exception):
failed.append({'index': i, 'error': result})
self.logger.error(f"Task {i} failed: {result}")
else:
successful.append(result)
if failed:
self.logger.warning(f"{len(failed)} out of {len(results)} tasks failed")
return {
'successful': successful,
'failed': failed,
'success_rate': len(successful) / len(results)
}
# Example usage with comprehensive error handling
class RobustAsyncService:
"""Production-ready async service with error handling"""
def __init__(self):
self.error_handler = AsyncErrorHandler()
self.session: Optional[aiohttp.ClientSession] = None
async def __aenter__(self):
self.session = aiohttp.ClientSession()
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
if self.session:
await self.session.close()
@AsyncErrorHandler().retry_async(max_retries=3, delay=1.0)
async def fetch_with_retries(self, url: str) -> dict:
"""HTTP request with automatic retries"""
if not self.session:
raise RuntimeError("Service not initialized")
async with self.session.get(url, timeout=10) as response:
if response.status >= 400:
raise aiohttp.ClientResponseError(
request_info=response.request_info,
history=response.history,
status=response.status
)
return await response.json()
async def process_urls_safely(self, urls: List[str]) -> dict:
"""Process URLs with comprehensive error handling"""
async with self.error_handler.error_context("batch_url_processing"):
tasks = [self.fetch_with_retries(url) for url in urls]
return await self.error_handler.safe_gather(*tasks)
Look closely at the @AsyncErrorHandler().retry_async(max_retries=3, delay=1.0) decorator on fetch_with_retries — it constructs a brand-new AsyncErrorHandler instance at class-definition time, entirely separate from self.error_handler used everywhere else in the class. Both instances log through the same module-level logging.getLogger(__name__), so the retry warnings still appear in the output and the bug is easy to miss by reading the logs alone. But error_counts on that throwaway instance is never read anywhere, and if you extended AsyncErrorHandler with per-instance state — a circuit breaker that opens after N failures, say — it would silently track failures against an object nothing else can see. The fix is to decorate with a reference to self.error_handler instead, which usually means moving the decoration into __init__ or using a plain module-level retry function that doesn’t carry instance state at all. It’s a useful reminder that Python decorators evaluate their arguments once, at definition time, not per-call — so any decorator argument that looks like it should be “the current instance” needs to actually be the current instance.
# Usage example
async def main_error_handling_demo():
"""Demonstrate error handling patterns"""
urls = [
"https://httpbin.org/json",
"https://httpbin.org/status/500", # Will fail
"https://httpbin.org/delay/2",
"https://invalid-url-that-will-fail.com" # Will fail
]
async with RobustAsyncService() as service:
results = await service.process_urls_safely(urls)
print(f"Success rate: {results['success_rate']:.1%}")
print(f"Successful requests: {len(results['successful'])}")
print(f"Failed requests: {len(results['failed'])}")
return results
if __name__ == "__main__":
# Run the comprehensive example
results = asyncio.run(main_error_handling_demo())
Integration with Web Frameworks
Modern Python web frameworks leverage asyncio for high-performance applications, but the framework only gets you half the benefit — the other half depends on whether the code inside your route handlers actually respects the async contract. FastAPI (and Starlette underneath it) runs async def endpoints directly on the event loop alongside every other concurrent request; a synchronous def endpoint, by contrast, is automatically dispatched to a thread pool so it doesn’t block the loop, which is a deliberate safety net but means you should default to async def and only fall back to a plain def when you’re intentionally calling blocking code that has no async equivalent.
The two endpoints below illustrate the same GIL-driven trade-off from the concurrency-model section, applied to a real API surface. process_data_endpoint awaits I/O directly and returns once the work is done, so its response time is bounded by however long the underlying requests take. create_background_task takes the opposite approach: it returns an HTTP response immediately and lets long_running_process continue running after the response has already been sent to the client. This matters for anything a caller shouldn’t have to wait on — sending a notification email, regenerating a cache, kicking off the CPU-bound benchmark shown earlier — but it comes with a real operational cost: BackgroundTasks runs in-process, so if the server restarts or crashes before the task finishes, that work is silently lost with no retry and no record that it was ever scheduled. For anything where losing the task would actually matter, a real task queue (Celery, arq, or a hosted equivalent) that persists the job outside the web process is worth the extra infrastructure.
# FastAPI integration example
from fastapi import FastAPI, BackgroundTasks
import asyncio
from typing import List
app = FastAPI()
@app.get("/process-data")
async def process_data_endpoint(urls: List[str]):
"""API endpoint demonstrating async processing"""
async with RobustAsyncService() as service:
results = await service.process_urls_safely(urls)
return {
"processed": len(results['successful']),
"failed": len(results['failed']),
"success_rate": results['success_rate']
}
@app.post("/background-task")
async def create_background_task(background_tasks: BackgroundTasks):
"""Demonstrate background task processing"""
background_tasks.add_task(long_running_process)
return {"message": "Background task started"}
async def long_running_process():
"""Example of long-running async process"""
optimizer = PerformanceOptimizer()
results = await optimizer.benchmark_approaches()
# Process results, update database, send notifications, etc.
return results
Python’s asyncio ecosystem provides powerful tools for building efficient, scalable applications. The key is understanding when to use async patterns versus traditional approaches, implementing proper error handling, and choosing the right concurrency model for your specific use case.
Success with asyncio comes from starting with clear I/O-bound use cases, gradually building complexity, and always measuring performance to ensure your async code delivers the expected benefits. The patterns shown here provide a solid foundation for production-ready async Python applications.
Related Articles
- Finding and Fixing Slow React Renders
- C++ vs Python: Which Language Should You Learn?
- C++ vs Go: Performance, Concurrency Models, and When to Choose Each