Async ETL în Python 2026: httpx, asyncio și pandas pentru pipeline-uri rapide de date

Construiește pipeline-uri ETL asincrone în Python cu httpx 0.28, asyncio.TaskGroup, Semaphore, tenacity 9 și pandas. Rate limiting, retry cu backoff exponential și validare Pydantic v2.

Async ETL Python 2026: httpx Guide

Actualizat: 14 septembrie 2026

ETL asincron în Python înseamnă să rulezi sute sau mii de apeluri I/O (HTTP, DB, S3) în paralel pe un singur thread folosind asyncio, apoi să materializezi rezultatele într-un pandas.DataFrame pentru transformare. E de trei până la zece ori mai rapid decât un pipeline sincron pentru sarcini I/O-bound. În 2026, combinația canonică este httpx 0.28 pentru client HTTP async, asyncio.TaskGroup din Python 3.11+ pentru orchestrare cu propagare corectă a excepțiilor, tenacity 9.0 pentru retry cu backoff exponential și Pydantic v2.10 pentru validarea răspunsurilor. Am migrat cinci pipeline-uri de producție pe această stivă în ultimele 18 luni și am strâns tiparele care contează.

  • httpx 0.28 (mai 2025) rămâne clientul HTTP async standard cu HTTP/2 nativ, timeout-uri configurabile per operațiune și un pool de conexiuni deterministic.
  • asyncio.TaskGroup (Python 3.11+) înlocuiește asyncio.gather în producție: anulează frații la prima excepție și expune ExceptionGroup.
  • Limitează concurența cu asyncio.Semaphore pentru resurse locale și cu aiolimiter pentru rate limiting per-secundă impus de API-ul remote.
  • tenacity 9.0 aduce decorator @retry compatibil async, oprire pe status HTTP specifice și logging structurat prin before_sleep.
  • Pandas 2.2+ acceptă construcție eficientă din liste de dict-uri Pydantic. Folosește DataFrame.from_records cu tipuri PyArrow pentru un footprint de memorie 30% mai mic.
  • PEP 703 (no-GIL, experimental în 3.13, opt-in în 3.14) nu înlocuiește asyncio pentru ETL; asyncio rămâne modelul corect pentru I/O concurent.

Ce este ETL asincron și când merită

Un pipeline ETL clasic (extrage, transformă, încarcă) petrece adesea 80–95% din timp așteptând I/O: apeluri HTTP către un API terț, interogări către un warehouse, upload-uri către S3. Cu requests sincron, threadul principal blochează pentru fiecare cerere; două mii de endpoint-uri × 400 ms latență fac ~13 minute doar de așteptare. Async schimbă modelul complet.

O singură buclă de evenimente comută între cereri în timp ce fiecare așteaptă rețeaua, colapsând acele minute în ~30 de secunde pentru același număr de apeluri. Regula pe care o aplic: dacă un pas al pipeline-ului petrece >40% din wall time în așteptare I/O, migrarea la asyncio dă câștig măsurabil.

Pentru CPU-bound (parsing intens, feature engineering pe milioane de rânduri) rămâne multiprocessing sau Polars. În practică, extragerea (E) e aproape întotdeauna I/O-bound și beneficiază; transformarea (T) rămâne sincronă în pandas; încărcarea (L) redevine async dacă scrii către S3, BigQuery sau un endpoint HTTP.

httpx 0.28 AsyncClient: setup corect pentru producție

httpx a devenit standardul de facto pentru clienți HTTP async în ecosistemul de date pentru că API-ul lui oglindește requests, dar suportă HTTP/2, timeouts granulare și un connection pool controlabil. Un anti-pattern pe care îl văd des: crearea unui AsyncClient nou în fiecare funcție. Fiecare instanță deschide propriile conexiuni TCP, iar TLS handshake-ul dominează latența. Instanțiază clientul o singură dată pentru toată rularea și trece-l explicit.

import httpx
import asyncio

# Configurare care se poate refolosi pentru toată aplicația
LIMITS = httpx.Limits(max_connections=200, max_keepalive_connections=50)
TIMEOUT = httpx.Timeout(10.0, connect=3.0, read=8.0)

async def fetch_json(client: httpx.AsyncClient, url: str) -> dict:
    response = await client.get(url)
    response.raise_for_status()
    return response.json()

async def main(urls: list[str]) -> list[dict]:
    async with httpx.AsyncClient(
        limits=LIMITS,
        timeout=TIMEOUT,
        http2=True,
        headers={"User-Agent": "pipelines-etl/1.4"},
    ) as client:
        results = await asyncio.gather(*(fetch_json(client, u) for u in urls))
    return results

