HardТеория8 min

Практические паттерны async

HTTP-клиенты, rate limiting, retry, graceful shutdown и продвинутые паттерны asyncio

Асинхронные HTTP-запросы с httpx

Библиотека httpx — современная замена requests с полной поддержкой async/await. Она предоставляет AsyncClient для эффективной работы с HTTP в асинхронном коде.

import asyncio
import httpx

async def fetch_users() -> list[dict]:
    """Fetch users from JSON API."""
    async with httpx.AsyncClient(timeout=10.0) as client:
        response = await client.get("https://jsonplaceholder.typicode.com/users")
        response.raise_for_status()
        return response.json()

async def fetch_multiple_endpoints() -> dict:
    """Fetch data from multiple endpoints concurrently."""
    async with httpx.AsyncClient(
        base_url="https://jsonplaceholder.typicode.com",
        timeout=10.0,
    ) as client:
        # All requests run concurrently
        users_task = asyncio.create_task(client.get("/users"))
        posts_task = asyncio.create_task(client.get("/posts"))
        albums_task = asyncio.create_task(client.get("/albums"))

        users_resp = await users_task
        posts_resp = await posts_task
        albums_resp = await albums_task

        return {
            "users": users_resp.json()[:5],
            "posts": posts_resp.json()[:5],
            "albums": albums_resp.json()[:5],
        }

async def main() -> None:
    data = await fetch_multiple_endpoints()
    for key, items in data.items():
        print(f"{key}: {len(items)} записей")

asyncio.run(main())

Важно: всегда используйте async with httpx.AsyncClient() — это гарантирует корректное закрытие соединений и переиспользование пула подключений.

Rate Limiting — ограничение частоты запросов

Большинство API имеют лимиты на количество запросов. Паттерн token bucket позволяет контролировать частоту:

import asyncio
import time

class RateLimiter:
    """Token bucket rate limiter for async operations."""

    def __init__(self, rate: float, burst: int = 1) -> None:
        self.rate = rate          # Requests per second
        self.burst = burst        # Max burst size
        self._tokens = float(burst)
        self._last_refill = time.monotonic()
        self._lock = asyncio.Lock()

    async def acquire(self) -> None:
        """Wait until a token is available."""
        async with self._lock:
            while True:
                now = time.monotonic()
                elapsed = now - self._last_refill
                self._tokens = min(
                    self.burst,
                    self._tokens + elapsed * self.rate,
                )
                self._last_refill = now

                if self._tokens >= 1.0:
                    self._tokens -= 1.0
                    return

                # Calculate wait time for next token
                wait = (1.0 - self._tokens) / self.rate
                await asyncio.sleep(wait)

    async def __aenter__(self):
        await self.acquire()
        return self

    async def __aexit__(self, *args):
        pass

async def fetch_with_rate_limit(
    url: str,
    limiter: RateLimiter,
) -> str:
    """Fetch URL respecting rate limits."""
    async with limiter:
        # Simulate API call
        await asyncio.sleep(0.1)
        return f"OK: {url}"

async def main() -> None:
    # Allow 5 requests per second with burst of 3
    limiter = RateLimiter(rate=5.0, burst=3)

    urls = [f"https://api.example.com/item/{i}" for i in range(15)]

    start = time.monotonic()
    results = await asyncio.gather(
        *[fetch_with_rate_limit(url, limiter) for url in urls]
    )
    elapsed = time.monotonic() - start

    print(f"Выполнено {len(results)} запросов за {elapsed:.2f}с")
    print(f"Средняя скорость: {len(results) / elapsed:.1f} req/s")

asyncio.run(main())

Retry с экспоненциальной задержкой

Сетевые запросы могут временно падать. Паттерн retry с backoff повторяет операцию с увеличивающимися интервалами:

import asyncio
import random
from collections.abc import Awaitable, Callable
from functools import wraps
from typing import TypeVar

T = TypeVar("T")

def async_retry(
    max_attempts: int = 3,
    base_delay: float = 1.0,
    max_delay: float = 60.0,
    exceptions: tuple[type[Exception], ...] = (Exception,),
):
    """Decorator for async retry with exponential backoff."""
    def decorator(func: Callable[..., Awaitable[T]]) -> Callable[..., Awaitable[T]]:
        @wraps(func)
        async def wrapper(*args, **kwargs) -> T:
            last_exception: Exception | None = None

            for attempt in range(1, max_attempts + 1):
                try:
                    return await func(*args, **kwargs)
                except exceptions as e:
                    last_exception = e
                    if attempt == max_attempts:
                        break

                    # Exponential backoff with jitter
                    delay = min(
                        base_delay * (2 ** (attempt - 1)),
                        max_delay,
                    )
                    jitter = random.uniform(0, delay * 0.1)
                    wait = delay + jitter

                    print(
                        f"Попытка {attempt}/{max_attempts} не удалась: {e}. "
                        f"Повтор через {wait:.1f}с"
                    )
                    await asyncio.sleep(wait)

            raise last_exception  # type: ignore[misc]

        return wrapper
    return decorator

