Async ETL v Pythonu 2026: httpx, asyncio a pandas pro rychlé datové pipeline

Praktický průvodce async ETL v Pythonu 2026: httpx AsyncClient, asyncio.TaskGroup, Semaphore, tenacity retry, Pydantic v2 a nasazení v produkci s příklady kódu.

Async ETL v Pythonu 2026: httpx a asyncio

Aktualizováno: 12. září 2026

Async ETL v Pythonu je způsob, jak paralelně stahovat data z desítek nebo stovek HTTP API pomocí jediného vlákna a event loopu, v roce 2026 typicky s knihovnou httpx 0.28, primitivem asyncio.TaskGroup (Python 3.11+) a následným zpracováním v pandas nebo Polars. Oproti klasickému requests + ThreadPoolExecutor dostanete 5–20× vyšší propustnost při stejném CPU rozpočtu, protože event loop se během I/O čekání věnuje dalším úlohám místo blokování vlákna. Přiznám se, že jsem tenhle návod sepsal na základě pipeline, kterou jsem loni migroval z threadingu na asyncio a která teď zvládne 40 000 requestů za minutu na jednom podu.

  • Async ETL se hodí, když jsou úzkým hrdlem síťové I/O (API, S3, databáze), nikoli CPU výpočty.
  • httpx 0.28 podporuje HTTP/2, connection pooling a je API-kompatibilní s requests, což usnadňuje migraci.
  • Bez asyncio.Semaphore nebo bounded worker poolu vyhoří target API za 10 sekund. Vždy omezujte souběžnost.
  • Blokující volání (pandas read_csv, requests, čtení souboru) v async funkci zablokuje celý event loop, používejte asyncio.to_thread.
  • Retry s exponenciálním backoffem přes tenacity plus validace odpovědí přes Pydantic v2 udělá pipeline produkčně použitelnou.
  • Pro pipeline nad ~100 GB dat kombinujte async fetch s Polars nebo DuckDB, ne s eager pandas.

Proč přejít na async ETL v roce 2026?

ETL pipeline většinou tráví 80–95 % času čekáním: na HTTP odpověď, na databázový cursor, na S3 PutObject. Když toto čekání děláte synchronně, jedno vlákno je celou dobu zaparkované. ThreadPoolExecutor to řeší tak, že si vezme 32 nebo 64 vláken. Funguje to, ale každé vlákno má vlastní stack (~8 MB v Linuxu), context switching není zadarmo a GIL vás stejně limituje při jakémkoli decode/parse kroku. Async model naproti tomu zvládne desítky tisíc "úloh" v jednom vlákně, protože přepínání mezi nimi je pouhá výměna callbacku v event loopu.

Konkrétně v roce 2026 tři věci posunuly async ETL do kategorie "výchozí volba" místo "hračka pro nadšence". Za prvé Python 3.11 přinesl asyncio.TaskGroup se strukturovanou konkurencí, takže konečně máte guarantee, že žádná coroutine neuteče nezachycená. Za druhé httpx 0.28 (vydáno v červnu 2026) plně podporuje HTTP/2 multiplexing, takže i 1000 requestů na jednu doménu jde přes hrstku TCP spojení. A za třetí PEP 703 (no-GIL Python) v beta buildech CPython 3.14t ukazuje, že async už není konkurent threadingu, ale komplementární technika: async pro I/O, free-threaded pro CPU.

Praktické měření z mé loňské migrace vypadalo takto. Pipeline stahující 8 000 článků z REST API, každý ~30 KB JSON, běžela synchronně 14 minut a se ThreadPoolExecutor(32) 96 sekund. Async verze s httpx a Semaphore(50) jela 41 sekund na stejném stroji, s poloviční RAM a bez race conditions v pandas writeru. To je typický řád zlepšení, který v I/O-heavy pipeline uvidíte.

httpx AsyncClient a asyncio.TaskGroup