http2=True necesită pachetul opțional h2 (pip install "httpx[http2]"). Beneficiul apare doar dacă serverul o suportă și cererile ating același host. HTTP/2 multiplexează mai multe cereri pe un singur TCP, ceea ce reduce dramatic overhead-ul TLS. Pentru API-uri diverse (multi-host), câștigul e marginal. Vezi documentația oficială httpx pentru referință completă a opțiunilor.

TaskGroup vs asyncio.gather: care e diferența în 2026

Timp de un deceniu, asyncio.gather(*coros, return_exceptions=False) a fost șablonul default. Problema lui e comportamentul la eroare: dacă o coroutină eșuează, celelalte continuă în background până se termină singure, iar excepția e propagată doar către caller. În ETL asta înseamnă cereri fantomă care încă lovesc API-ul remote și logging incomplet. asyncio.TaskGroup, disponibil din Python 3.11, rezolvă problema: la prima excepție anulează cooperativ toate task-urile fraților și emite un ExceptionGroup care conține toate erorile.

async def extract_all(client: httpx.AsyncClient, urls: list[str]) -> list[dict]:
    results: list[dict] = []
    try:
        async with asyncio.TaskGroup() as tg:
            tasks = [tg.create_task(fetch_json(client, u)) for u in urls]
    except* httpx.HTTPStatusError as eg:
        # Fiecare eroare este disponibilă separat prin ExceptionGroup
        for err in eg.exceptions:
            logger.warning("HTTP error: %s", err)
        raise
    return [t.result() for t in tasks]

Sintaxa except* (PEP 654, GA în 3.11) prinde selectiv un tip din ExceptionGroup, iar celelalte se propagă. Pentru pipeline-uri care trebuie să continue chiar și cu eșecuri parțiale (de exemplu, 998 din 1000 de rânduri sunt suficiente), rămâne gather(*, return_exceptions=True) plus filtrare manuală. Pentru all-or-nothing folosește TaskGroup.

Cum limităm concurența și rate limiting cu Semaphore și aiolimiter

Zece mii de cereri simultane către un endpoint terț înseamnă doar un lucru: 429 Too Many Requests, urmat de ban temporar. Sincer, e cea mai frecventă cauză a call-urilor de escaladare pe care le-am văzut. Există două cadrane de limitare care se combină în producție: concurrency (câte cereri se pot afla simultan în zbor) și rate (câte cereri pornesc pe secundă). Pentru concurrency, asyncio.Semaphore e primitiva corectă; pentru rate, aiolimiter implementează algoritmul leaky-bucket și funcționează impecabil ca context manager async.

from aiolimiter import AsyncLimiter

# 30 cereri max în zbor, 100 cereri/secundă
concurrency = asyncio.Semaphore(30)
rate_limit = AsyncLimiter(max_rate=100, time_period=1.0)

async def fetch_bounded(client, url):
    async with concurrency, rate_limit:
        return await fetch_json(client, url)

async def main(urls):
    async with httpx.AsyncClient(limits=LIMITS, timeout=TIMEOUT) as client:
        async with asyncio.TaskGroup() as tg:
            tasks = [tg.create_task(fetch_bounded(client, u)) for u in urls]
    return [t.result() for t in tasks]

Regulă empirică: setează Semaphore cam la 60–80% din max_connections din httpx.Limits. Dacă permiți mai mult decât are pool-ul, cererile stau în coadă interior în httpx și nu ai un metric ușor pentru saturare. Pentru API-urile care declară limite documentate (Stripe: 100/sec, GitHub: 5000/oră autenticat), folosește AsyncLimiter cu valorile lor exacte. Leaky bucket-ul netezește vârfurile mai bine decât un contor manual.

Retry, timeout și circuit breaker cu tenacity 9.0

Rețelele au ganduri proprii. Un endpoint stabil răspunde 200 timp de 6 luni, apoi returnează 502 trei minute la rând. tenacity 9.0 (mai 2025) e biblioteca de retry de facto și oferă decorator nativ pentru async, backoff exponential cu jitter și cârlige de logging structurat. Cheia e să nu re-încerci la orice: retry pe 500/502/503/504 și timeout, dar niciodată pe 400/401/403/404. Repetarea unei cereri malformate doar risipește budget-ul.

from tenacity import (
    AsyncRetrying, retry, retry_if_exception_type,
    stop_after_attempt, wait_exponential_jitter,
    before_sleep_log,
)
import logging

logger = logging.getLogger("etl")

RETRYABLE = (httpx.TimeoutException, httpx.NetworkError)

def is_retryable_http(exc: BaseException) -> bool:
    if isinstance(exc, httpx.HTTPStatusError):
        return exc.response.status_code in {500, 502, 503, 504, 429}
    return isinstance(exc, RETRYABLE)