@async_retry(max_attempts=3, base_delay=0.5, exceptions=(ConnectionError, TimeoutError))
async def fetch_data(url: str) -> dict:
    """Fetch data with automatic retry on failure."""
    # Simulate flaky API
    if random.random() < 0.6:
        raise ConnectionError(f"Не удалось подключиться к {url}")
    await asyncio.sleep(0.2)
    return {"url": url, "status": "ok"}

async def main() -> None:
    try:
        result = await fetch_data("https://api.example.com/data")
        print(f"Успех: {result}")
    except ConnectionError:
        print("Все попытки исчерпаны")

asyncio.run(main())

Circuit Breaker — предохранитель

Если сервис стабильно недоступен, нет смысла продолжать запросы. Circuit breaker «размыкает цепь» после серии ошибок:

import asyncio
import time
import random
from enum import Enum

class CircuitState(Enum):
    CLOSED = "closed"        # Normal operation
    OPEN = "open"            # Failing, reject requests
    HALF_OPEN = "half_open"  # Testing if service recovered

class CircuitBreaker:
    """Async circuit breaker pattern."""

    def __init__(
        self,
        failure_threshold: int = 5,
        recovery_timeout: float = 30.0,
    ) -> None:
        self.failure_threshold = failure_threshold
        self.recovery_timeout = recovery_timeout
        self._state = CircuitState.CLOSED
        self._failure_count = 0
        self._last_failure_time = 0.0

    @property
    def state(self) -> CircuitState:
        if self._state == CircuitState.OPEN:
            elapsed = time.monotonic() - self._last_failure_time
            if elapsed >= self.recovery_timeout:
                self._state = CircuitState.HALF_OPEN
        return self._state

    async def call(self, coro):
        """Execute coroutine through circuit breaker."""
        if self.state == CircuitState.OPEN:
            raise RuntimeError(
                f"Circuit breaker OPEN — сервис недоступен"
            )

        try:
            result = await coro
            self._on_success()
            return result
        except Exception:
            self._on_failure()
            raise

    def _on_success(self) -> None:
        self._failure_count = 0
        self._state = CircuitState.CLOSED

    def _on_failure(self) -> None:
        self._failure_count += 1
        self._last_failure_time = time.monotonic()
        if self._failure_count >= self.failure_threshold:
            self._state = CircuitState.OPEN
            print(f"Circuit OPEN после {self._failure_count} ошибок")

async def unreliable_service() -> str:
    """Simulate unreliable service."""
    if random.random() < 0.7:
        raise ConnectionError("Сервис недоступен")
    return "Успех!"

async def main() -> None:
    breaker = CircuitBreaker(failure_threshold=3, recovery_timeout=5.0)

    for i in range(10):
        try:
            result = await breaker.call(unreliable_service())
            print(f"[{i}] {result}")
        except (ConnectionError, RuntimeError) as e:
            print(f"[{i}] Ошибка: {e}")
        await asyncio.sleep(0.5)

asyncio.run(main())

Graceful Shutdown — корректное завершение

Асинхронное приложение должно корректно завершаться: дождаться текущих задач, закрыть соединения, сохранить состояние.

import asyncio
import signal

class GracefulApp:
    """Application with graceful shutdown support."""

    def __init__(self) -> None:
        self._shutdown_event = asyncio.Event()
        self._tasks: set[asyncio.Task] = set()

    async def worker(self, name: str) -> None:
        """Long-running worker."""
        try:
            while not self._shutdown_event.is_set():
                print(f"[{name}] Обрабатываю...")
                try:
                    await asyncio.wait_for(
                        self._shutdown_event.wait(),
                        timeout=2.0,
                    )
                except TimeoutError:
                    pass  # Continue working
        except asyncio.CancelledError:
            pass
        finally:
            print(f"[{name}] Завершён")

    def _signal_handler(self) -> None:
        """Handle OS signals for graceful shutdown."""
        print("\nПолучен сигнал завершения...")
        self._shutdown_event.set()

    async def run(self) -> None:
        """Run the application."""
        loop = asyncio.get_running_loop()

        # Register signal handlers
        for sig in (signal.SIGINT, signal.SIGTERM):
            loop.add_signal_handler(sig, self._signal_handler)

        # Start workers
        for i in range(3):
            task = asyncio.create_task(self.worker(f"W{i}"))
            self._tasks.add(task)
            task.add_done_callback(self._tasks.discard)

        # Wait for shutdown signal
        await self._shutdown_event.wait()

        # Wait for all tasks to finish (with timeout)
        print("Ожидаю завершения задач...")
        if self._tasks:
            await asyncio.wait(self._tasks, timeout=5.0)

            # Force cancel remaining tasks
            for task in self._tasks:
                task.cancel()

        print("Приложение завершено")

# async def main():
#     app = GracefulApp()
#     await app.run()
#
# asyncio.run(main())

Паттерн Fan-Out / Fan-In