Základní stavební kámen je httpx.AsyncClient. Udržuje connection pool, spravuje HTTP/2 sessions a jeho API je téměř shodné se synchronním httpx.Client. Klíčové pravidlo: jeden AsyncClient na celou pipeline, nikdy nevytvářejte nový per-request. Bez sdíleného poolu ztratíte HTTP/2 multiplexing i keep-alive a latence poletí o 3–5× nahoru. Tohle jsem si zjistil až po druhém profilování, kdy mi jeden kolega ukázal, že vytváříme klienta v smyčce.

import asyncio
import httpx
import pandas as pd

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

async def fetch_all(urls: list[str]) -> list[dict]:
    limits = httpx.Limits(max_connections=100, max_keepalive_connections=50)
    async with httpx.AsyncClient(http2=True, limits=limits) as client:
        async with asyncio.TaskGroup() as tg:
            tasks = [tg.create_task(fetch_one(client, u)) for u in urls]
        return [t.result() for t in tasks]

urls = [f"https://api.example.com/items/{i}" for i in range(500)]
records = asyncio.run(fetch_all(urls))
df = pd.DataFrame.from_records(records)
print(df.head())

asyncio.TaskGroup je oproti staršímu asyncio.gather() lepší ve dvou věcech. Pokud jedna task vyhodí výjimku, TaskGroup zruší všechny sourozence a propaguje ExceptionGroup ven, takže žádná coroutine nezůstane osiřelá v pozadí. A díky async with je jasně ohraničeno, kde konkurence začíná a končí, což při čtení kódu za rok šetří hodiny.

Jak omezit počet souběžných requestů

Bez omezení souběžnosti asyncio pipeline pošle všechny requesty najednou. Když jich je 5 000, target API vrátí HTTP 429 nebo prostě zavře spojení. Řešením je asyncio.Semaphore, který funguje jako bounded permit a povolí jen N současně běžících operací.

import asyncio
import httpx

MAX_CONCURRENT = 50

async def fetch_bounded(sem: asyncio.Semaphore, client: httpx.AsyncClient, url: str) -> dict:
    async with sem:
        response = await client.get(url, timeout=15.0)
        response.raise_for_status()
        return response.json()

async def fetch_all(urls: list[str]) -> list[dict]:
    sem = asyncio.Semaphore(MAX_CONCURRENT)
    limits = httpx.Limits(max_connections=MAX_CONCURRENT * 2)
    async with httpx.AsyncClient(http2=True, limits=limits) as client:
        async with asyncio.TaskGroup() as tg:
            tasks = [tg.create_task(fetch_bounded(sem, client, u)) for u in urls]
        return [t.result() for t in tasks]

Hodnotu MAX_CONCURRENT ladíme empiricky. Dobré startovní číslo je 50 pro veřejné API a 200 pro interní služby ve stejném VPC. Pokud target API dokumentuje rate limit (třeba 100 req/s), přidejte na jistotu ještě token bucket. Knihovna aiolimiter to řeší v pěti řádcích:

from aiolimiter import AsyncLimiter

rate_limiter = AsyncLimiter(max_rate=100, time_period=1)  # 100 requestů/sekundu

async def fetch_ratelimited(client: httpx.AsyncClient, url: str) -> dict:
    async with rate_limiter:
        response = await client.get(url, timeout=15.0)
        response.raise_for_status()
        return response.json()

Kombinace Semaphore (limit souběžnosti) plus AsyncLimiter (limit rychlosti) pokryje 95 % produkčních případů. Kdo si to nehlídá, ten pravidelně dostane e-mail od poskytovatele API o "abuse" a v horším případě rate-ban na 24 hodin. Sám jsem takový e-mail dostal, není to legrace.

Retry a rate limiting s tenacity

Ať děláte cokoli, sítě padají. HTTP 502, RST connection, timeout, a při 10 000 requestech to dostanete na 20 z nich. Bez retry logiky vaše ETL pipeline selže na první chybě a přijdete o celou dávku. Nejčistší řešení je dekorátor z knihovny tenacity 9.0, která podporuje async funkce out-of-the-box:

from tenacity import (
    retry,
    stop_after_attempt,
    wait_exponential,
    retry_if_exception_type,
)
import httpx

