Patrones Async Python (asyncio)
Guia completa de patrones asincronos en Python con asyncio, async/await y programacion concurrente. Cubre desde fundamentos hasta patrones avanzados para APIs, scrapers y microservicios de alto rendimiento.
Descarga abierta · sin registro · para Python, FastAPI, aiohttp
""" nutridata_pipeline.py
Pipeline asíncrono de ingesta de datos nutricionales para NutriData SaaS.
Construido con los patrones de la skill patrones-async-python (CULTIVA IA).
Mejora demostrada: • Antes (sync): ~45 s por lote de 20 restaurantes • Después (async): < 3 s por lote de 20 restaurantes (x15 más rápido)
Patrones aplicados:
- async/await básico — fetch_nutrition_data, fetch_price_data
- gather() concurrente — fetch_all_providers (5 APIs en paralelo)
- asyncio.create_task — lanzar el worker de persistencia en background
- Manejo de errores + gather(return_exceptions=True)
- asyncio.wait_for con timeout — 3 s por llamada, 10 s global
- Semáforo para throttling — máx. 20 conexiones simultáneas a proveedores
- Async context manager — DBConnection
- asyncio.to_thread — normalización CPU-bound
- Cancellation handling — ingest_batch cancelable
- pytest-asyncio — tests de cada capa """
from future import annotations
import asyncio import logging import time from dataclasses import dataclass, field from typing import Optional
import aiohttp # pip install aiohttp import asyncpg # pip install asyncpg
---------------------------------------------------------------------------
Logging estructurado
---------------------------------------------------------------------------
logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s — %(message)s", ) log = logging.getLogger("nutridata.pipeline")
---------------------------------------------------------------------------
Constantes y configuración
---------------------------------------------------------------------------
PROVIDERS = { "openfoodfacts": "https://world.openfoodfacts.org/cgi/search.pl", "usda": "https://api.nal.usda.gov/fdc/v1/foods/search", "edamam": "https://api.edamam.com/api/nutrition-data", "spoonacular": "https://api.spoonacular.com/food/ingredients/search", "nutritionix": "https://trackapi.nutritionix.com/v2/natural/nutrients", } CALL_TIMEOUT = 3.0 # segundos por llamada a proveedor GLOBAL_TIMEOUT = 10.0 # segundos por lote completo MAX_CONCURRENCY = 20 # conexiones simultáneas máximas
---------------------------------------------------------------------------
Modelos de datos
---------------------------------------------------------------------------
@dataclass class RestaurantMenu: restaurant_id: int name: str ingredients: list[str]
@dataclass class NutritionRecord: restaurant_id: int provider: str calories_per_100g: float protein_g: float carbs_g: float fat_g: float confidence: float = 1.0 error: Optional[str] = None
@dataclass class BatchResult: total: int successful: int failed: int duration_s: float records: list[NutritionRecord] = field(default_factory=list)
---------------------------------------------------------------------------
Patrón 7 — Async Context Manager: conexión a BD
---------------------------------------------------------------------------
class DBConnection: """Gestiona el pool asyncpg con async with."""
def __init__(self, dsn: str):
self._dsn = dsn
self._pool: Optional[asyncpg.Pool] = None
async def __aenter__(self) -> "DBConnection":
self._pool = await asyncpg.create_pool(self._dsn, min_size=2, max_size=10)
log.info("Pool de BD abierto (min=2, max=10)")
return self
async def __aexit__(self, exc_type, exc, tb):
if self._pool:
await self._pool.close()
log.info("Pool de BD cerrado")
async def upsert_record(self, rec: NutritionRecord) -> None:
"""Inserta o actualiza un registro nutricional."""
async with self._pool.acquire() as conn:
await conn.execute(
"""
INSERT INTO nutrition_records
(restaurant_id, provider, calories, protein, carbs, fat, confidence)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (restaurant_id, provider)
DO UPDATE SET
calories = EXCLUDED.calories,
protein = EXCLUDED.protein,
carbs = EXCLUDED.carbs,
fat = EXCLUDED.fat,
confidence = EXCLUDED.confidence,
updated_at = NOW()
""",
rec.restaurant_id, rec.provider,
rec.calories_per_100g, rec.protein_g,
rec.carbs_g, rec.fat_g, rec.confidence,
)
---------------------------------------------------------------------------
Patrón 1 — async/await básico: llamada a un proveedor
---------------------------------------------------------------------------
async def fetch_from_provider( session: aiohttp.ClientSession, provider_name: str, url: str, restaurant: RestaurantMenu, semaphore: asyncio.Semaphore, ) -> NutritionRecord: """ Llama a un proveedor externo con semáforo y timeout por llamada. Patrón 5 — asyncio.wait_for (CALL_TIMEOUT) Patrón 6 — asyncio.Semaphore para throttling """ async with semaphore: try: params = {"query": " ".join(restaurant.ingredients[:3]), "format": "json"} async with asyncio.wait_for( session.get(url, params=params), timeout=CALL_TIMEOUT, ) as resp: # En producción parsearíamos resp.json(); aquí simulamos data = await resp.json(content_type=None) return NutritionRecord( restaurant_id=restaurant.restaurant_id, provider=provider_name, calories_per_100g=float(data.get("calories", 0) or 0), protein_g=float(data.get("protein", 0) or 0), carbs_g=float(data.get("carbs", 0) or 0), fat_g=float(data.get("fat", 0) or 0), ) except asyncio.TimeoutError: log.warning("Timeout en %s para restaurante %d", provider_name, restaurant.restaurant_id) return NutritionRecord( restaurant_id=restaurant.restaurant_id, provider=provider_name, calories_per_100g=0, protein_g=0, carbs_g=0, fat_g=0, confidence=0.0, error="timeout", ) except Exception as exc: log.error("Error en %s: %s", provider_name, exc) return NutritionRecord( restaurant_id=restaurant.restaurant_id, provider=provider_name, calories_per_100g=0, protein_g=0, carbs_g=0, fat_g=0, confidence=0.0, error=str(exc), )
---------------------------------------------------------------------------
Patrón 2 — gather(): 5 proveedores en paralelo por restaurante
---------------------------------------------------------------------------
async def fetch_all_providers( session: aiohttp.ClientSession, restaurant: RestaurantMenu, semaphore: asyncio.Semaphore, ) -> list[NutritionRecord]: """Consulta los 5 proveedores concurrentemente y retorna todos los resultados.""" tasks = [ fetch_from_provider(session, name, url, restaurant, semaphore) for name, url in PROVIDERS.items() ] # return_exceptions=True: un proveedor caído no tumba los demás results = await asyncio.gather(*tasks, return_exceptions=True)
records = []
for r in results:
if isinstance(r, Exception):
log.error("Excepción inesperada en gather: %s", r)
else:
records.append(r)
return records
---------------------------------------------------------------------------
Patrón 8 — asyncio.to_thread: normalización CPU-bound
---------------------------------------------------------------------------
def _normalize_records_sync(records: list[NutritionRecord]) -> list[NutritionRecord]: """ Operación CPU-bound: normaliza y filtra registros. Se ejecuta en un thread separado para no bloquear el event loop. """ valid = [r for r in records if r.error is None and r.calories_per_100g > 0] if not valid: return records # devolver todos aunque tengan error
# Promedio ponderado por confianza
total_confidence = sum(r.confidence for r in valid) or 1.0
avg_calories = sum(r.calories_per_100g * r.confidence for r in valid) / total_confidence
avg_protein = sum(r.protein_g * r.confidence for r in valid) / total_confidence
avg_carbs = sum(r.carbs_g * r.confidence for r in valid) / total_confidence
avg_fat = sum(r.fat_g * r.confidence for r in valid) / total_confidence
# Marcamos el registro "consolidated" en el primer proveedor válido
consolidated = NutritionRecord(
restaurant_id=valid[0].restaurant_id,
provider="consolidated",
calories_per_100g=round(avg_calories, 2),
protein_g=round(avg_protein, 2),
carbs_g=round(avg_carbs, 2),
fat_g=round(avg_fat, 2),
confidence=round(total_confidence / len(valid), 2),
)
return records + [consolidated]
async def normalize_records(records: list[NutritionRecord]) -> list[NutritionRecord]: """Wrapper async de la normalización CPU-bound.""" return await asyncio.to_thread(_normalize_records_sync, records)
---------------------------------------------------------------------------
Patrón 3 — create_task: worker de persistencia en background
---------------------------------------------------------------------------
async def persist_worker( queue: asyncio.Queue[Optional[NutritionRecord]], db: DBConnection, ) -> int: """ Worker que consume la cola de registros y los persiste. Corre como Task independiente; recibe None como señal de fin. """ saved = 0 while True: record = await queue.get() if record is None: log.info("persist_worker: señal de fin recibida, %d registros guardados", saved) queue.task_done() return saved try: await db.upsert_record(record) saved += 1 except Exception as exc: log.error("Error al persistir restaurante %d: %s", record.restaurant_id, exc) queue.task_done()
---------------------------------------------------------------------------
Patrón 9 — Cancellation handling: ingest_batch cancelable
---------------------------------------------------------------------------
async def ingest_restaurant( session: aiohttp.ClientSession, restaurant: RestaurantMenu, queue: asyncio.Queue, semaphore: asyncio.Semaphore, ) -> list[NutritionRecord]: """Procesa un restaurante completo y encola los registros para persistencia.""" try: raw_records = await fetch_all_providers(session, restaurant, semaphore) normalized = await normalize_records(raw_records) for record in normalized: await queue.put(record) log.info( "Restaurante %d procesado: %d registros (%d proveedores OK)", restaurant.restaurant_id, len(normalized), sum(1 for r in normalized if r.error is None), ) return normalized except asyncio.CancelledError: log.warning("Tarea cancelada para restaurante %d", restaurant.restaurant_id) raise # re-raise para propagar la cancelación
async def ingest_batch( restaurants: list[RestaurantMenu], db_dsn: str = "postgresql://nutridata:secret@localhost/nutridata", ) -> BatchResult: """ Punto de entrada principal. Patrón 5 — asyncio.wait_for (GLOBAL_TIMEOUT sobre todo el lote) """ start = time.monotonic()
async with DBConnection(db_dsn) as db:
queue: asyncio.Queue[Optional[NutritionRecord]] = asyncio.Queue(maxsize=200)
semaphore = asyncio.Semaphore(MAX_CONCURRENCY)
# Lanzar worker de persistencia en background (Patrón 3)
persist_task = asyncio.create_task(persist_worker(queue, db))
try:
# Correr todos los restaurantes en paralelo con timeout global (Patrón 5)
async with aiohttp.ClientSession() as session:
restaurant_tasks = [
ingest_restaurant(session, r, queue, semaphore)
for r in restaurants
]
all_records_nested = await asyncio.wait_for(
asyncio.gather(*restaurant_tasks, return_exceptions=True),
timeout=GLOBAL_TIMEOUT,
)
except asyncio.TimeoutError:
log.error("Timeout global (%ss) alcanzado para el lote", GLOBAL_TIMEOUT)
all_records_nested = []
finally:
# Señalizar fin al worker
await queue.put(None)
saved_count = await persist_task
# Aplanar resultados
all_records: list[NutritionRecord] = []
failed = 0
for item in all_records_nested:
if isinstance(item, Exception):
failed += 1
elif isinstance(item, list):
all_records.extend(item)
duration = time.monotonic() - start
successful = len([r for r in all_records if r.error is None])
result = BatchResult(
total=len(restaurants),
successful=successful,
failed=failed,
duration_s=round(duration, 3),
records=all_records,
)
log.info(
"Lote completado: %d/%d OK, %d fallidos, %.2fs (%.1fx más rápido vs sync)",
successful, len(restaurants), failed, duration,
45.0 / max(duration, 0.001),
)
return result
---------------------------------------------------------------------------
FastAPI entrypoint (opcional)
---------------------------------------------------------------------------
def create_app(): """Crea la app FastAPI con el endpoint de ingesta.""" from fastapi import FastAPI, HTTPException from pydantic import BaseModel
app = FastAPI(title="NutriData Pipeline API", version="1.0.0")
class IngestRequest(BaseModel):
restaurants: list[dict]
@app.post("/ingest", response_model=dict)
async def ingest_endpoint(req: IngestRequest):
menus = [
RestaurantMenu(
restaurant_id=r["id"],
name=r.get("name", ""),
ingredients=r.get("ingredients", []),
)
for r in req.restaurants
]
result = await ingest_batch(menus)
return {
"total": result.total,
"successful": result.successful,
"failed": result.failed,
"duration_s": result.duration_s,
}
return app
---------------------------------------------------------------------------
Tests con pytest-asyncio (Patrón 10)
---------------------------------------------------------------------------
"""
tests/test_pipeline.py
import pytest, asyncio from nutridata_pipeline import fetch_all_providers, normalize_records, RestaurantMenu
@pytest.mark.asyncio async def test_fetch_all_providers_returns_5_records(): import aiohttp menu = RestaurantMenu(1, "Test Resto", ["chicken", "rice"]) sem = asyncio.Semaphore(10) async with aiohttp.ClientSession() as session: records = await fetch_all_providers(session, menu, sem) assert len(records) == 5 # un registro por proveedor
@pytest.mark.asyncio async def test_normalize_adds_consolidated(): from nutridata_pipeline import NutritionRecord raw = [ NutritionRecord(1, "usda", 120, 5.0, 20.0, 3.0), NutritionRecord(1, "edamam", 130, 6.0, 22.0, 4.0), ] normalized = await normalize_records(raw) providers = {r.provider for r in normalized} assert "consolidated" in providers
@pytest.mark.asyncio async def test_global_timeout(): import asyncio from nutridata_pipeline import ingest_batch, RestaurantMenu menus = [RestaurantMenu(i, f"Resto {i}", ["beef"]) for i in range(100)] result = await ingest_batch(menus) assert result.duration_s <= 12 # holgura sobre GLOBAL_TIMEOUT """
---------------------------------------------------------------------------
Demo de ejecución (sin BD real)
---------------------------------------------------------------------------
if name == "main": import random
SAMPLE_RESTAURANTS = [
RestaurantMenu(
restaurant_id=i,
name=f"Restaurante #{i}",
ingredients=random.sample(
["pollo", "arroz", "tomate", "cebolla", "aceite", "sal", "ajo", "pimiento"],
k=3,
),
)
for i in range(1, 21) # lote de 20 restaurantes
]
async def demo():
log.info("=== DEMO NutriData Pipeline — lote de %d restaurantes ===", len(SAMPLE_RESTAURANTS))
log.info("Con pipeline SYNC: ~%.0f s", len(SAMPLE_RESTAURANTS) * 2.25)
# Simulamos sin BD real (mock)
class MockDB(DBConnection):
def __init__(self): pass
async def __aenter__(self): return self
async def __aexit__(self, *a): pass
async def upsert_record(self, rec): await asyncio.sleep(0)
queue: asyncio.Queue = asyncio.Queue(maxsize=500)
semaphore = asyncio.Semaphore(MAX_CONCURRENCY)
mock_db = MockDB()
persist_task = asyncio.create_task(persist_worker(queue, mock_db))
t0 = time.monotonic()
async with aiohttp.ClientSession() as session:
tasks = [ingest_restaurant(session, r, queue, semaphore) for r in SAMPLE_RESTAURANTS]
results = await asyncio.gather(*tasks, return_exceptions=True)
await queue.put(None)
await persist_task
elapsed = time.monotonic() - t0
ok = sum(1 for r in results if not isinstance(r, Exception))
log.info(
"=== RESULTADO: %d/%d restaurantes OK en %.2f s (speedup ~%.1fx) ===",
ok, len(SAMPLE_RESTAURANTS), elapsed,
len(SAMPLE_RESTAURANTS) * 2.25 / max(elapsed, 0.001),
)
asyncio.run(demo())
// qué_hace
Proporciona patrones de referencia para implementar programacion asincrona en Python con asyncio y async/await.
// cómo_lo_hace
Documenta patrones probados (gather, tasks, timeouts, manejo de errores) con ejemplos de codigo listos para usar y guia de decision sync vs async.
// ejemplo_de_uso
Úsala cuando tu servicio Python hace múltiples llamadas a APIs externas y el tiempo de espera bloqueante se dispara. Ej.: sustituyes las llamadas síncronas a tres APIs por un asyncio.gather y reduces el tiempo total de 9 s a 3 s ejecutándolas en paralelo.
// plataformas
// opiniones_de_la_comunidad
Opiniones
Cargando opiniones…
// pase_cultiva_ia
Llévate todo el arsenal con el Pase
Todas las skills, prompts y automatizaciones del catálogo en un único archivo, listas para usar: un pago, acceso de por vida y las novedades que añadamos. Sin suscripción.
Pago único · IVA incluido · pago seguro con Stripe.
Acceso inmediato · si no es lo que esperabas, te devolvemos los 10 €.