Claude
Skill
python-async-ops
Python asyncio patterns for concurrent programming. Triggers on: asyncio, async, await, coroutine, gather, semaphore, TaskGroup, event loop, aiohttp, concurrent.
Virus-scanned
Reviewed automatically before listing.
Download
0xdarkmatter-claude-mods-skills_python-async-ops-3dfaf0b.zip · 22 KB
Install
skills CLI
npx skills add https://github.com/0xDarkMatter/claude-mods/tree/main/skills/python-async-ops
Claude Code
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install 0xdarkmatter-claude-mods@llmmart
Git
git clone https://github.com/0xDarkMatter/claude-mods.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole 0xdarkmatter/claude-mods collection as a plugin from our marketplace. Git is the plain clone.
Skill manifest
Python Async Patterns
Asyncio patterns for concurrent Python programming.
When to Use Async vs Sync
| Use Async When | Use Sync When |
|---|---|
| I/O-bound operations (HTTP, DB, files) | CPU-bound computations |
| High concurrency (100s+ connections) | Simple scripts, one-off tasks |
| WebSocket/streaming connections | Small data processing |
| Microservices with network calls | Single sequential operations |
Decision tree:
- Is it CPU-bound? → Sync (or multiprocessing)
- Is it I/O-bound with high concurrency? → Async
- Is it simple I/O with few connections? → Sync is fine
Core Concepts
import asyncio
# Coroutine (must be awaited)
async def fetch(url: str) -> str:
async with aiohttp.ClientSession() as session:
async with session.get(url) as response:
return await response.text()
# Entry point
async def main():
result = await fetch("https://example.com")
return result
asyncio.run(main())
Pattern 1: Concurrent with gather
async def fetch_all(urls: list[str]) -> list[str]:
"""Fetch multiple URLs concurrently."""
async with aiohttp.ClientSession() as session:
tasks = [fetch_one(session, url) for url in urls]
return await asyncio.gather(*tasks, return_exceptions=True)
Pattern 2: Bounded Concurrency
async def fetch_with_limit(urls: list[str], limit: int = 10):
"""Limit concurrent requests."""
semaphore = asyncio.Semaphore(limit)
async def bounded_fetch(url):
async with semaphore:
return await fetch_one(url)
return await asyncio.gather(*[bounded_fetch(url) for url in urls])
Pattern 3: TaskGroup (Python 3.11+)
async def process_items(items):
"""Structured concurrency with automatic cleanup."""
async with asyncio.TaskGroup() as tg:
for item in items:
tg.create_task(process_one(item))
# All tasks complete here, or exception raised
Pattern 4: Timeout
async def with_timeout():
try:
async with asyncio.timeout(5.0): # Python 3.11+
result = await slow_operation()
except asyncio.TimeoutError:
result = None
return result
Critical Warnings
# WRONG - blocks event loop
async def bad():
time.sleep(5) # Never use time.sleep!
requests.get(url) # Blocking I/O!
# CORRECT
async def good():
await asyncio.sleep(5)
async with aiohttp.ClientSession() as s:
await s.get(url)
# WRONG - orphaned task
async def bad():
asyncio.create_task(work()) # May be garbage collected!
# CORRECT - keep reference
async def good():
task = asyncio.create_task(work())
await task
Quick Reference
| Pattern | Use Case |
|---|---|
gather(*tasks) |
Multiple independent operations |
Semaphore(n) |
Rate limiting, resource constraints |
TaskGroup() |
Structured concurrency (3.11+) |
Queue() |
Producer-consumer |
timeout(s) |
Timeout wrapper (3.11+) |
Lock() |
Shared mutable state |
Async Context Manager
from contextlib import asynccontextmanager
@asynccontextmanager
async def managed_connection():
conn = await create_connection()
try:
yield conn
finally:
await conn.close()
Additional Resources
For detailed patterns, load:
./references/concurrency-patterns.md- Queue, Lock, producer-consumer./references/aiohttp-patterns.md- HTTP client/server patterns./references/mixing-sync-async.md- run_in_executor, thread pools./references/debugging-async.md- Debug mode, profiling, finding issues./references/production-patterns.md- Graceful shutdown, health checks, signal handling./references/error-handling.md- Retry with backoff, circuit breakers, partial failures./references/performance.md- uvloop, connection pooling, buffer sizing
Scripts
./scripts/find-blocking-calls.sh- Scan code for blocking calls in async functions
Assets
./assets/async-project-template.py- Production-ready async app skeleton
See Also
Prerequisites:
python-typing-ops- Type hints for async functions
Related Skills:
python-fastapi-ops- Async web APIspython-observability-ops- Async logging and tracingpython-database-ops- Async database access
Files (claude-mods)
-
assets
-
async-project-template.py 3.6 KB
#!/usr/bin/env python3 """ Async Python Project Template Production-ready async application structure with: - Proper session management - Graceful shutdown - Structured concurrency - Error handling - Logging """ import asyncio import logging import signal from contextlib import asynccontextmanager from typing import Any import aiohttp # Configure logging logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s" ) logger = logging.getLogger(__name__) class AsyncApp: """Main application class with lifecycle management.""" def __init__(self): self.session: aiohttp.ClientSession | None = None self.running = False self._tasks: set[asyncio.Task] = set() async def start(self): """Initialize resources.""" logger.info("Starting application...") self.session = aiohttp.ClientSession( timeout=aiohttp.ClientTimeout(total=30), connector=aiohttp.TCPConnector(limit=100), ) self.running = True logger.info("Application started") async def stop(self): """Cleanup resources.""" logger.info("Stopping application...") self.running = False # Cancel background tasks for task in self._tasks: task.cancel() if self._tasks: await asyncio.gather(*self._tasks, return_exceptions=True) # Close session if self.session: await self.session.close() logger.info("Application stopped") def create_task(self, coro) -> asyncio.Task: """Create a tracked background task.""" task = asyncio.create_task(coro) self._tasks.add(task) task.add_done_callback(self._tasks.discard) return task async def fetch(self, url: str) -> dict[str, Any] | None: """Fetch URL with error handling.""" if not self.session: raise RuntimeError("App not started") try: async with self.session.get(url) as response: response.raise_for_status() return await response.json() except aiohttp.ClientError as e: logger.error(f"Request failed: {e}") return None async def fetch_many( self, urls: list[str], concurrency: int = 10 ) -> list[dict[str, Any] | None]: """Fetch multiple URLs with bounded concurrency.""" semaphore = asyncio.Semaphore(concurrency) async def bounded_fetch(url: str): async with semaphore: return await self.fetch(url) return await asyncio.gather(*[bounded_fetch(url) for url in urls]) @asynccontextmanager async def create_app(): """Context manager for app lifecycle.""" app = AsyncApp() try: await app.start() yield app finally: await app.stop() async def main(): """Main entry point.""" # Setup signal handlers for graceful shutdown loop = asyncio.get_running_loop() stop_event = asyncio.Event() def signal_handler(): logger.info("Received shutdown signal") stop_event.set() for sig in (signal.SIGTERM, signal.SIGINT): loop.add_signal_handler(sig, signal_handler) async with create_app() as app: # Example: Fetch some URLs urls = [ "https://httpbin.org/json", "https://httpbin.org/uuid", ] results = await app.fetch_many(urls) for url, result in zip(urls, results): logger.info(f"{url}: {result}") # Keep running until shutdown signal # await stop_event.wait() if __name__ == "__main__": asyncio.run(main())
-
-
references
-
aiohttp-patterns.md 6.1 KB
# aiohttp Patterns HTTP client and server patterns with aiohttp. ## Client Session Best Practices ```python import aiohttp # WRONG - creates session per request async def bad_fetch(url): async with aiohttp.ClientSession() as session: async with session.get(url) as response: return await response.text() # CORRECT - reuse session async def fetch_all(urls: list[str]) -> list[str]: async with aiohttp.ClientSession() as session: tasks = [fetch_one(session, url) for url in urls] return await asyncio.gather(*tasks) async def fetch_one(session: aiohttp.ClientSession, url: str) -> str: async with session.get(url) as response: return await response.text() ``` ## Connection Pooling ```python import aiohttp # Configure connection pool connector = aiohttp.TCPConnector( limit=100, # Max connections limit_per_host=10, # Max per host ttl_dns_cache=300, # DNS cache TTL ) async with aiohttp.ClientSession(connector=connector) as session: # Use session pass ``` ## Timeout Configuration ```python import aiohttp timeout = aiohttp.ClientTimeout( total=30, # Total timeout connect=10, # Connection timeout sock_read=10, # Read timeout sock_connect=10, # Socket connect timeout ) async with aiohttp.ClientSession(timeout=timeout) as session: async with session.get(url) as response: return await response.text() ``` ## Request Methods ```python async with aiohttp.ClientSession() as session: # GET async with session.get(url, params={'key': 'value'}) as r: data = await r.json() # POST JSON async with session.post(url, json={'key': 'value'}) as r: data = await r.json() # POST form data async with session.post(url, data={'key': 'value'}) as r: data = await r.text() # PUT async with session.put(url, json={'key': 'value'}) as r: pass # DELETE async with session.delete(url) as r: pass # With headers headers = {'Authorization': 'Bearer token'} async with session.get(url, headers=headers) as r: pass ``` ## Response Handling ```python async with session.get(url) as response: # Status print(response.status) # 200 print(response.reason) # OK # Headers print(response.headers['Content-Type']) # Body text = await response.text() json_data = await response.json() bytes_data = await response.read() # Streaming async for chunk in response.content.iter_chunked(1024): process(chunk) ``` ## Error Handling ```python import aiohttp async def safe_fetch(session, url): try: async with session.get(url) as response: response.raise_for_status() return await response.json() except aiohttp.ClientResponseError as e: print(f"HTTP error: {e.status}") except aiohttp.ClientConnectionError: print("Connection error") except aiohttp.ClientTimeout: print("Request timed out") except Exception as e: print(f"Unexpected error: {e}") return None ``` ## Retry with Backoff ```python async def fetch_with_retry( session: aiohttp.ClientSession, url: str, max_retries: int = 3 ) -> dict | None: for attempt in range(max_retries): try: async with session.get(url) as response: response.raise_for_status() return await response.json() except (aiohttp.ClientError, asyncio.TimeoutError): if attempt == max_retries - 1: raise await asyncio.sleep(2 ** attempt) # Exponential backoff return None ``` ## File Upload ```python async def upload_file(session, url, file_path): with open(file_path, 'rb') as f: data = aiohttp.FormData() data.add_field('file', f, filename='upload.txt') async with session.post(url, data=data) as response: return await response.json() ``` ## File Download ```python async def download_file(session, url, dest_path): async with session.get(url) as response: with open(dest_path, 'wb') as f: async for chunk in response.content.iter_chunked(8192): f.write(chunk) ``` ## WebSocket Client ```python async def websocket_client(url): async with aiohttp.ClientSession() as session: async with session.ws_connect(url) as ws: # Send message await ws.send_str("Hello") # Receive messages async for msg in ws: if msg.type == aiohttp.WSMsgType.TEXT: print(f"Received: {msg.data}") elif msg.type == aiohttp.WSMsgType.ERROR: break ``` ## Simple aiohttp Server ```python from aiohttp import web async def handle_get(request): name = request.match_info.get('name', 'World') return web.json_response({'message': f'Hello, {name}'}) async def handle_post(request): data = await request.json() return web.json_response({'received': data}) app = web.Application() app.router.add_get('/', handle_get) app.router.add_get('/{name}', handle_get) app.router.add_post('/data', handle_post) if __name__ == '__main__': web.run_app(app, port=8080) ``` ## Server Middleware ```python from aiohttp import web @web.middleware async def error_middleware(request, handler): try: response = await handler(request) return response except web.HTTPException: raise except Exception as e: return web.json_response( {'error': str(e)}, status=500 ) @web.middleware async def logging_middleware(request, handler): print(f"{request.method} {request.path}") response = await handler(request) print(f"Response: {response.status}") return response app = web.Application(middlewares=[logging_middleware, error_middleware]) ``` ## Session State ```python from aiohttp import web async def init_db(app): app['db'] = await create_db_pool() async def cleanup_db(app): await app['db'].close() app = web.Application() app.on_startup.append(init_db) app.on_cleanup.append(cleanup_db) async def handler(request): db = request.app['db'] # Use db connection ``` -
concurrency-patterns.md 6 KB
# Async Concurrency Patterns Advanced concurrency patterns for Python asyncio. ## Producer-Consumer with Queue ```python import asyncio async def producer(queue: asyncio.Queue, items): """Produce items to queue.""" for item in items: await queue.put(item) await queue.put(None) # Sentinel to signal completion async def consumer(queue: asyncio.Queue, name: str): """Consume items from queue.""" while True: item = await queue.get() if item is None: queue.task_done() break await process(item) queue.task_done() async def main(): queue = asyncio.Queue(maxsize=100) # Backpressure # Run producer and multiple consumers await asyncio.gather( producer(queue, items), consumer(queue, "worker-1"), consumer(queue, "worker-2"), consumer(queue, "worker-3"), ) ``` ## Sharing State with Lock ```python import asyncio # WRONG - race condition even in async! counter = 0 async def increment(): global counter temp = counter await asyncio.sleep(0) # Context switch point! counter = temp + 1 # CORRECT - use Lock lock = asyncio.Lock() async def safe_increment(): global counter async with lock: counter += 1 ``` ## Event Signaling ```python import asyncio async def waiter(event: asyncio.Event): print("Waiting for event...") await event.wait() print("Event received!") async def setter(event: asyncio.Event): await asyncio.sleep(2) event.set() print("Event set!") async def main(): event = asyncio.Event() await asyncio.gather( waiter(event), setter(event), ) ``` ## Condition Variable ```python import asyncio async def consumer(condition: asyncio.Condition, data: list): async with condition: await condition.wait_for(lambda: len(data) > 0) item = data.pop(0) return item async def producer(condition: asyncio.Condition, data: list): async with condition: data.append("new item") condition.notify() ``` ## Barrier (Python 3.11+) ```python import asyncio async def worker(barrier: asyncio.Barrier, name: str): print(f"{name}: Starting work") await asyncio.sleep(1) print(f"{name}: Waiting at barrier") await barrier.wait() print(f"{name}: Continuing after barrier") async def main(): barrier = asyncio.Barrier(3) await asyncio.gather( worker(barrier, "A"), worker(barrier, "B"), worker(barrier, "C"), ) ``` ## Cancellation Handling ```python async def cancellable_task(): try: while True: await do_work() except asyncio.CancelledError: # Cleanup on cancellation await cleanup() raise # Re-raise to propagate # Cancel a task task = asyncio.create_task(cancellable_task()) task.cancel() try: await task except asyncio.CancelledError: print("Task was cancelled") ``` ## Task Completion Callbacks ```python def on_complete(task: asyncio.Task): if task.exception(): print(f"Task failed: {task.exception()}") else: print(f"Task result: {task.result()}") task = asyncio.create_task(some_work()) task.add_done_callback(on_complete) ``` ## Running Tasks in Background ```python # Keep track of background tasks background_tasks = set() async def start_background_task(coro): task = asyncio.create_task(coro) background_tasks.add(task) task.add_done_callback(background_tasks.discard) return task async def cleanup(): for task in background_tasks: task.cancel() await asyncio.gather(*background_tasks, return_exceptions=True) ``` ## Async Iterator/Generator ```python async def async_range(n: int): """Async generator.""" for i in range(n): await asyncio.sleep(0.1) yield i # Usage async for value in async_range(10): print(value) # Async comprehension results = [x async for x in async_range(10)] ``` ## Streaming with AsyncIterator ```python class AsyncIterator: def __init__(self, items): self.items = items self.index = 0 def __aiter__(self): return self async def __anext__(self): if self.index >= len(self.items): raise StopAsyncIteration item = self.items[self.index] self.index += 1 await asyncio.sleep(0) # Yield control return item ``` ## Shield from Cancellation ```python async def critical_operation(): # This operation must complete even if outer task is cancelled try: result = await asyncio.shield(important_work()) except asyncio.CancelledError: # Shield was cancelled, but important_work continues result = await important_work() # Wait for it return result ``` ## Wait with First Completed ```python async def first_response(urls: list[str]): """Return first successful response.""" tasks = [asyncio.create_task(fetch(url)) for url in urls] done, pending = await asyncio.wait( tasks, return_when=asyncio.FIRST_COMPLETED ) # Cancel remaining tasks for task in pending: task.cancel() return done.pop().result() ``` ## Debounce Pattern ```python class Debouncer: def __init__(self, delay: float): self.delay = delay self.task: asyncio.Task | None = None async def debounce(self, coro): if self.task: self.task.cancel() try: await self.task except asyncio.CancelledError: pass async def delayed(): await asyncio.sleep(self.delay) await coro self.task = asyncio.create_task(delayed()) ``` ## Retry Pattern ```python async def retry(coro_func, max_retries: int = 3, delay: float = 1.0): """Retry coroutine with exponential backoff.""" for attempt in range(max_retries): try: return await coro_func() except Exception as e: if attempt == max_retries - 1: raise await asyncio.sleep(delay * (2 ** attempt)) ``` -
debugging-async.md 5.8 KB
# Debugging Async Python Techniques for debugging asyncio applications. ## Enable Debug Mode ```python import asyncio # Option 1: Environment variable # PYTHONASYNCIODEBUG=1 python script.py # Option 2: In code asyncio.run(main(), debug=True) # Option 3: On running loop loop = asyncio.get_running_loop() loop.set_debug(True) ``` ## Debug Mode Features When enabled: - Slow callbacks (>100ms) are logged - Unawaited coroutines are detected - Resource warnings for unclosed resources - More detailed tracebacks ## Finding Slow Callbacks ```python import asyncio import logging # Enable asyncio debug logging logging.getLogger("asyncio").setLevel(logging.DEBUG) # Custom slow callback threshold loop = asyncio.get_event_loop() loop.slow_callback_duration = 0.05 # 50ms ``` ## Detecting Unawaited Coroutines ```python import warnings warnings.filterwarnings("error", category=RuntimeWarning) # Now this will raise instead of warn: async def main(): some_coroutine() # RuntimeWarning -> Exception! ``` ## Task Introspection ```python import asyncio async def debug_tasks(): # Get all tasks all_tasks = asyncio.all_tasks() print(f"Total tasks: {len(all_tasks)}") for task in all_tasks: print(f"Task: {task.get_name()}") print(f" Done: {task.done()}") print(f" Cancelled: {task.cancelled()}") # Get stack if not task.done(): stack = task.get_stack() for frame in stack: print(f" {frame}") # Get current task current = asyncio.current_task() ``` ## Tracing Coroutines ```python import sys def trace_coroutines(frame, event, arg): if event == "call" and frame.f_code.co_flags & 0x80: # CO_COROUTINE print(f"Coroutine called: {frame.f_code.co_name}") return trace_coroutines sys.settrace(trace_coroutines) ``` ## asyncio Debug Logger ```python import logging # Detailed asyncio logging logging.basicConfig(level=logging.DEBUG) logger = logging.getLogger("asyncio") logger.setLevel(logging.DEBUG) # Custom handler handler = logging.StreamHandler() handler.setFormatter(logging.Formatter( "%(asctime)s - %(name)s - %(levelname)s - %(message)s" )) logger.addHandler(handler) ``` ## Profiling Async Code ### With cProfile ```python import asyncio import cProfile import pstats async def main(): await some_work() # Profile profiler = cProfile.Profile() profiler.enable() asyncio.run(main()) profiler.disable() # Print stats stats = pstats.Stats(profiler) stats.sort_stats("cumtime") stats.print_stats(20) ``` ### With yappi (async-aware) ```python import yappi import asyncio yappi.set_clock_type("wall") # or "cpu" yappi.start() asyncio.run(main()) yappi.stop() # Get stats for coroutines func_stats = yappi.get_func_stats() func_stats.print_all() # Async-specific stats asyncio_stats = yappi.get_func_stats( filter_callback=lambda x: asyncio.iscoroutinefunction(x.full_name) ) ``` ## Finding Memory Leaks ```python import asyncio import gc import tracemalloc tracemalloc.start() async def main(): # ... your code ... pass asyncio.run(main()) # Get memory snapshot snapshot = tracemalloc.take_snapshot() top_stats = snapshot.statistics("lineno") print("Top 10 memory allocations:") for stat in top_stats[:10]: print(stat) # Find leaking tasks gc.collect() for obj in gc.get_objects(): if isinstance(obj, asyncio.Task): print(f"Leaked task: {obj}") ``` ## Common Issues and Solutions ### Issue: "Task was destroyed but it is pending" ```python # WRONG async def bad(): asyncio.create_task(background_work()) # Orphaned! # CORRECT background_tasks = set() async def good(): task = asyncio.create_task(background_work()) background_tasks.add(task) task.add_done_callback(background_tasks.discard) ``` ### Issue: "Event loop is closed" ```python # WRONG - reusing closed loop loop = asyncio.get_event_loop() loop.run_until_complete(coro1()) loop.close() loop.run_until_complete(coro2()) # Error! # CORRECT - use asyncio.run() asyncio.run(coro1()) asyncio.run(coro2()) # New loop each time ``` ### Issue: "Cannot schedule new futures after shutdown" ```python # Happens when creating tasks during shutdown async def cleanup(): # DON'T create new tasks here await existing_task ``` ### Issue: Hung program (blocked event loop) ```python # Find the blocking call import asyncio async def debug_blocking(): loop = asyncio.get_running_loop() loop.slow_callback_duration = 0.001 # 1ms threshold # Enable debug mode loop.set_debug(True) # Your code here ``` ## Testing Async Code ```python import pytest import asyncio @pytest.mark.asyncio async def test_async_function(): result = await async_function() assert result == expected # Test timeouts @pytest.mark.asyncio async def test_with_timeout(): with pytest.raises(asyncio.TimeoutError): async with asyncio.timeout(0.1): await slow_function() # Mock async functions from unittest.mock import AsyncMock async def test_with_mock(): mock = AsyncMock(return_value="mocked") result = await mock() assert result == "mocked" ``` ## Visualization Tools ### aiomonitor ```python import aiomonitor async def main(): with aiomonitor.start_monitor(): # Connect via: nc localhost 50101 # or: python -m aiomonitor.cli await long_running_task() ``` ### aiodebug ```python from aiodebug import log_slow_callbacks log_slow_callbacks.enable(0.05) # Log callbacks > 50ms ``` ## Quick Debug Checklist 1. [ ] Enable debug mode: `asyncio.run(main(), debug=True)` 2. [ ] Check for unawaited coroutines: warnings -> errors 3. [ ] Look for blocking calls: time.sleep, requests, open() 4. [ ] Verify all tasks are awaited or tracked 5. [ ] Check for proper resource cleanup (sessions, connections) 6. [ ] Monitor task count: `len(asyncio.all_tasks())` 7. [ ] Profile with yappi for async-aware profiling -
error-handling.md 11.9 KB
# Async Error Handling Patterns Error handling patterns for resilient async applications. ## Retry with Exponential Backoff ```python import asyncio import random from typing import TypeVar, Callable, Awaitable T = TypeVar("T") async def retry_with_backoff( func: Callable[[], Awaitable[T]], max_retries: int = 3, base_delay: float = 1.0, max_delay: float = 60.0, exponential_base: float = 2.0, jitter: bool = True, retryable_exceptions: tuple = (Exception,), ) -> T: """ Retry async function with exponential backoff. Args: func: Async function to retry max_retries: Maximum number of retry attempts base_delay: Initial delay between retries (seconds) max_delay: Maximum delay between retries (seconds) exponential_base: Base for exponential calculation jitter: Add randomness to prevent thundering herd retryable_exceptions: Exceptions that trigger retry """ last_exception = None for attempt in range(max_retries + 1): try: return await func() except retryable_exceptions as e: last_exception = e if attempt == max_retries: raise # Calculate delay with exponential backoff delay = min( base_delay * (exponential_base ** attempt), max_delay ) # Add jitter (±25%) if jitter: delay *= 0.75 + random.random() * 0.5 await asyncio.sleep(delay) raise last_exception # Should never reach here # Usage async def fetch_with_retry(url: str) -> str: return await retry_with_backoff( lambda: fetch(url), max_retries=3, retryable_exceptions=(aiohttp.ClientError, asyncio.TimeoutError) ) ``` ## Retry Decorator ```python import functools from typing import Type def async_retry( max_retries: int = 3, base_delay: float = 1.0, exceptions: tuple[Type[Exception], ...] = (Exception,) ): """Decorator for async retry with backoff.""" def decorator(func): @functools.wraps(func) async def wrapper(*args, **kwargs): return await retry_with_backoff( lambda: func(*args, **kwargs), max_retries=max_retries, base_delay=base_delay, retryable_exceptions=exceptions ) return wrapper return decorator # Usage @async_retry(max_retries=3, exceptions=(aiohttp.ClientError,)) async def fetch_data(url: str) -> dict: async with aiohttp.ClientSession() as session: async with session.get(url) as response: return await response.json() ``` ## Circuit Breaker ```python import asyncio import time from enum import Enum from dataclasses import dataclass class CircuitState(Enum): CLOSED = "closed" # Normal operation OPEN = "open" # Failing, reject calls HALF_OPEN = "half_open" # Testing if recovered @dataclass class CircuitBreakerConfig: failure_threshold: int = 5 # Failures before opening success_threshold: int = 3 # Successes to close timeout: float = 60.0 # Seconds before half-open half_open_max_calls: int = 1 # Calls allowed in half-open class CircuitBreaker: """ Circuit breaker pattern for async operations. States: - CLOSED: Normal operation, tracking failures - OPEN: Rejecting calls, waiting for timeout - HALF_OPEN: Testing with limited calls """ def __init__(self, config: CircuitBreakerConfig | None = None): self.config = config or CircuitBreakerConfig() self._state = CircuitState.CLOSED self._failure_count = 0 self._success_count = 0 self._last_failure_time: float = 0 self._half_open_calls = 0 self._lock = asyncio.Lock() @property def state(self) -> CircuitState: return self._state async def call(self, func, *args, **kwargs): """Execute function through circuit breaker.""" async with self._lock: self._check_state_transition() if self._state == CircuitState.OPEN: raise CircuitBreakerOpen( f"Circuit open, retry after {self._retry_after():.1f}s" ) if self._state == CircuitState.HALF_OPEN: if self._half_open_calls >= self.config.half_open_max_calls: raise CircuitBreakerOpen("Half-open limit reached") self._half_open_calls += 1 try: result = await func(*args, **kwargs) await self._record_success() return result except Exception as e: await self._record_failure() raise def _check_state_transition(self): """Check if state should transition.""" if self._state == CircuitState.OPEN: if time.time() - self._last_failure_time >= self.config.timeout: self._state = CircuitState.HALF_OPEN self._half_open_calls = 0 self._success_count = 0 async def _record_success(self): async with self._lock: if self._state == CircuitState.HALF_OPEN: self._success_count += 1 if self._success_count >= self.config.success_threshold: self._state = CircuitState.CLOSED self._failure_count = 0 else: self._failure_count = 0 async def _record_failure(self): async with self._lock: self._failure_count += 1 self._last_failure_time = time.time() if self._state == CircuitState.HALF_OPEN: self._state = CircuitState.OPEN elif self._failure_count >= self.config.failure_threshold: self._state = CircuitState.OPEN def _retry_after(self) -> float: elapsed = time.time() - self._last_failure_time return max(0, self.config.timeout - elapsed) class CircuitBreakerOpen(Exception): """Raised when circuit breaker is open.""" pass # Usage breaker = CircuitBreaker(CircuitBreakerConfig( failure_threshold=5, timeout=30.0 )) async def fetch_with_breaker(url: str): return await breaker.call(fetch, url) ``` ## Partial Failure Handling ```python async def fetch_all_with_partial_failure( urls: list[str], max_failures: int | None = None ) -> tuple[list[str], list[Exception]]: """ Fetch all URLs, collecting both successes and failures. Args: urls: URLs to fetch max_failures: If set, abort after this many failures Returns: Tuple of (successful_results, exceptions) """ results = await asyncio.gather( *[fetch(url) for url in urls], return_exceptions=True ) successes = [] failures = [] for result in results: if isinstance(result, Exception): failures.append(result) if max_failures and len(failures) >= max_failures: # Cancel remaining work if too many failures break else: successes.append(result) return successes, failures # With structured handling @dataclass class FetchResult: url: str data: str | None = None error: Exception | None = None @property def success(self) -> bool: return self.error is None async def fetch_with_result(url: str) -> FetchResult: """Wrap fetch in result object.""" try: data = await fetch(url) return FetchResult(url=url, data=data) except Exception as e: return FetchResult(url=url, error=e) async def fetch_all_structured(urls: list[str]) -> list[FetchResult]: """Fetch all URLs with structured results.""" return await asyncio.gather(*[fetch_with_result(url) for url in urls]) ``` ## Exception Groups (Python 3.11+) ```python async def process_with_exception_groups(): """Handle multiple exceptions from TaskGroup.""" try: async with asyncio.TaskGroup() as tg: tg.create_task(task1()) tg.create_task(task2()) tg.create_task(task3()) except* ValueError as eg: # Handle all ValueError instances for exc in eg.exceptions: logger.error(f"ValueError: {exc}") except* TypeError as eg: # Handle all TypeError instances for exc in eg.exceptions: logger.error(f"TypeError: {exc}") # Filtering exception groups def handle_exception_group(eg: ExceptionGroup): """Process exception group by type.""" critical = [] recoverable = [] for exc in eg.exceptions: if isinstance(exc, (ConnectionError, TimeoutError)): recoverable.append(exc) else: critical.append(exc) # Retry recoverable errors for exc in recoverable: logger.warning(f"Recoverable error: {exc}") # Raise critical errors if critical: raise ExceptionGroup("Critical errors", critical) ``` ## Fallback Pattern ```python async def with_fallback( primary: Callable[[], Awaitable[T]], fallback: Callable[[], Awaitable[T]], exceptions: tuple = (Exception,) ) -> T: """Try primary, fall back on failure.""" try: return await primary() except exceptions as e: logger.warning(f"Primary failed, using fallback: {e}") return await fallback() # With multiple fallbacks async def with_fallback_chain( *funcs: Callable[[], Awaitable[T]] ) -> T: """Try functions in order until one succeeds.""" last_error = None for func in funcs: try: return await func() except Exception as e: last_error = e continue raise last_error or RuntimeError("No fallbacks provided") # Usage result = await with_fallback_chain( lambda: fetch_from_primary_api(), lambda: fetch_from_secondary_api(), lambda: fetch_from_cache(), ) ``` ## Bulkhead Pattern ```python class Bulkhead: """ Bulkhead pattern to isolate failures. Limits concurrent calls to protect resources. """ def __init__( self, max_concurrent: int, max_waiting: int = 0, timeout: float | None = None ): self._semaphore = asyncio.Semaphore(max_concurrent) self._max_waiting = max_waiting self._waiting = 0 self._timeout = timeout self._lock = asyncio.Lock() async def call(self, func, *args, **kwargs): """Execute function within bulkhead.""" async with self._lock: if self._waiting >= self._max_waiting: raise BulkheadFull("Bulkhead queue full") self._waiting += 1 try: if self._timeout: async with asyncio.timeout(self._timeout): async with self._semaphore: return await func(*args, **kwargs) else: async with self._semaphore: return await func(*args, **kwargs) finally: async with self._lock: self._waiting -= 1 class BulkheadFull(Exception): """Raised when bulkhead cannot accept more calls.""" pass # Usage - isolate external service calls external_api_bulkhead = Bulkhead( max_concurrent=10, # Max 10 concurrent calls max_waiting=50, # Max 50 in queue timeout=30.0 # 30s timeout ) async def call_external_api(data): return await external_api_bulkhead.call( lambda: http_client.post("/api", json=data) ) ``` ## Quick Reference | Pattern | Use Case | Behavior | |---------|----------|----------| | Retry + backoff | Transient failures | Retry with increasing delays | | Circuit breaker | Cascading failures | Fast-fail when service down | | Fallback | Degraded operation | Use backup on failure | | Bulkhead | Resource isolation | Limit concurrent access | | Exception groups | Multiple failures | Handle 3.11+ TaskGroup errors | | Partial failure | Best-effort batch | Collect successes and failures | -
mixing-sync-async.md 5.7 KB
# Mixing Sync and Async Patterns for bridging synchronous and asynchronous Python code. ## Running Sync Code from Async ### run_in_executor ```python import asyncio from concurrent.futures import ThreadPoolExecutor async def run_blocking(): """Run blocking I/O in thread pool.""" loop = asyncio.get_running_loop() # Using default executor (ThreadPoolExecutor) result = await loop.run_in_executor( None, # Default executor blocking_function, arg1, arg2 ) return result # With custom executor executor = ThreadPoolExecutor(max_workers=4) async def run_with_custom_executor(): loop = asyncio.get_running_loop() result = await loop.run_in_executor( executor, blocking_function, arg1 ) return result ``` ### CPU-bound with ProcessPoolExecutor ```python from concurrent.futures import ProcessPoolExecutor executor = ProcessPoolExecutor(max_workers=4) async def run_cpu_bound(): """Run CPU-bound code in process pool.""" loop = asyncio.get_running_loop() result = await loop.run_in_executor( executor, cpu_intensive_function, data ) return result ``` ### Decorator Pattern ```python import asyncio import functools def run_in_executor(func): """Decorator to run sync function in executor.""" @functools.wraps(func) async def wrapper(*args, **kwargs): loop = asyncio.get_running_loop() return await loop.run_in_executor( None, functools.partial(func, *args, **kwargs) ) return wrapper @run_in_executor def blocking_io_operation(path): with open(path) as f: return f.read() # Usage async def main(): content = await blocking_io_operation("file.txt") ``` ## Running Async Code from Sync ### asyncio.run() ```python import asyncio async def async_function(): await asyncio.sleep(1) return "done" # From sync code def sync_wrapper(): return asyncio.run(async_function()) ``` ### Nested Event Loops (nest_asyncio) ```python # For Jupyter notebooks or nested contexts import nest_asyncio nest_asyncio.apply() # Now asyncio.run() works even if event loop is running ``` ### Thread with Event Loop ```python import asyncio import threading def run_in_new_thread(coro): """Run coroutine in a new thread with its own event loop.""" result = None exception = None def runner(): nonlocal result, exception try: result = asyncio.run(coro) except Exception as e: exception = e thread = threading.Thread(target=runner) thread.start() thread.join() if exception: raise exception return result ``` ## Common Pitfalls ### DON'T: Call asyncio.run() from async ```python # WRONG - nested asyncio.run() async def bad(): result = asyncio.run(other_async()) # RuntimeError! # CORRECT - just await async def good(): result = await other_async() ``` ### DON'T: Use time.sleep() in async ```python # WRONG - blocks event loop async def bad(): time.sleep(5) # Blocks entire event loop! # CORRECT async def good(): await asyncio.sleep(5) ``` ### DON'T: Use blocking I/O directly ```python # WRONG - blocks event loop async def bad(): with open("file.txt") as f: # Blocking! return f.read() # CORRECT - use executor async def good(): loop = asyncio.get_running_loop() return await loop.run_in_executor(None, read_file, "file.txt") # OR use async file library import aiofiles async def better(): async with aiofiles.open("file.txt") as f: return await f.read() ``` ## Synchronization Primitives ### Threading Lock vs asyncio Lock ```python import threading import asyncio # For sync code sync_lock = threading.Lock() # For async code async_lock = asyncio.Lock() # DON'T mix them! # threading.Lock() in async code blocks event loop # asyncio.Lock() in sync code doesn't work ``` ### Thread-Safe Queue for Sync/Async Bridge ```python import asyncio import queue import threading def sync_producer(q: queue.Queue): """Sync code putting items.""" for i in range(10): q.put(i) q.put(None) # Sentinel async def async_consumer(q: queue.Queue): """Async code getting items from sync queue.""" loop = asyncio.get_running_loop() while True: # Non-blocking get in executor item = await loop.run_in_executor(None, q.get) if item is None: break await process(item) async def main(): q = queue.Queue() # Start sync producer in thread thread = threading.Thread(target=sync_producer, args=(q,)) thread.start() # Consume async await async_consumer(q) thread.join() ``` ## Async-First Database Access ```python # Instead of sync database drivers, use async versions # SQLite import aiosqlite async def query_db(): async with aiosqlite.connect("db.sqlite") as db: async with db.execute("SELECT * FROM users") as cursor: return await cursor.fetchall() # PostgreSQL import asyncpg async def query_postgres(): conn = await asyncpg.connect("postgresql://...") rows = await conn.fetch("SELECT * FROM users") await conn.close() return rows # HTTP import aiohttp async def fetch_api(): async with aiohttp.ClientSession() as session: async with session.get(url) as response: return await response.json() ``` ## Best Practices 1. **Prefer async libraries** - Use aiohttp, aiosqlite, asyncpg over sync versions 2. **Use run_in_executor for blocking** - Never block the event loop 3. **Keep sync/async boundaries clean** - Don't mix unnecessarily 4. **Use ProcessPoolExecutor for CPU-bound** - ThreadPool for I/O 5. **Don't nest event loops** - Use a single asyncio.run() entry point 6. **Profile before threading** - Async is often enough -
performance.md 10.1 KB
# Async Performance Optimization Performance patterns for high-throughput async applications. ## uvloop - Drop-in Event Loop Replacement ```python # Install: uv add uvloop # Option 1: Install as default (before any asyncio calls) import uvloop uvloop.install() # Then use asyncio normally import asyncio async def main(): # Now using uvloop pass asyncio.run(main()) # Option 2: Use explicitly import asyncio import uvloop async def main(): pass with asyncio.Runner(loop_factory=uvloop.new_event_loop) as runner: runner.run(main()) # Option 3: Check if available def get_event_loop_policy(): try: import uvloop return uvloop.EventLoopPolicy() except ImportError: return asyncio.DefaultEventLoopPolicy() asyncio.set_event_loop_policy(get_event_loop_policy()) ``` **Performance gains:** - 2-4x faster than default asyncio event loop - Significant improvement for I/O-bound workloads - Based on libuv (same as Node.js) ## Connection Pool Tuning ```python import aiohttp # Optimal connector settings connector = aiohttp.TCPConnector( limit=100, # Total connection limit limit_per_host=30, # Per-host limit (prevents overwhelming one server) ttl_dns_cache=300, # DNS cache TTL (seconds) use_dns_cache=True, # Enable DNS caching keepalive_timeout=30, # Keep connections alive (seconds) enable_cleanup_closed=True, # Clean up closed connections ) async with aiohttp.ClientSession(connector=connector) as session: # Use session pass # Database connection pool (asyncpg) import asyncpg pool = await asyncpg.create_pool( dsn="postgresql://user:pass@localhost/db", min_size=5, # Minimum connections to keep max_size=20, # Maximum connections allowed max_inactive_connection_lifetime=300.0, # Close idle connections command_timeout=60.0, # Query timeout ) # Redis connection pool (aioredis/redis-py) import redis.asyncio as redis pool = redis.ConnectionPool.from_url( "redis://localhost", max_connections=50, decode_responses=True, ) client = redis.Redis(connection_pool=pool) ``` ### Pool Sizing Guidelines | Service Type | Min Size | Max Size | Notes | |--------------|----------|----------|-------| | Database (heavy) | 10 | 50 | Match CPU cores × 2-4 | | Database (light) | 5 | 20 | Standard web apps | | HTTP external API | N/A | 100 | Limited by rate limits | | HTTP per-host | N/A | 30 | Prevent overwhelming | | Redis | 10 | 50 | Very fast, less critical | ## Buffer Sizing ```python # aiohttp response reading async def fetch_large(session, url): async with session.get(url) as response: # Default: reads entire response into memory data = await response.read() # For large responses, stream: chunks = [] async for chunk in response.content.iter_chunked(8192): chunks.append(chunk) # Custom buffer sizes for TCP import asyncio async def create_connection(): reader, writer = await asyncio.open_connection( "localhost", 8888, limit=2**20, # 1MB read buffer (default is 64KB) ) return reader, writer # aiohttp server with custom limits from aiohttp import web app = web.Application( client_max_size=1024 * 1024 * 100, # 100MB max request body ) ``` ## Batching Requests ```python import asyncio from collections import defaultdict from typing import TypeVar, Callable T = TypeVar("T") class BatchProcessor: """Batch multiple requests into single operations.""" def __init__( self, batch_func: Callable[[list[str]], dict[str, T]], max_batch_size: int = 100, max_delay: float = 0.01 # 10ms ): self._batch_func = batch_func self._max_batch_size = max_batch_size self._max_delay = max_delay self._pending: dict[str, asyncio.Future] = {} self._batch: list[str] = [] self._lock = asyncio.Lock() self._timer: asyncio.Task | None = None async def get(self, key: str) -> T: """Get single item (batched with other requests).""" async with self._lock: if key in self._pending: return await self._pending[key] future = asyncio.get_event_loop().create_future() self._pending[key] = future self._batch.append(key) if len(self._batch) >= self._max_batch_size: await self._flush() elif not self._timer: self._timer = asyncio.create_task(self._delayed_flush()) return await future async def _delayed_flush(self): await asyncio.sleep(self._max_delay) async with self._lock: await self._flush() async def _flush(self): if not self._batch: return batch = self._batch pending = self._pending self._batch = [] self._pending = {} self._timer = None try: results = await self._batch_func(batch) for key in batch: if key in results: pending[key].set_result(results[key]) else: pending[key].set_exception(KeyError(key)) except Exception as e: for key in batch: pending[key].set_exception(e) # Usage async def batch_fetch_users(user_ids: list[str]) -> dict[str, User]: # Single database query for multiple users return {u.id: u for u in await db.fetch_users(user_ids)} user_batcher = BatchProcessor(batch_fetch_users, max_batch_size=50) # These will be batched together: user1 = await user_batcher.get("user-1") user2 = await user_batcher.get("user-2") ``` ## Task Prioritization ```python import asyncio import heapq from dataclasses import dataclass, field from typing import Any @dataclass(order=True) class PrioritizedTask: priority: int item: Any = field(compare=False) class PriorityQueue: """Async priority queue for task ordering.""" def __init__(self): self._queue: list[PrioritizedTask] = [] self._condition = asyncio.Condition() async def put(self, priority: int, item: Any): async with self._condition: heapq.heappush(self._queue, PrioritizedTask(priority, item)) self._condition.notify() async def get(self) -> Any: async with self._condition: while not self._queue: await self._condition.wait() return heapq.heappop(self._queue).item # Usage queue = PriorityQueue() # Lower number = higher priority await queue.put(1, "critical task") await queue.put(10, "low priority task") await queue.put(5, "normal task") ``` ## Memory Optimization ```python import asyncio from weakref import WeakValueDictionary # Use weak references for caches class AsyncCache: """Memory-efficient async cache using weak references.""" def __init__(self, fetch_func): self._cache = WeakValueDictionary() self._fetch_func = fetch_func self._locks: dict[str, asyncio.Lock] = {} async def get(self, key: str): if key in self._cache: return self._cache[key] if key not in self._locks: self._locks[key] = asyncio.Lock() async with self._locks[key]: if key in self._cache: return self._cache[key] value = await self._fetch_func(key) self._cache[key] = value return value # Limit concurrent operations to prevent memory spikes async def process_large_dataset(items: list, concurrency: int = 10): """Process items with limited concurrency.""" semaphore = asyncio.Semaphore(concurrency) async def process_one(item): async with semaphore: result = await heavy_processing(item) return result # Process in chunks to avoid memory issues with huge lists chunk_size = 1000 all_results = [] for i in range(0, len(items), chunk_size): chunk = items[i:i + chunk_size] results = await asyncio.gather(*[process_one(item) for item in chunk]) all_results.extend(results) return all_results ``` ## Profiling Async Code ```python import asyncio import time from contextlib import asynccontextmanager @asynccontextmanager async def async_timer(name: str): """Context manager to time async operations.""" start = time.perf_counter() try: yield finally: elapsed = time.perf_counter() - start print(f"{name}: {elapsed:.3f}s") # Usage async with async_timer("fetch_all"): results = await fetch_all(urls) # Detailed profiling class AsyncProfiler: def __init__(self): self.timings: dict[str, list[float]] = {} @asynccontextmanager async def profile(self, name: str): start = time.perf_counter() try: yield finally: elapsed = time.perf_counter() - start if name not in self.timings: self.timings[name] = [] self.timings[name].append(elapsed) def report(self): for name, times in self.timings.items(): avg = sum(times) / len(times) total = sum(times) print(f"{name}: avg={avg:.3f}s, total={total:.3f}s, count={len(times)}") # Use yappi for comprehensive profiling # uv add --dev yappi import yappi yappi.set_clock_type("wall") # For async code yappi.start() asyncio.run(main()) yappi.stop() yappi.get_func_stats().print_all() ``` ## Quick Reference | Optimization | Impact | When to Use | |--------------|--------|-------------| | uvloop | 2-4x throughput | Always (production) | | Connection pooling | Reduce latency | Any external service | | Request batching | N requests → 1 | Database, APIs | | Semaphore limiting | Memory control | Large datasets | | Streaming | Memory efficiency | Large responses | | Priority queue | Latency SLAs | Mixed workloads | ## Performance Checklist ```markdown - [ ] uvloop installed and configured - [ ] Connection pools properly sized - [ ] Timeouts on all external calls - [ ] Semaphores limiting concurrency - [ ] Large responses streamed - [ ] DNS caching enabled - [ ] Connection keep-alive configured - [ ] Profiling in place for hot paths ``` -
production-patterns.md 10 KB
# Production Async Patterns Production-ready patterns for deploying async Python applications. ## Graceful Shutdown ```python import asyncio import signal from contextlib import asynccontextmanager class GracefulShutdown: """Handle graceful shutdown with signal handlers.""" def __init__(self): self._shutdown = asyncio.Event() self._tasks: set[asyncio.Task] = set() @property def should_exit(self) -> bool: return self._shutdown.is_set() async def wait_for_shutdown(self): """Block until shutdown signal received.""" await self._shutdown.wait() def trigger_shutdown(self): """Signal shutdown to all waiting coroutines.""" self._shutdown.set() def register_task(self, task: asyncio.Task): """Track task for cleanup on shutdown.""" self._tasks.add(task) task.add_done_callback(self._tasks.discard) async def cleanup(self, timeout: float = 30.0): """Cancel and await all tracked tasks.""" for task in self._tasks: task.cancel() if self._tasks: await asyncio.wait( self._tasks, timeout=timeout, return_when=asyncio.ALL_COMPLETED ) async def main(): shutdown = GracefulShutdown() loop = asyncio.get_running_loop() # Register signal handlers for sig in (signal.SIGTERM, signal.SIGINT): loop.add_signal_handler(sig, shutdown.trigger_shutdown) try: # Start background services worker = asyncio.create_task(background_worker(shutdown)) shutdown.register_task(worker) # Run until shutdown await shutdown.wait_for_shutdown() finally: # Cleanup await shutdown.cleanup(timeout=30.0) # Remove signal handlers for sig in (signal.SIGTERM, signal.SIGINT): loop.remove_signal_handler(sig) async def background_worker(shutdown: GracefulShutdown): """Worker that respects shutdown signals.""" while not shutdown.should_exit: try: await process_next_item() except asyncio.CancelledError: # Finish current work before exiting await finish_current_work() raise ``` ## Lifespan Context Manager ```python from contextlib import asynccontextmanager @asynccontextmanager async def lifespan(): """Application lifespan manager for startup/shutdown.""" # Startup db_pool = await create_db_pool() redis = await create_redis_client() try: yield {"db": db_pool, "redis": redis} finally: # Shutdown (always runs) await redis.close() await db_pool.close() # Usage with FastAPI from fastapi import FastAPI @asynccontextmanager async def lifespan(app: FastAPI): # Startup app.state.db = await create_db_pool() yield # Shutdown await app.state.db.close() app = FastAPI(lifespan=lifespan) ``` ## Health Check Endpoints ```python import asyncio from dataclasses import dataclass from enum import Enum class HealthStatus(str, Enum): HEALTHY = "healthy" DEGRADED = "degraded" UNHEALTHY = "unhealthy" @dataclass class ComponentHealth: name: str status: HealthStatus latency_ms: float | None = None error: str | None = None async def check_database(pool) -> ComponentHealth: """Check database connectivity.""" try: start = asyncio.get_event_loop().time() async with pool.acquire() as conn: await conn.execute("SELECT 1") latency = (asyncio.get_event_loop().time() - start) * 1000 return ComponentHealth("database", HealthStatus.HEALTHY, latency) except Exception as e: return ComponentHealth("database", HealthStatus.UNHEALTHY, error=str(e)) async def check_redis(client) -> ComponentHealth: """Check Redis connectivity.""" try: start = asyncio.get_event_loop().time() await client.ping() latency = (asyncio.get_event_loop().time() - start) * 1000 return ComponentHealth("redis", HealthStatus.HEALTHY, latency) except Exception as e: return ComponentHealth("redis", HealthStatus.UNHEALTHY, error=str(e)) async def health_check(pool, redis) -> dict: """Aggregate health check for all components.""" checks = await asyncio.gather( check_database(pool), check_redis(redis), return_exceptions=True ) components = [] overall = HealthStatus.HEALTHY for check in checks: if isinstance(check, Exception): components.append(ComponentHealth("unknown", HealthStatus.UNHEALTHY, error=str(check))) overall = HealthStatus.UNHEALTHY else: components.append(check) if check.status == HealthStatus.UNHEALTHY: overall = HealthStatus.UNHEALTHY elif check.status == HealthStatus.DEGRADED and overall == HealthStatus.HEALTHY: overall = HealthStatus.DEGRADED return { "status": overall.value, "components": [ {"name": c.name, "status": c.status.value, "latency_ms": c.latency_ms, "error": c.error} for c in components ] } ``` ## Liveness vs Readiness Probes ```python class HealthProbes: """Kubernetes-style health probes.""" def __init__(self): self._ready = asyncio.Event() self._alive = True def set_ready(self): """Mark application as ready to receive traffic.""" self._ready.set() def set_not_ready(self): """Mark application as not ready (drain traffic).""" self._ready.clear() def set_not_alive(self): """Mark application as dead (trigger restart).""" self._alive = False async def liveness(self) -> bool: """ Liveness probe - is the process healthy? Failing this triggers a container restart. """ return self._alive async def readiness(self) -> bool: """ Readiness probe - can the app handle traffic? Failing this removes the pod from service. """ return self._ready.is_set() async def startup(self) -> bool: """ Startup probe - has the app finished initializing? Prevents liveness checks during slow startup. """ return self._ready.is_set() # Usage probes = HealthProbes() async def startup(): await initialize_db() await warm_caches() probes.set_ready() # Now accept traffic async def shutdown(): probes.set_not_ready() # Stop accepting new requests await drain_connections() # Finish in-flight requests ``` ## Resource Cleanup on Cancellation ```python async def process_with_cleanup(): """Ensure cleanup even when cancelled.""" resource = await acquire_resource() try: await do_work(resource) except asyncio.CancelledError: # Perform essential cleanup before re-raising await resource.flush() raise finally: # Always close resource await resource.close() async def shielded_cleanup(): """Protect critical cleanup from cancellation.""" resource = await acquire_resource() try: await do_work(resource) finally: # Shield cleanup from cancellation await asyncio.shield(resource.close()) ``` ## Background Task Management ```python class BackgroundTaskManager: """Manage long-running background tasks.""" def __init__(self): self._tasks: dict[str, asyncio.Task] = {} self._shutdown = asyncio.Event() def start(self, name: str, coro): """Start a named background task.""" if name in self._tasks: raise ValueError(f"Task {name} already running") task = asyncio.create_task(coro, name=name) task.add_done_callback(lambda t: self._task_done(name, t)) self._tasks[name] = task return task def _task_done(self, name: str, task: asyncio.Task): """Handle task completion.""" self._tasks.pop(name, None) if not task.cancelled(): exc = task.exception() if exc: # Log error, potentially restart logger.error(f"Task {name} failed: {exc}") async def stop(self, name: str, timeout: float = 10.0): """Stop a specific task.""" if task := self._tasks.get(name): task.cancel() try: await asyncio.wait_for(task, timeout=timeout) except (asyncio.CancelledError, asyncio.TimeoutError): pass async def shutdown(self, timeout: float = 30.0): """Stop all background tasks.""" self._shutdown.set() for task in self._tasks.values(): task.cancel() if self._tasks: await asyncio.wait( self._tasks.values(), timeout=timeout, return_when=asyncio.ALL_COMPLETED ) ``` ## Periodic Tasks ```python async def periodic_task( interval: float, coro_func, shutdown_event: asyncio.Event | None = None ): """Run a coroutine periodically.""" while True: if shutdown_event and shutdown_event.is_set(): break try: await coro_func() except asyncio.CancelledError: raise except Exception as e: logger.error(f"Periodic task error: {e}") # Wait for interval or shutdown if shutdown_event: try: await asyncio.wait_for( shutdown_event.wait(), timeout=interval ) break # Shutdown signaled except asyncio.TimeoutError: pass # Continue loop else: await asyncio.sleep(interval) ``` ## Quick Reference | Pattern | Use Case | |---------|----------| | `GracefulShutdown` | SIGTERM/SIGINT handling | | `lifespan` context | Startup/shutdown resources | | `HealthProbes` | Kubernetes health checks | | `asyncio.shield()` | Protect critical cleanup | | `BackgroundTaskManager` | Long-running task lifecycle | | `periodic_task` | Scheduled background work |
-
-
scripts
-
find-blocking-calls.sh 1.4 KB
#!/bin/bash # Find potentially blocking calls in async Python code # Usage: ./find-blocking-calls.sh [directory] DIR="${1:-.}" echo "=== Scanning for blocking calls in async code ===" echo "Directory: $DIR" echo # time.sleep() in async functions echo "--- time.sleep() in async functions ---" rg -n "async def" -A 20 "$DIR" | rg "time\.sleep\(" || echo "None found" echo # requests library (blocking HTTP) echo "--- requests library usage ---" rg -n "import requests|from requests" "$DIR" --type py || echo "None found" echo # Blocking file operations echo "--- Blocking file operations in async ---" rg -n "async def" -A 30 "$DIR" | rg "open\(|\.read\(\)|\.write\(" || echo "None found" echo # subprocess without asyncio echo "--- Blocking subprocess calls ---" rg -n "subprocess\.(run|call|check_output)" "$DIR" --type py || echo "None found" echo # socket operations echo "--- Blocking socket operations ---" rg -n "socket\.(socket|create_connection)" "$DIR" --type py || echo "None found" echo # input() calls echo "--- Blocking input() calls ---" rg -n "\binput\(" "$DIR" --type py || echo "None found" echo echo "=== Recommendations ===" echo "- Replace time.sleep() with asyncio.sleep()" echo "- Replace requests with aiohttp or httpx" echo "- Replace open() with aiofiles" echo "- Replace subprocess with asyncio.create_subprocess_exec()" echo "- Use asyncio.get_event_loop().run_in_executor() for blocking code"
-
-
SKILL.md 4.8 KB
--- name: python-async-ops description: "Python asyncio patterns for concurrent programming. Triggers on: asyncio, async, await, coroutine, gather, semaphore, TaskGroup, event loop, aiohttp, concurrent." license: MIT compatibility: "Python 3.10+ recommended. Some patterns require 3.11+ (TaskGroup, timeout)." allowed-tools: "Read Write" metadata: author: claude-mods depends-on: python-typing-ops related-skills: python-fastapi-ops, python-observability-ops --- # Python Async Patterns Asyncio patterns for concurrent Python programming. ## When to Use Async vs Sync | Use Async When | Use Sync When | |----------------|---------------| | I/O-bound operations (HTTP, DB, files) | CPU-bound computations | | High concurrency (100s+ connections) | Simple scripts, one-off tasks | | WebSocket/streaming connections | Small data processing | | Microservices with network calls | Single sequential operations | **Decision tree:** 1. Is it CPU-bound? → Sync (or multiprocessing) 2. Is it I/O-bound with high concurrency? → Async 3. Is it simple I/O with few connections? → Sync is fine ## Core Concepts ```python import asyncio # Coroutine (must be awaited) async def fetch(url: str) -> str: async with aiohttp.ClientSession() as session: async with session.get(url) as response: return await response.text() # Entry point async def main(): result = await fetch("https://example.com") return result asyncio.run(main()) ``` ## Pattern 1: Concurrent with gather ```python async def fetch_all(urls: list[str]) -> list[str]: """Fetch multiple URLs concurrently.""" async with aiohttp.ClientSession() as session: tasks = [fetch_one(session, url) for url in urls] return await asyncio.gather(*tasks, return_exceptions=True) ``` ## Pattern 2: Bounded Concurrency ```python async def fetch_with_limit(urls: list[str], limit: int = 10): """Limit concurrent requests.""" semaphore = asyncio.Semaphore(limit) async def bounded_fetch(url): async with semaphore: return await fetch_one(url) return await asyncio.gather(*[bounded_fetch(url) for url in urls]) ``` ## Pattern 3: TaskGroup (Python 3.11+) ```python async def process_items(items): """Structured concurrency with automatic cleanup.""" async with asyncio.TaskGroup() as tg: for item in items: tg.create_task(process_one(item)) # All tasks complete here, or exception raised ``` ## Pattern 4: Timeout ```python async def with_timeout(): try: async with asyncio.timeout(5.0): # Python 3.11+ result = await slow_operation() except asyncio.TimeoutError: result = None return result ``` ## Critical Warnings ```python # WRONG - blocks event loop async def bad(): time.sleep(5) # Never use time.sleep! requests.get(url) # Blocking I/O! # CORRECT async def good(): await asyncio.sleep(5) async with aiohttp.ClientSession() as s: await s.get(url) ``` ```python # WRONG - orphaned task async def bad(): asyncio.create_task(work()) # May be garbage collected! # CORRECT - keep reference async def good(): task = asyncio.create_task(work()) await task ``` ## Quick Reference | Pattern | Use Case | |---------|----------| | `gather(*tasks)` | Multiple independent operations | | `Semaphore(n)` | Rate limiting, resource constraints | | `TaskGroup()` | Structured concurrency (3.11+) | | `Queue()` | Producer-consumer | | `timeout(s)` | Timeout wrapper (3.11+) | | `Lock()` | Shared mutable state | ## Async Context Manager ```python from contextlib import asynccontextmanager @asynccontextmanager async def managed_connection(): conn = await create_connection() try: yield conn finally: await conn.close() ``` ## Additional Resources For detailed patterns, load: - `./references/concurrency-patterns.md` - Queue, Lock, producer-consumer - `./references/aiohttp-patterns.md` - HTTP client/server patterns - `./references/mixing-sync-async.md` - run_in_executor, thread pools - `./references/debugging-async.md` - Debug mode, profiling, finding issues - `./references/production-patterns.md` - Graceful shutdown, health checks, signal handling - `./references/error-handling.md` - Retry with backoff, circuit breakers, partial failures - `./references/performance.md` - uvloop, connection pooling, buffer sizing ## Scripts - `./scripts/find-blocking-calls.sh` - Scan code for blocking calls in async functions ## Assets - `./assets/async-project-template.py` - Production-ready async app skeleton --- ## See Also **Prerequisites:** - `python-typing-ops` - Type hints for async functions **Related Skills:** - `python-fastapi-ops` - Async web APIs - `python-observability-ops` - Async logging and tracing - `python-database-ops` - Async database access
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.