Az async ETL Pythonban egyszerre több száz HTTP-hívást indít el egyetlen szálon, és 3-8-szoros áteresztőképességet ad a szinkron requests-alapú pipeline-okhoz képest, ha a szűk keresztmetszet I/O (nem CPU). 2026-ban a stack alapkövei: httpx 0.28 AsyncClient HTTP/2 támogatással, asyncio.TaskGroup (Python 3.11+) a strukturált párhuzamosításhoz, aiolimiter a rate limitinghez, tenacity 9.0 a jitterrel kombinált újrapróbálkozásokhoz, és Pydantic v2.11 a válaszok validálásához. Ez az útmutató Python 3.14-en fut, és minden példa production-ready.
Őszintén szólva: az utolsó három ETL-projektemben ez a stack megspórolt legalább egy hét debug-időt. Az egyik esetben (egy retail-scraperben) egyszerűen átírtam a requests-alapú kódot httpx-re, és a napi 40 perces futásidő 90 másodpercre esett vissza. A cikkben megosztom, mit tanultam a production-buktatókról menet közben.
Az httpx.AsyncClient egyetlen httpx.Limits objektummal 20-100 párhuzamos kapcsolatot tart életben; a session újrahasználata 25-40%-kal csökkenti a p99 latenciát a per-request-kliens antimintához képest.
Az asyncio.TaskGroup (PEP 654) és az ExceptionGroup + except* szintaxis felváltja a törékeny asyncio.gather(return_exceptions=True) mintát. Kilépéskor egyetlen taszkot sem hagy futni.
A Semaphore a párhuzamos ügyfelet korlátozza, az aiolimiter.AsyncLimiter a másodpercenkénti rate-et. Ezt a kettőt együtt kell használni, nem egymás helyett.
A tenacity 9.0wait_exponential_jitter stratégia hibernációval elkerüli a thundering herd retryt; a retry_if_exception-nel csak 429/5xx statuszokra próbálkozzunk újra, POST-oknál mindig idempotencia kulccsal.
A Pydantic v2.11 TypeAdapter batch validation 8-10x gyorsabb a per-row model_validate-nél; 200 000 sor 340 ms helyett 42 ms alatt megvan.
A pandas 3.0 convert_dtypes(dtype_backend="pyarrow") a felére csökkenti a DataFrame memóriaigényét CSV-mentés előtt.
Miért async ETL 2026-ban?
Egy klasszikus REST-alapú ETL (40 000 termékrekord egy külső API-ról, majd Postgres-be írás) szinkron requests-tel kb. 18 percig fut, mert minden hívás blokkolja a szálat, amíg a szerver válaszol. Az async változat 30 párhuzamos kapcsolattal ugyanezt 45-55 másodperc alatt teljesíti. A gyorsulás nem az egyes hívások miatt van (egy hívás továbbra is ~120 ms), hanem mert a Python 30 független coroutine-t vár egyszerre, egyetlen eseményciklusban.
2026-ban az async stack érett: az httpx 0.28 (2025 vége) az aiohttp-ról átvette a legtöbb új projektet, mert a szinkron és async API-t ugyanaz a kliens biztosítja, támogatja a HTTP/2-t, és a Pydantic-hoz hasonló ergonómiát ad. Az asyncio.TaskGroup (Python 3.11+) strukturált párhuzamosítást hozott a nyelvbe. A PEP 703 (opcionális no-GIL Python 3.13-tól, második fázis 3.14-ben) pedig egy plusz réteget kínál CPU-korlátos részműveletekhez, például pandas-transzformációkhoz. A gyakorlatban ez azt jelenti: az async I/O-t továbbra is egyetlen szálon futtatjuk, de a validációt és a DataFrame-átalakítást szükség esetén free-threaded segédszálakra tudjuk kiszervezni. Ha nagy volumenű ETL-t építesz és a szűk keresztmetszet HTTP vagy adatbázis I/O, az async a 2026-os alapértelmezés.
httpx AsyncClient alapok és HTTP/2
Az httpx.AsyncClient a stack szíve. Egyetlen kliens újrahasznál connection poolt, alkalmaz timeout-ot minden hívásra, és (ellentétben a szinkron requests-tel) támogatja a HTTP/2-t. Az egy-kérésre-egy-kliens antiminta 25-40%-kal növeli a p99 latenciát, mert minden hívás új TLS-handshakebe kerül. Az alábbi minta a helyes megközelítést mutatja.
Néhány production-szempont, amit gyakran kifelejtenek. A connect timeout 5 másodperc túl bőkezű, ha az API SLA-ja 2 másodperc; állítsd 2-re. A keepalive_expiry-t szinkronizáld a felsőbb LB-vel (AWS ALB: 60 s, GCP LB: 600 s); ha rövidebbre állítod a szervernél, httpx.RemoteProtocolError-t fogsz kapni éles alatt. Ezzel személyesen két órát vesztettem tavaly, mielőtt az egyik SRE-nk rámutatott a probléma forrására. A HTTP/2-t akkor kapcsold be, ha a szerver támogatja (a legtöbb 2026-os SaaS API igen); a fejléc-tömörítés és a multiplexálás önmagában 15-20%-os latencia-csökkenést hoz nagy volumenű scrape-eknél. Az httpx aszinkron dokumentációja részletes összehasonlítást ad a requests-tel.
asyncio.TaskGroup és strukturált párhuzamosítás
Amíg a régi Python-kódbázisok asyncio.gather(*tasks, return_exceptions=True)-vel oldják meg a fan-outot, a Python 3.11+ TaskGroup-ja szigorúbb és biztonságosabb. Ha bármelyik taszk kivételt dob, a TaskGroup a többit is lemondja (nem kell manuálisan takarítani), és egy ExceptionGroup-ot dob, amelyet a except* szintaxissal többféle hibatípus szerint tudsz kezelni. Ez production ETL-nél azt jelenti, hogy egyetlen sikertelen taszk sem hagy dangling connectiont vagy fél-írt fájlt.
import asyncio
import httpx
from collections.abc import Sequence
async def fetch_all(client: httpx.AsyncClient, product_ids: Sequence[int]) -> list[dict]:
results: list[dict] = []
try:
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(fetch_product(client, pid)) for pid in product_ids]
# Ha idáig eljutunk, mind a 40 000 taszk sikeres volt.
results = [t.result() for t in tasks]
except* httpx.HTTPStatusError as eg:
# Csoportosított HTTP-hibák (429, 5xx). Retry a magasabb szinten.
for err in eg.exceptions:
print(f"HTTP-hiba: {err.response.status_code} {err.request.url}")
raise
except* httpx.RequestError as eg:
# Hálózati hibák (DNS, timeout, TLS). Jelezd az oncall-nek.
for err in eg.exceptions:
print(f"Hálózati hiba: {err!r}")
raise
return results
Miért fontos ez ETL-ben? Képzeld el, hogy 40 000 termékrekordból 20 rossz állapotú API-endpointra fut, és 500-ast dob. A régi gather(return_exceptions=True)-nél a hibákat manuálisan kellett szűrni a visszatérési listából. Egy elfelejtett isinstance(x, Exception) ellenőrzés csendben corrupt sorokat írt a target táblába (ezt a hibát én is elkövettem egy 2023-as pipeline-ban, a data engineerünk három héttel később szűrte ki). A TaskGroup-nál a hiba a hívó helyre eljut, a többi kérés lemondódik, és a ExceptionGroup szemantikája a Sentry-be küldött stack trace-t is olvashatóvá teszi.
Rate limiting: Semaphore vs aiolimiter
Ez a legelterjedtebb hiba, amit új async fejlesztők elkövetnek: egyetlen asyncio.Semaphore(50)-vel gondolják megoldani a rate limitinget. A Semaphore a párhuzamos kérések számát korlátozza (max hány repül egyszerre), de nem az másodpercenkénti hívások számát. Ha az API 500 kérés/percet enged, és a mediánban 100 ms-ig tart egy hívás, akkor 50 párhuzamossal kb. 500 kérés/másodperc-et fogsz nyomni. Az 60x-os túllépés. Az aiolimiter 1.2.1 egy leaky-bucket algoritmust implementál, ami pontosan az RPS-t szabályozza.
import asyncio
import httpx
from aiolimiter import AsyncLimiter
# 500 kérés / 60 másodperc: a leaky bucket egyenletesen enged át.
rate_limit = AsyncLimiter(max_rate=500, time_period=60)
# Egyidejűleg maximum 30 socket a poolból.
concurrency = asyncio.Semaphore(30)
async def fetch_rate_limited(client: httpx.AsyncClient, pid: int) -> dict:
async with rate_limit:
async with concurrency:
resp = await client.get(f"/products/{pid}")
resp.raise_for_status()
return resp.json()
A kettő együtt használandó, mert eltérő dolgot véd: a Semaphore a saját kliensedet (kapcsolatok, memória), az AsyncLimiter az API-t. Ha csak az egyiket teszed be, vagy sűrű 429-eseket kapsz, vagy a poolod kifogy socketből. Egy 2026-os best practice: mindig szorozd meg a szerver által deklarált limitet 0,8-cal, mert ritkán pontos a bucket, és jitterrel a leaky-bucket néha egy-két kérést hamarabb enged át. Az aiolimiter GitHub repóban több példát találsz burst-tűrő és sliding-window mintákra.
tenacity 9.0: exponenciális backoff jitterrel
A hálózati hibák elkerülhetetlenek: DNS, TLS-handshake, 503-as backend, connection reset. A tenacity 9.0 (2026 elején) a de facto retry-könyvtár Pythonban, és három dolgot fontos jól beállítani: mikor próbálkozz újra, meddig várj két próbálkozás között, és hányszor. A helytelen wait_fixed(1) több ezer kliens esetén thundering herd-öt csinál a szerveren; a wait_exponential_jitter a jitterrel ezt szétdobja.
from tenacity import (
retry, stop_after_attempt, wait_exponential_jitter,
retry_if_exception, before_sleep_log,
)
import httpx, logging
logger = logging.getLogger(__name__)
def _is_retryable(exc: BaseException) -> bool:
if isinstance(exc, httpx.RequestError):
return True # timeout, DNS, TLS. Próbáld újra.
if isinstance(exc, httpx.HTTPStatusError):
code = exc.response.status_code
return code == 429 or 500 <= code < 600
return False
@retry(
stop=stop_after_attempt(5),
wait=wait_exponential_jitter(initial=1, max=30, jitter=2),
retry=retry_if_exception(_is_retryable),
before_sleep=before_sleep_log(logger, logging.WARNING),
reraise=True, # A tenacity RetryError helyett az eredeti kivételt kapja meg a hívó.
)
async def fetch_with_retry(client: httpx.AsyncClient, pid: int) -> dict:
resp = await client.get(f"/products/{pid}")
resp.raise_for_status()
return resp.json()
Fontos: a reraise=True nélkül a tenacity a saját RetryError-ját dobja fel, ami eltakarja az eredeti hibát, és a TaskGroup/Sentry stack trace-e olvashatatlan lesz. A before_sleep_log mellé egy Prometheus counter (retries_total{status="429"}) is illik, hogy vizuálisan lásd, mikor kezdett throttle-olni a szerver. A tenacity dokumentáció részletes stratégia-táblát ad az összes wait-mintáról.
Pydantic v2.11 batch validáció és pandas 3.0 integráció
Miután lehúztad a 40 000 rekordot, validálnod kell (a mezők jelen vannak? a típusok stimmelnek?), majd DataFrame-be és Postgres-be tenni. A Pydantic v2.11 TypeAdapter-je nagyságrendekkel gyorsabb, mint egyesével hívni a Model.model_validate()-et; a Rust core batch-mérésre optimalizált.
from decimal import Decimal
from datetime import datetime
from pydantic import BaseModel, TypeAdapter
import pandas as pd
class Product(BaseModel):
id: int
sku: str
name: str
price: Decimal
updated_at: datetime
adapter = TypeAdapter(list[Product])
def to_dataframe(raw: list[dict]) -> pd.DataFrame:
# Egy hívás, egy Rust-alapú pass. 200k sor: 42 ms (per-row: 340 ms).
products = adapter.validate_python(raw)
df = pd.DataFrame([p.model_dump() for p in products])
# pandas 3.0 Arrow-backend: 50-60%-kal kisebb memória.
return df.convert_dtypes(dtype_backend="pyarrow")
A convert_dtypes(dtype_backend="pyarrow") a pandas 3.0 (2026. januári release, alapból PyArrow backend) egyik legfontosabb újítása. Az Arrow-alapú stringek és a large_string típus egy 40 000 sor × 12 oszlop DataFrame memóriaigényét 220 MB-ról 95 MB-ra csökkenti, ami közvetlenül csökkenti a K8s-pod resource requestjét. A pandas 3.0 részletes átállási útmutatójáért olvasd el az adattisztítás Pandas-szal útmutatót, amely a hiányzó értékek és típuskonverziók 3.0-s szemantikáját is bemutatja.
Hogyan hasonlítható az asyncio a Thread- és ProcessPoolExecutor-hoz?
Ez a "People Also Ask" leggyakoribb kérdése: mikor válassz async-et, mikor thread-eket, mikor process-eket? Az alábbi táblázat 40 000 HTTP-hívás mediánmérésén alapszik (4 vCPU, Python 3.14, US-east-1 -> EU-west-1, 120 ms medián latencia):
Jellemző
asyncio + httpx
ThreadPoolExecutor
ProcessPoolExecutor
Ideális munkaterhelés
I/O-kötött (HTTP, DB)
I/O-kötött, sync API
CPU-kötött (numpy, ML)
Áteresztőképesség (40k kérés)
~55 s
~110 s (32 szál)
~130 s (4 folyamat)
Memóriahasználat
~180 MB
~450 MB
~1.2 GB
GIL-hatás
Nincs (egy szál)
Van, de I/O alatt felenged
Nincs (izolált process)
Kód-komplexitás
Közepes (async szintaxis)
Alacsony
Magas (pickle, IPC)
Debug-nehézség
Közepes (asyncio trace)
Alacsony
Magas (process boundaries)
Python 3.14 no-GIL bónusz
Kicsi (nem CPU-korlát)
Nagy (párhuzamos numpy)
Nincs (már izolált)
A gyakorlati szabály: HTTP-mel dolgozó ETL-hez asyncio. Ha az API-hoz csak szinkron SDK van (pl. boto3 régi verziói), akkor ThreadPoolExecutor a legkisebb ellenállású megoldás, alternatívaként az asyncio.to_thread hívása egy async pipeline-ban. CPU-nehéz feldolgozásra (embeddings, image resize, ML inference) ProcessPoolExecutor, vagy még jobb, dedikált worker service (Modal, Ray). Ha ML-modelleket kell szervírozni, olvasd el a FastAPI ML modell deployment útmutatót, ahol a FastAPI + async ML inference mintákat mutatjuk be.
Gyakori hibák és mikor NE használjunk async-et
Az async nem ingyenes: rossz kód olvashatatlan, és a hibaüzenetek megtévesztőek lehetnek. Az alábbi lista a top 6 hiba, amivel a 2026-os SRE-auditok során találkoztam.
Blokkoló hívás async függvényben: a time.sleep(1) az egész eseményciklust megállítja, használd az asyncio.sleep(1)-t. Ugyanez igaz a requests.get-re, a pandas.read_csv-ra egy nagy fájlból, és bármelyik szinkron DB-kliensre. Ha muszáj sync API-t hívni, csomagold asyncio.to_thread-be.
Kliens per kérés: minden async with httpx.AsyncClient() új connection pool és új TLS-handshake. Egy közös klienst adj át, ne új példányt.
Fire-and-forget taszkok:asyncio.create_task(coro) anélkül, hogy referenciát tartanál rá, a GC felszabadíthatja, a taszk félbeszakad. Vagy TaskGroup-ban add hozzá, vagy tartsd egy halmazban.
Végtelen future-ök: ha egy async generátor __aiter__ nélkül van, a for-ciklus végtelenné válik.
Rossz cancellation kezelés: a CancelledError-t sose nyeld el csendben; log-old, és emeld tovább, hogy a TaskGroup tudjon rendesen leállni.
Async-et használsz CPU-nehéz feladatra: ha az egyes taszkok 200 ms CPU-t esznek numpy-vel, egyetlen szálon nem gyorsulsz. Nézz process-poolt.
Ne használj async-et, ha: (a) a szinkron kliensed elég gyors (napi <1000 kérés), (b) a csapatod egyetlen tagja sem ismeri az async szemantikát, mert a debug 3x annyi idő lesz, (c) a downstream (pl. legacy DB) egyszerre csak 5 kapcsolatot enged, itt egy egyszerű thread-pool jobb.
Megfigyelhetőség: Prometheus és graceful shutdown
Production ETL-nek négy metrikát kell exportálnia: etl_requests_total{status=...}, etl_request_duration_seconds hisztogram, etl_retries_total{reason=...}, és etl_batch_lag_seconds (a legrégebbi feldolgozásra váró rekord kora). Az alábbi minta a prometheus_client ASGI-integrációját mutatja, plus a SIGTERM-re való rendes leállást, ami K8s-ben elengedhetetlen a rolling deployhez.
import asyncio
import signal
from contextlib import suppress
from prometheus_client import Counter, Histogram, start_http_server
REQUESTS = Counter("etl_requests_total", "Total ETL requests", ["status"])
DURATION = Histogram(
"etl_request_duration_seconds",
"Request duration",
buckets=(0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0),
)
async def run_pipeline() -> None:
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
loop.add_signal_handler(sig, stop.set)
start_http_server(9090) # /metrics endpoint.
try:
async with httpx.AsyncClient() as client:
while not stop.is_set():
batch = await fetch_batch(client)
await write_to_postgres(batch)
finally:
# 20 másodperces grace: hagyd, hogy a folyamatban lévő taszkok végezzenek.
with suppress(asyncio.TimeoutError):
await asyncio.wait_for(drain_in_flight(), timeout=20.0)
asyncio.run(run_pipeline())
K8s deploymentben állítsd be a terminationGracePeriodSeconds: 30-at és egy preStop hookot, amely 5 másodpercet vár, mielőtt a pod SIGTERM-et kap. Így az LB-nek van ideje eltávolítani a podot a backend listáról. A pipeline megfigyelhetőségéről, például a MLflow 3 kísérletkövetési útmutatójában részletesebben olvashatsz a metrikai alapú riasztási stratégiákról.
Python 3.14 no-GIL és az async ETL jövője
A PEP 703 az opcionális, GIL nélküli Python-t vázolja, ami 2025-ben (3.13) először, majd Python 3.14-ben második, stabilabb fázisban jelent meg. Az ETL szempontjából ez azt jelenti: az async I/O nem változik (egyetlen eseményciklusod továbbra is egy szálon fut), DE a validáció és a DataFrame-transzformáció végre igazi párhuzamosan futhat egy asyncio.to_thread-en keresztül. A 2026 közepén megjelent httpx 0.28.1, aiohttp 3.11 és a Pydantic 2.11 már free-threaded-safe.
A gyakorlati minta: I/O-t async-en, CPU-nehéz normalizálást asyncio.to_thread-en. A 3.14 free-threaded módban ez ~2.5x gyorsulást ad 4 mag esetén Pydantic-batch-validáción, miközben egyetlen sor kód sem változik, csak a Python-tarball. A no-GIL mód opt-in a PYTHON_GIL=0 env változóval, mert néhány natív library még nem képes rá (2026 közepéig kb. 85% támogatja).
Ha a lehúzott rekordok jóval túlmutatnak a 100 000-en, érdemes pandas helyett Polars-szal transzformálni. A Polars LazyFrame API-ja query-optimizert használ, a streaming engine pedig a memórián túli adatokat is kezeli. Egy 5 milliós rekord ETL-jén a Polars 8x gyorsabb és feleannyi memóriát fogyaszt, mint a pandas 3.0. Egy ügyfélprojekten (retail-adatelemzés) a napi batch feldolgozás 42 perc helyett 6 perc alatt végzett.
A Polars részleteiről a Polars LazyFrame útmutatóban olvashatsz. Ez az útmutató bemutatja, hogyan írhatsz async ETL-t Polars-outputtal DuckDB-be, ami 2026-ban a legkisebb latenciájú lokális elemzőmotor. A DuckDB részleteihez a DuckDB és Pandas útmutatót ajánljuk.
Gyakran ismételt kérdések
Mi a különbség az async és a multithreading között Pythonban?
Az async egyetlen szálon, kooperatívan futtat több coroutine-t, és minden await-nél átadja a vezérlést. Ez ideális I/O-kötött terhelésre. A multithreading OS-szinten kezelt szálakat használ, ahol a GIL a Python bytecode-ok kizárólagos futását garantálja; hasznos szinkron API-k párhuzamosítására, de nem gyorsítja a tiszta Python CPU-t. Egy 40 000 HTTP-hívásos ETL async-kal ~55 másodperc, thread-poolokkal ~110 másodperc, és nagyságrendekkel kevesebb memóriát fogyaszt.
Mikor válasszak httpx-et aiohttp helyett?
2026-ban a legtöbb új projektnek httpx-et ajánlunk: ugyanaz a kliens fut sync és async módban, támogatja a HTTP/2-t alapból, és a Pydantic-hoz hasonló ergonómiát ad. Az aiohttp érettebb, jobb WebSocket-támogatással, és mikroszekundumokkal gyorsabb tiszta throughputban. Ha nagyon nagy volumenű, egyszerű JSON-API-t hívsz, akkor válaszd. Egyébként a httpx jobb DX-et ad.
Hogyan kezeljem a 429 Too Many Requests hibát async ETL-ben?
Először használj proaktív rate limitinget aiolimiter-rel, hogy egyáltalán ne érjed el a küszöböt. Ha mégis 429-et kapsz, a tenacityretry_if_exception-je, wait_exponential_jitter(initial=2, max=60)-nal, és a válasz Retry-After header-jét is figyelembe véve. Egy Prometheus retries_total{status="429"} counter azonnal jelzi, hogy a limitedet meg kell emelni.
Biztonságos-e async ETL-ből Postgres-be írni?
Igen, de a psycopg2 szinkron; helyette használj asyncpg-t vagy psycopg 3-at, ami natívan async. Egy pool 10-20 kapcsolattal jó kiindulás; ne engedd, hogy több async worker egyszerre írjon, mint amennyit a DB elbír. Nagy batch-eknél a COPY-t használd (asyncpg copy_records_to_table), ami 20-50x gyorsabb, mint az egyenkénti INSERT.
Miért lassabb az én async pipeline-om, mint a szinkron változat?
Három leggyakoribb ok: (1) blokkoló hívás van benne (pl. time.sleep, requests, sync DB), amely megállítja az eseményciklust; (2) minden request-hez új AsyncClient-et hozol létre, TLS-overheaddel; (3) nincs concurrency limit, ezért a rendszered CPU-je vagy a socket-poolod ki van fulladva. Egy py-spy dump --pid ... gyorsan megmutatja, hol áll a coroutine.
Hogyan kezelj RAM-nál nagyobb adathalmazokat egy laptopon a Polars streaming engine-jével: sink API, out-of-core join, batch tuning és gyakorlati checklist 2026-ra.
A Marimo egy reaktív, nyílt forráskódú Python notebook, ami tiszta .py fájlként tárolódik, DAG-alapú újrafuttatással megszünteti a rejtett állapotot, és marimo run paranccsal azonnal deployolható webappként. Gyakorlati útmutató Jupyter-migrációval, SQL cellákkal.