@retry(
    retry=retry_if_exception_type((httpx.HTTPError,)) & is_retryable_http,
    stop=stop_after_attempt(5),
    wait=wait_exponential_jitter(initial=0.5, max=10),
    before_sleep=before_sleep_log(logger, logging.WARNING),
    reraise=True,
)
async def fetch_with_retry(client, url):
    response = await client.get(url)
    response.raise_for_status()
    return response.json()

Câteva puncte pe care le-am învățat pe pielea mea. Jitterul e obligatoriu; fără el, o mie de clienți care re-încearcă exact la 1s produc un thundering herd care doboară endpoint-ul din nou. reraise=True propagă excepția originală după ultima încercare, altfel primești un RetryError care pierde stack trace-ul util. Pentru 429 specifically, mulți API răspund cu header Retry-After; tenacity nu îl citește automat, dar poți implementa un wait custom care extrage headerul din exc.response.headers. Documentația completă e la tenacity.readthedocs.io.

Validarea răspunsurilor cu Pydantic v2 și materializare în pandas

Un răspuns JSON nu e un contract; până nu îl validezi, nu ai date, ai text. Pydantic v2.10 (august 2026) validează liste de mii de obiecte în milisecunde grație core-ului scris în Rust, iar TypeAdapter îți permite să validezi un list[Model] direct fără să declari un wrapper. Am pățit-o pe un pipeline de facturare: downstream așteptam amount: float și primeam accidental un string. Mai bine crapi la validare decât după 40 de minute în transformare.

from pydantic import BaseModel, Field, TypeAdapter
from datetime import datetime
import pandas as pd

class OrderRow(BaseModel):
    order_id: str
    customer_id: str
    amount: float = Field(ge=0)
    currency: str = Field(min_length=3, max_length=3)
    created_at: datetime

OrderList = TypeAdapter(list[OrderRow])

async def load_orders(client, urls) -> pd.DataFrame:
    raw = await extract_all(client, urls)
    # raw este list[dict]; validăm în bloc
    orders = OrderList.validate_python(raw)
    df = pd.DataFrame.from_records([o.model_dump() for o in orders])
    return df.convert_dtypes(dtype_backend="pyarrow")

dtype_backend="pyarrow" transformă coloanele în tipuri Arrow, cu un footprint de memorie 30–40% mai mic pentru string-uri și suport nativ pentru NA. Combinația Pydantic + pandas cu backend Arrow e stiva mea default pentru pipeline-uri noi în 2026. Pentru context suplimentar pe partea de validare, vezi ghidul de calitate a datelor cu Great Expectations.

Streaming pentru payload-uri mari cu ijson

Un răspuns JSON de 800 MB nu încape confortabil în memorie de trei ori (bytes → text → dict). Pentru API-uri care returnează liste imense (rapoarte financiare, export-uri de utilizatori), fac streaming cu httpx.stream plus ijson 3.3, care parseoză iterativ elementele unui array fără să încarce totul.

import ijson

async def stream_records(client, url, path="records.item"):
    async with client.stream("GET", url) as response:
        response.raise_for_status()
        async for record in ijson.items_async(response.aiter_bytes(), path):
            yield record

async def build_frame(client, url):
    rows = []
    async for rec in stream_records(client, url):
        rows.append(rec)
        if len(rows) >= 50_000:
            yield pd.DataFrame(rows)
            rows.clear()
    if rows:
        yield pd.DataFrame(rows)

Batch-uri de 50k rânduri sunt un compromis rezonabil: destul de mari să amortizeze costul construcției DataFrame-ului, destul de mici încât un peak să nu te dea afară din containerul de 2 GB. Cheia async aici este că aiter_bytes() nu blochează event loop-ul, așa că poți continua să extragi din alte surse în paralel. Pentru pași downstream analitici pe volume mari, ia în calcul DuckDB ca engine de analiză; poate citi direct un fișier Parquet scris din stream.

asyncio vs ThreadPoolExecutor vs ProcessPoolExecutor

Cea mai comună întrebare pe care o primesc: „De ce nu folosesc doar concurrent.futures?" Depinde de natura sarcinii. Tabelul de mai jos rezumă când fiecare model câștigă în 2026:

CriteriuasyncioThreadPoolExecutorProcessPoolExecutor
Tip sarcinăI/O non-blocant (HTTP, sockets)I/O blocant (biblioteci sincrone, filesystem)CPU-bound (parsing, algoritmi)
Overhead per task~2–5 μs (coroutină)~50–200 μs (thread)~1–5 ms (proces)
Concurență practică10.000+ task-uri~100–500 threads≤ nucleee CPU
Cooperare cu httpx/aiohttpNativNecesită wrappersNu suportă
DebuggingMediu (traceback-uri prin ExceptionGroup)Ușor (stack-uri per thread)Ușor (procese separate)
Impact GILIgnoră (single-threaded)Blocat pentru CPUNu are (procese distincte)

