Асинхронные 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())