@retry(
    stop=stop_after_attempt(5),
    wait=wait_exponential(multiplier=1, min=1, max=30),
    retry=retry_if_exception_type((httpx.TransportError, httpx.HTTPStatusError)),
    reraise=True,
)
async def fetch_with_retry(client: httpx.AsyncClient, url: str) -> dict:
    response = await client.get(url, timeout=15.0)
    if response.status_code == 429:
        raise httpx.HTTPStatusError(
            "Rate limited", request=response.request, response=response
        )
    response.raise_for_status()
    return response.json()

Exponenciální backoff (1s, 2s, 4s, 8s, 16s) je důležitý, protože bez něj váš retry jen zesílí problém na přetíženém serveru. Speciálně na HTTP 429 respektujte hlavičku Retry-After, pokud ji API posílá:

async def fetch_respecting_retry_after(client: httpx.AsyncClient, url: str) -> dict:
    for attempt in range(5):
        response = await client.get(url, timeout=15.0)
        if response.status_code == 429:
            wait = float(response.headers.get("Retry-After", 2 ** attempt))
            await asyncio.sleep(wait)
            continue
        response.raise_for_status()
        return response.json()
    raise RuntimeError(f"Vzdal jsem to po 5 pokusech: {url}")

Validace odpovědí pomocí Pydantic v2

Když API vrací JSON, syrový dict je bomba s časovanou rozbuškou. Chybějící klíč nebo string tam, kde má být int, zabije pipeline třeba až v pandas groupby. Pydantic v2 (v 2026 verze 2.10) validuje odpověď hned po fetch a při chybě spadne s přesnou informací, který field a proč. Bonus: Pydantic model dokumentuje očekávané schéma přímo v kódu.

from datetime import datetime
from pydantic import BaseModel, Field, ValidationError
import httpx

class Article(BaseModel):
    id: int
    title: str = Field(min_length=1, max_length=500)
    author_id: int
    published_at: datetime
    view_count: int = Field(ge=0, default=0)

async def fetch_article(client: httpx.AsyncClient, article_id: int) -> Article | None:
    try:
        response = await client.get(f"https://api.example.com/articles/{article_id}")
        response.raise_for_status()
        return Article.model_validate_json(response.content)
    except ValidationError as exc:
        print(f"Nevalidní odpověď pro článek {article_id}: {exc.errors()[:2]}")
        return None
    except httpx.HTTPError as exc:
        print(f"HTTP chyba: {exc}")
        return None

model_validate_json je v Pydantic v2 rychlejší než model_validate(response.json()), protože parser přeskočí Python dict jako mezikrok. Detailní rozbor validace najdete v článku Pydantic v2 pro validaci dat v Pythonu, na tomto místě jen zdůrazním jedno: neposílejte do pandas surové dict. Vždycky přes validovaný model, i za cenu pár mikrosekund navíc na řádek.

Streaming velkých odpovědí do pandas

Pokud API vrací 200 MB JSON pole, načtení celého do paměti a následné pd.DataFrame(records) vám nafoukne heap na 1,5 GB (JSON má overhead ~7× nad DataFrame díky Python objektům). Řešením je async streaming přes httpx.stream a inkrementální parsování s ijson:

import httpx
import ijson
import pandas as pd

async def stream_to_dataframe(url: str, batch_size: int = 5000) -> pd.DataFrame:
    batches: list[pd.DataFrame] = []
    buffer: list[dict] = []
    async with httpx.AsyncClient(timeout=None) as client:
        async with client.stream("GET", url) as response:
            response.raise_for_status()
            async for chunk in response.aiter_bytes(chunk_size=65536):
                for item in ijson.items(chunk, "items.item"):
                    buffer.append(item)
                    if len(buffer) >= batch_size:
                        batches.append(pd.DataFrame.from_records(buffer))
                        buffer.clear()
    if buffer:
        batches.append(pd.DataFrame.from_records(buffer))
    return pd.concat(batches, ignore_index=True)

Pro opravdu velké výsledky (nad ~2 GB) doporučuju přeskočit pandas úplně a psát rovnou do Parquet přes pyarrow.parquet.ParquetWriter nebo do DuckDB in-memory tabulky. Pandas do 2 GB je pohodová, nad to začíná bolet. Polars nebo DuckDB tam mají 5–10× lepší throughput.