Regulă practică: extragere (E) → asyncio + httpx; transformare (T) → pandas sincron sau Polars; fan-out CPU-bound → ProcessPoolExecutor. Pentru un layer legacy care încă folosește requests sincron într-un pipeline nou, pune-l pe asyncio.to_thread. Nu rescrie totul într-o singură iterație.

Observabilitate, shutdown grațios și PEP 703

Un pipeline care rulează 12 ore fără metrici e o cutie neagră. Instrumentez trei semnale minime: (1) latency histogram per endpoint, (2) contor de retry pe status HTTP, (3) task-uri în zbor la un moment dat. Cu prometheus_client și un Gauge care crește la începutul fiecărui fetch_bounded și scade la sfârșit, primesc grafic pe saturare fără să adaug OpenTelemetry complet.

from prometheus_client import Counter, Histogram, Gauge, start_http_server

request_latency = Histogram("etl_request_seconds", "HTTP latency", ["endpoint"])
request_retries = Counter("etl_request_retries_total", "Retries by status", ["status"])
inflight = Gauge("etl_inflight_requests", "Requests in flight")

async def fetch_measured(client, url, endpoint_label):
    inflight.inc()
    try:
        with request_latency.labels(endpoint=endpoint_label).time():
            return await fetch_with_retry(client, url)
    finally:
        inflight.dec()

Shutdown grațios contează la fel de mult. Când Kubernetes trimite SIGTERM (de exemplu, la re-deploy), pipeline-ul are ~30 s să încheie ce e în zbor. Prinzi semnalul cu loop.add_signal_handler, setezi un asyncio.Event, iar coroutinele verifică event.is_set() înainte să lanseze cereri noi. TaskGroup-ul curent lasă cererile deja pornite să termine natural. Pentru contextul de deploy în producție, vezi ghidul de servire ML cu FastAPI.

În fine, PEP 703 (no-GIL, opt-in în Python 3.14 lansat octombrie 2025) generează multă confuzie. Fără GIL, threadurile pot rula în paralel cod Python, dar asta nu înlocuiește asyncio pentru ETL. Threadurile rămân scumpe (kB-uri de stack per thread), au overhead de context switch, iar bibliotecile async au ecosisteme mature (httpx, asyncpg, aioboto3). Modelul concurrent corect pentru I/O rămâne single-threaded event loop; no-GIL îmbunătățește doar sarcinile CPU-bound când preferi threads în locul proceselor. Vezi textul complet al PEP 703.

Întrebări frecvente

Care este diferența dintre asyncio.gather și TaskGroup în Python?

asyncio.gather continuă să ruleze task-urile fraților chiar și după o excepție, în timp ce asyncio.TaskGroup (Python 3.11+) le anulează automat la prima eroare și le grupează într-un ExceptionGroup. Pentru pipeline-uri de producție, TaskGroup e alegerea corectă.

Este httpx mai rapid decât aiohttp în 2026?

Pentru trafic tipic de ETL, diferența e sub 5% și nu decide alegerea. httpx câștigă pe API (identic cu requests), suport HTTP/2 nativ și integrare cu instrumentele de test; aiohttp rămâne marginal mai rapid la thousands-per-second și mai matur pentru servere. Pentru clienți ETL, httpx e standardul recomandat.

Cum limitez rata cererilor async pentru a evita HTTP 429?

Folosește aiolimiter.AsyncLimiter(max_rate, time_period) pentru rate limiting bazat pe leaky bucket, combinat cu un asyncio.Semaphore pentru concurență maximă. Setează valorile la limitele documentate ale API-ului remote (de exemplu 100 cereri/secundă).

Pandas funcționează cu asyncio direct?

Nu. Operațiile pandas sunt sincrone și blochează event loop-ul. Extrage și validează datele async, apoi construiește DataFrame-ul într-un pas sincron. Pentru transformări grele, folosește asyncio.to_thread sau mută-le pe un ProcessPoolExecutor.

Când NU trebuie să folosesc asyncio pentru ETL?

Când pipeline-ul e CPU-bound (parsing complex, calcule pe milioane de rânduri), asyncio nu ajută. Folosește multiprocessing sau Polars/DuckDB. La fel dacă toate bibliotecile din stivă sunt sincrone și nu ai timp să le înfășori, un pipeline sincron simplu bate un pipeline async prost scris.

Tomás Oliveira
Despre Autor Tomás Oliveira

Python backend developer who came to data work via FastAPI. Bridges the messy world between APIs and pipelines.