Распределение работы по нескольким воркерам (fan-out) и сбор результатов (fan-in):

import asyncio
from dataclasses import dataclass

@dataclass
class ProcessResult:
    """Result of processing."""
    item_id: int
    value: str
    worker: str

async def fan_out_fan_in(
    items: list[int],
    num_workers: int = 4,
) -> list[ProcessResult]:
    """Distribute work across workers and collect results."""
    input_queue: asyncio.Queue[int | None] = asyncio.Queue()
    output_queue: asyncio.Queue[ProcessResult] = asyncio.Queue()

    async def worker(name: str) -> None:
        """Process items from input queue."""
        while True:
            item = await input_queue.get()
            if item is None:
                input_queue.task_done()
                break

            # Simulate variable processing time
            await asyncio.sleep(0.1 * (item % 5 + 1))
            result = ProcessResult(
                item_id=item,
                value=f"processed-{item}",
                worker=name,
            )
            await output_queue.put(result)
            input_queue.task_done()

    # Fan-out: enqueue all items
    for item in items:
        await input_queue.put(item)

    # Add poison pills for each worker
    for _ in range(num_workers):
        await input_queue.put(None)

    # Start workers
    workers = [
        asyncio.create_task(worker(f"W{i}"))
        for i in range(num_workers)
    ]

    # Wait for all work to complete
    await input_queue.join()
    await asyncio.gather(*workers)

    # Fan-in: collect results
    results: list[ProcessResult] = []
    while not output_queue.empty():
        results.append(await output_queue.get())

    return results

async def main() -> None:
    items = list(range(20))
    results = await fan_out_fan_in(items, num_workers=4)

    print(f"Обработано {len(results)} элементов:")
    for r in sorted(results, key=lambda x: x.item_id):
        print(f"  #{r.item_id} -> {r.value} [{r.worker}]")

asyncio.run(main())

Интеграция sync и async кода

Иногда нужно вызвать sync-код из async-контекста или наоборот:

import asyncio
from concurrent.futures import ProcessPoolExecutor

def cpu_heavy_task(n: int) -> int:
    """CPU-bound computation (runs in process pool)."""
    total = 0
    for i in range(n):
        total += i * i
    return total

def blocking_io_task(path: str) -> str:
    """Blocking I/O operation (runs in thread pool)."""
    import time
    time.sleep(1)  # Simulate file read
    return f"Данные из {path}"

async def main() -> None:
    loop = asyncio.get_running_loop()

    # Run blocking I/O in thread pool (default executor)
    result = await loop.run_in_executor(
        None,  # Default ThreadPoolExecutor
        blocking_io_task,
        "/tmp/data.txt",
    )
    print(f"I/O результат: {result}")

    # Run CPU-bound task in process pool
    with ProcessPoolExecutor() as pool:
        result = await loop.run_in_executor(
            pool,
            cpu_heavy_task,
            10_000_000,
        )
        print(f"CPU результат: {result}")

asyncio.run(main())

Практический пример: асинхронный веб-скрапер

import asyncio
from dataclasses import dataclass, field
from collections.abc import AsyncIterator

@dataclass
class ScrapedPage:
    """Result of scraping a page."""
    url: str
    title: str
    links: list[str] = field(default_factory=list)

@dataclass
class AsyncScraper:
    """Web scraper with rate limiting and retry."""
    max_concurrent: int = 5
    rate_limit: float = 2.0  # Requests per second

    async def scrape_page(self, url: str) -> ScrapedPage:
        """Scrape a single page."""
        await asyncio.sleep(0.3)  # Simulate HTTP request
        return ScrapedPage(
            url=url,
            title=f"Страница: {url}",
            links=[f"{url}/link{i}" for i in range(3)],
        )

    async def scrape_all(self, urls: list[str]) -> AsyncIterator[ScrapedPage]:
        """Scrape multiple URLs with concurrency control."""
        semaphore = asyncio.Semaphore(self.max_concurrent)

        async def limited_scrape(url: str) -> ScrapedPage:
            async with semaphore:
                return await self.scrape_page(url)

        tasks = [
            asyncio.create_task(limited_scrape(url))
            for url in urls
        ]

        for task in asyncio.as_completed(tasks):
            yield await task

async def main() -> None:
    scraper = AsyncScraper(max_concurrent=3)
    urls = [f"https://example.com/page/{i}" for i in range(10)]

    results: list[ScrapedPage] = []
    async for page in scraper.scrape_all(urls):
        results.append(page)
        print(f"Скрапнуто: {page.title} ({len(page.links)} ссылок)")

    print(f"\nВсего обработано: {len(results)} страниц")

asyncio.run(main())

Проверь себя

Что такое паттерн Fan-Out / Fan-In?

Какой паттерн используется для ограничения частоты запросов к API?

Как выполнить CPU-bound задачу в async-коде без блокировки event loop?

Когда circuit breaker переходит в состояние OPEN?

Почему важно использовать `async with httpx.AsyncClient()` вместо создания клиента без контекстного менеджера?