Async vs threading vs multiprocessing

Otázka "kdy použít async místo threadingu" má krátkou odpověď: pokud pipeline tráví většinu času v await, jděte do asyncio; pokud v CPU výpočtech, jděte do multiprocessing. Threading v Pythonu je kvůli GIL užitečný jen tam, kde volaná knihovna GIL uvolňuje (většina I/O, NumPy vektorizace, síťové sockety), a i tam ho async v roce 2026 překonává v propustnosti i paměti.

Dimenzeasyncio + httpxThreadPoolExecutorProcessPoolExecutor
Nejlepší proSíťové I/O, mnoho HTTP callsBlokující I/O ve staré knihovněCPU-heavy výpočty (parse, hash)
Souběžné úlohy na 4-core / 16 GB10 000+~500~4–8
Paměť per úloha~10 KB~8 MB (stack vlákna)~50 MB (proces + interpreter)
GIL kontenceŽádná (jedno vlákno)Vysoká pod zátěžíŽádná (více procesů)
Sdílení stavuSnadné (jeden proces)Snadné, ale potřeba lockůIPC přes Queue/Pipe
DebuggingNáročnější (stack traces)StandardníStandardní
Křivka učeníStřední (async/await)NízkáNízká, ale IPC bolí

Nejčastější hybrid, který v produkci vidím: async fetch přes httpx pro I/O plus asyncio.to_thread(pd.read_parquet, path) pro sync knihovny, které volání zvenku nemají. Pandas je stále synchronní. Nesnažte se to obejít pseudo-async wrapperem. Klasický ETL pattern s pandas a SQLAlchemy pořád funguje, jen fetch fázi obalte do async a zbytek nechte být.

Rozdíly mezi httpx a aiohttp

Historicky byl aiohttp defaultní volbou pro async HTTP v Pythonu, ale od roku 2024 se váha přesouvá k httpx. Důvod: httpx umí sync i async se stejným API, což je zásadní pro postupnou migraci existující codebase. Aiohttp je čistě async a nutí vás přepsat i moduly, které async vůbec nepotřebují.

Konkrétní srovnání pro ETL workload: httpx 0.28 podporuje HTTP/2, aiohttp 3.10 zatím ne (v roadmapě je pro 4.0). HTTP/2 znamená multiplexing přes jedno TCP spojení, a pro API, které vrací 1000 malých odpovědí, to dělá 30–50 % latenci navíc, když ho nemáte. Naopak aiohttp má o 10–15 % nižší per-request overhead v pure benchmarcích, takže na extrémně vysokoobjemových systémech (100k+ req/s) může být rychlejší. Pro 99 % ETL pipeline je to jedno.

import aiohttp
import asyncio

async def fetch_aiohttp(urls: list[str]) -> list[dict]:
    async with aiohttp.ClientSession() as session:
        async def get(url):
            async with session.get(url) as response:
                return await response.json()
        return await asyncio.gather(*(get(u) for u in urls))

Moje doporučení pro nový projekt: vždycky httpx. Migrace z requests je triviální (stejné API), HTTP/2 je zdarma a oficiální dokumentace je čitelnější než aiohttpová. Aiohttp má smysl jen v projektech, které už na něm stojí a byla by drahá migrace.

Nasazení a monitoring v produkci

Async pipeline v produkci potřebuje tři věci, které lokálně přeskočíte a v produkci vás pak štípnou: graceful shutdown, metriky a izolaci selhání. Bez nich se pipeline zasekne při SIGTERM, ztratíte přehled o latencích a jedna špatná odpověď potopí celou dávku.

import asyncio
import signal
import time
from contextlib import suppress

async def run_pipeline(urls: list[str], shutdown: asyncio.Event) -> None:
    limits = httpx.Limits(max_connections=100)
    async with httpx.AsyncClient(http2=True, limits=limits) as client:
        sem = asyncio.Semaphore(50)
        async with asyncio.TaskGroup() as tg:
            for url in urls:
                if shutdown.is_set():
                    break
                tg.create_task(fetch_bounded(sem, client, url))

