← Volver al catálogo
IA, ingeniería y MLOpsReferenciaIntermedioGratis

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.

Descargar SKILL.md

Descarga abierta · sin registro · para Python, FastAPI, aiohttp

// resultado_de_ejemplo

""" 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:

  1. async/await básico — fetch_nutrition_data, fetch_price_data
  2. gather() concurrente — fetch_all_providers (5 APIs en paralelo)
  3. asyncio.create_task — lanzar el worker de persistencia en background
  4. Manejo de errores + gather(return_exceptions=True)
  5. asyncio.wait_for con timeout — 3 s por llamada, 10 s global
  6. Semáforo para throttling — máx. 20 conexiones simultáneas a proveedores
  7. Async context manager — DBConnection
  8. asyncio.to_thread — normalización CPU-bound
  9. Cancellation handling — ingest_batch cancelable
  10. 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

PythonFastAPIaiohttpasyncio
Categoría
IA, ingeniería y MLOps
Tipo
Referencia
Nivel
Intermedio
Licencia
MIT
Seguridad
seguro · riesgo bajo
Versión
1.0.0

// 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 €.