async def main() -> None:
    shutdown = asyncio.Event()
    loop = asyncio.get_running_loop()
    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, shutdown.set)
    start = time.perf_counter()
    with suppress(asyncio.CancelledError):
        await run_pipeline(urls, shutdown)
    print(f"Hotovo za {time.perf_counter() - start:.2f}s")

asyncio.run(main())

Pro metriky doporučuju prometheus-client s Histogram pro fetch latenci a Counter pro HTTP kódy, export přes malý FastAPI endpoint, který k tomu servisu jako side-car přidáte. Pokud servisujete přes FastAPI, můžete se inspirovat článkem Nasazení ML modelů s FastAPI, kde ukazuju stejný pattern pro model-serving.

Poslední tip z praxe: async pipeline v Kubernetes běží nejlépe jako Deployment s preStop hookem, který nastaví shutdown flag a počká 15 sekund. Bez toho pod dostane SIGKILL a osiřelé requesty visí na targetu, který si nezvykl na disconnect uprostřed chunked odpovědi. Kdo neposkytne graceful shutdown, ten platí za rozbité idempotence klíče. Já jsem to loni řešil dvakrát, podruhé už s alertem v Grafaně.

Často kladené otázky

Jak zrychlit ETL pipeline v Pythonu?

Nejprve zjistěte, kde je úzké hrdlo. Pokud čekáte na síť (HTTP, DB, S3), přejděte na async I/O s httpx a asyncio, typicky 5–20× zrychlení. Pokud je pomalá CPU část (parse, transformace), použijte Polars nebo multiprocessing. Nejhorší je slepě přepisovat do asyncio kód, který je pomalý kvůli neefektivnímu SQL nebo N+1 dotazu.

Kdy použít asyncio místo threadingu?

Použijte asyncio, když máte tisíce souběžných I/O operací (HTTP calls, databázové cursory) a knihovny nabízející async API (httpx, asyncpg, aiofiles). Threading zvolte, když jste zamčeni v synchronní knihovně bez async varianty (starý SDK) a stačí vám nižší souběžnost do ~500 vláken. Pro CPU-heavy práci ani jedno, jděte do multiprocessing nebo Ray.

Jak omezit počet souběžných requestů v asyncio?

Použijte asyncio.Semaphore(N) a obalte fetch coroutine do async with sem:. Pro rate limiting v req/s přidejte aiolimiter.AsyncLimiter(max_rate=100, time_period=1). Semaphore hlídá souběžnost (kolik naráz), limiter hlídá rychlost (kolik za sekundu). V praxi potřebujete oba, jinak buď přetížíte target, nebo se zasekne pipeline.

Jaký je rozdíl mezi httpx a aiohttp?

httpx podporuje sync i async se stejným API a má nativní HTTP/2, aiohttp je jen async a HTTP/2 ho čeká teprve ve verzi 4.0. Pro nový projekt vždycky httpx, protože migrace z requests je snadná, HTTP/2 multiplexing snižuje latenci a dokumentace je čitelnější. Aiohttp má výhodu jen v extrémně vysokém průtoku (100k+ req/s), kde je 10–15 % rychlejší v per-request overheadu.

Můžu volat synchronní pandas z async funkce?

Můžete, ale zablokujete celý event loop. Vždy obalte volání do await asyncio.to_thread(pd.read_parquet, path), které pandas spustí v thread poolu a event loop mezitím obsluhuje ostatní úlohy. Ještě lepší je nechat pandas až na konci pipeline (po fetch fázi) a spustit ho jako samostatný synchronní blok, držte async a sync části jasně oddělené.

Jak řešit chyby v asyncio.TaskGroup?

TaskGroup při jakékoli výjimce zruší všechny sourozence a propaguje ExceptionGroup ven, což je bezpečné, ale znamená, že celou dávku ztratíte. Pro tolerantní pipeline zachyťte výjimky uvnitř coroutine (například vraťte None místo vyhození) a v hlavní smyčce filtrujte úspěšné výsledky. Nikdy nepolykejte výjimky bez logování, v produkci se ztratí signál o rozbité API endpoint.

Tomás Oliveira
O Autorovi Tomás Oliveira

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