ruta rag en profundidad / módulo 5

M5 — Async, concurrencia y sistemas

Outcome del módulo: entender qué se bloquea primero en scholar-rag bajo carga concurrente y por qué, midiendo throughput y latencia contra el número de peticiones simultáneas. Aquí bajas del RAG al sistema que lo sostiene.


Concepto 1 — El event loop: un hilo que atiende a muchos

Idea núcleo: asyncio no usa muchos hilos; usa uno que salta entre tareas cada vez que una se pone a esperar I/O. Por eso un backend async atiende miles de conexiones con recursos mínimos.

Entender (capa 1)

Texto: en un servidor síncrono clásico, cada petición ocupa un hilo del principio al fin, y cuando ese hilo espera la base de datos (milisegundos eternos para un CPU), se queda ahí sentado sin hacer nada. Async cambia el modelo: hay un solo hilo con un event loop que mantiene muchas tareas a medio hacer. Cuando la tarea A hace await db.fetch(...) y se pone a esperar la red, el loop la suspende y corre la tarea B, luego C, y vuelve a A cuando su dato llegó. El truco es que casi todo el tiempo de un backend es esperar (BD, APIs, disco): mientras una espera, otra avanza.

La palabra clave es await: marca los puntos donde una tarea cede el control al loop. Un backend async escala porque en esos puntos de espera atiende a otros. Pero eso esconde una trampa que es el Concepto 2: si una tarea no cede —porque está haciendo cálculo puro de CPU, sin await— el loop se congela y nadie más avanza.

Visual:

flowchart TB
    L((event loop)) --> A["tarea A: await db - esperando"]
    L --> B["tarea B: corriendo"]
    L --> C["tarea C: await api - esperando"]
    A -.dato llego.-> L
    B -.await.-> L
    L --> N[atiende la siguiente lista]

Ejemplo:

# cada await es un punto donde el loop puede atender a otra peticion
async def answer(question: str) -> RagAnswer:
    vec = await embed(question)             # cede si embed es I/O
    ctx = await retrieval.hybrid_search(vec) # cede esperando Postgres
    return await llm.generate(question, ctx) # cede esperando el LLM

Fijar (capa 2)

Nota atómica:

Feynman: “El event loop es un mesero para 20 mesas. No se planta en la mesa 1 hasta que terminen de comer; toma la orden (await), va a la cocina, y mientras se cocina atiende la mesa 2 y 3. Un solo mesero sirve a todos porque el tiempo real se va esperando la cocina.”

Aplicar (capa 3)

Recorre el flujo de una pregunta en scholar-rag y marca cada await. Esos son los puntos donde el sistema puede atender a otra petición. Anota si hay algún paso pesado sin await (cálculo puro): ese es sospechoso de bloquear el loop, y es lo que investigas en el Concepto 2.


Concepto 2 — El GIL y el embedding local: dónde se congela todo

Idea núcleo: el embedding local (FastEmbed) es cálculo de CPU. Correrlo en el event loop lo bloquea: por el GIL, mientras Python calcula el vector, ninguna otra tarea avanza. Es el cuello más traicionero de un RAG async.

Entender (capa 1)

Texto: el GIL (Global Interpreter Lock) de CPython permite ejecutar bytecode de Python en un solo hilo a la vez. Para trabajo I/O-bound no importa: cuando esperas la red, el GIL se libera y otras tareas corren. Pero el trabajo CPU-bound —generar un embedding con FastEmbed es matemática pura sobre el texto— agarra el GIL y no lo suelta hasta terminar. Si eso ocurre dentro del event loop, congela el loop entero: durante esos milisegundos, todas las demás peticiones esperan, aunque solo necesitaran leer la BD.

scholar-rag genera embeddings localmente. Si esa generación corre en el hilo del event loop, tu servidor “async” se comporta como síncrono bajo carga: una petición pesada de embedding traba a todas. La solución es sacar el cálculo del loop: asyncio.to_thread(...) o un pool de procesos, para que el loop siga atendiendo mientras un worker aparte hace el número. Medir esto —latencia de las otras peticiones con y sin el embedding en el loop— es de los experimentos que más enseñan.

Visual:

flowchart TB
    subgraph malo["embed en el event loop"]
    E1[peticion A: calcula embedding - agarra GIL] --> BLK[loop congelado]
    BLK --> W[peticiones B,C,D esperan sin razon]
    end
    subgraph bueno["embed fuera del loop"]
    E2[peticion A: to_thread embedding] --> LOOP[loop sigue libre]
    LOOP --> OK[B,C,D avanzan]
    end

Ejemplo:

# MAL: cálculo CPU-bound dentro del loop, congela a todos
def embed_sync(text): ...              # FastEmbed, puro CPU
vec = embed_sync(question)             # bloquea el event loop

# BIEN: sácalo a un hilo, el loop sigue atendiendo
vec = await asyncio.to_thread(embed_sync, question)

Fijar (capa 2)

Nota atómica:

Feynman: “El GIL es que en la cocina solo cabe un cocinero. Esperar el horno (I/O) no lo ocupa, así que puede atender otras órdenes. Pero picar cebolla a mano (CPU) sí lo ocupa entero: mientras pica, nadie más cocina. El embedding local es picar cebolla en plena cocina.”

Aplicar (capa 3)

Averigua en scholar-rag si el embedding se genera dentro del flujo async directo o si ya está aislado. Si está en el loop, envuélvelo en asyncio.to_thread y mide la latencia p95 de peticiones concurrentes antes y después. Ese antes/después es tu prueba de que entiendes el GIL, no solo que lo puedes definir.


Concepto 3 — El pool de conexiones: cuántas queries a la vez

Idea núcleo: el pool_size define cuántas queries concurrentes van a la base antes de que empiecen a hacer cola. Es un dial de capacidad, y su valor correcto se mide, no se copia.

Entender (capa 1)

Texto: abrir una conexión a Postgres es caro, así que asyncpg mantiene un pool de conexiones abiertas y reutilizables. scholar-rag usa pool_size=20. Ese número es un techo: si llegan 50 peticiones que necesitan la BD a la vez, 20 corren y 30 esperan a que se libere una. Un pool muy chico crea cola artificial (tienes CPU y BD ociosas pero peticiones esperando un slot); uno muy grande puede saturar la base (Postgres también tiene su límite de conexiones, y pasarlo lo degrada a él).

El valor correcto no es “20 porque venía así”. Depende de cuántas conexiones aguanta tu Postgres, cuánto dura cada query, y tu patrón de carga. Se encuentra midiendo: subes la concurrencia y observas dónde la latencia se dispara: si es antes de saturar CPU/BD, el pool es el cuello; si el pool está holgado y aun así se degrada, el cuello está en otro lado (la BD, el LLM, el embedding).

Visual:

flowchart TB
    R[50 peticiones concurrentes] --> P{pool_size = 20}
    P --> A[20 corriendo contra Postgres]
    P --> Q[30 en cola esperando slot]
    Q -.pool muy chico.-> W[latencia sube sin saturar recursos]
    A -.pool muy grande.-> S[Postgres saturado, se degrada la BD]

Ejemplo:

# repositories/db.py — el pool y su techo
pool = await asyncpg.create_pool(dsn, min_size=5, max_size=20)
# max_size=20 es el dial: mídelo, no lo heredes.
# Regla de arranque, no de fe: max_size <= (conexiones que aguanta tu Postgres)

Fijar (capa 2)

Nota atómica:

Feynman: “El pool son las cajas abiertas en el supermercado. Dos cajas y 50 clientes: fila larga con góndolas vacías. Cincuenta cajas y un cajero por caja que no existe: caos. El número correcto sale de contar clientes reales por minuto, no de copiar al supermercado de al lado.”

Aplicar (capa 3)

Corre scholar-rag con pool_size en 5, 20 y 50 bajo una carga concurrente fija y mide la latencia p95 en cada caso. Observa si el número mueve la aguja. Anota el valor que da mejor latencia sin degradar Postgres: ahora ese 20 (o el que sea) está justificado con datos.


Concepto 4 — Timeouts y backpressure: qué pasa cuando llega de más

Idea núcleo: un sistema maduro decide a propósito qué hacer cuando la carga supera su capacidad: cortar lo que tarda demasiado (timeout) y rechazar rápido en vez de encolar infinito (backpressure). Sin eso, un pico no te ralentiza, te tumba.

Entender (capa 1)

Texto: dos fallas clásicas bajo carga. Primera: una llamada al LLM o a la BD que se cuelga y, sin timeout, retiene su slot del pool y su tarea del loop para siempre; unas cuantas así y agotas la capacidad con peticiones zombie. Cada llamada a un tercero necesita un límite de tiempo tras el cual se corta y se maneja el error. Segunda: cuando llegan más peticiones de las que puedes atender, encolarlas todas parece amable pero es fatal: la cola crece, la latencia de todos se dispara, la memoria sube, y el sistema colapsa entero en vez de degradar suave. Backpressure es la decisión explícita de rechazar rápido el exceso (responder “ocupado, reintenta”) para proteger a los que ya estás atendiendo.

La mentalidad senior: el sistema no solo funciona cuando todo va bien, tiene un comportamiento diseñado para cuando va mal. “¿Qué pasa con la petición número 1001 cuando aguanto 1000?” tiene que tener respuesta, y esa respuesta no puede ser “se cae todo”.

Visual:

flowchart TB
    P[pico de carga] --> D{sobre capacidad?}
    D -->|sin backpressure| Q[encolar todo: latencia y memoria explotan, colapso total]
    D -->|con backpressure| R[rechazar rapido el exceso: los demas siguen bien]
    L[llamada a LLM/BD] --> T{timeout?}
    T -->|no| Z[peticion zombie retiene el slot para siempre]
    T -->|si| C[corta y maneja el error]

Ejemplo:

# timeout en toda llamada a un tercero
async def generate(question, ctx):
    async with asyncio.timeout(20):        # corta a los 20s, no zombie
        return await llm.complete(question, ctx)

# backpressure simple: limita in-flight y rechaza el exceso
sem = asyncio.Semaphore(1000)
if sem.locked():
    raise HTTPException(503, "Servicio ocupado, reintenta")  # rechaza rapido
async with sem:
    ...

Fijar (capa 2)

Nota atómica:

Feynman: “Sin backpressure, la discoteca deja entrar a todos: se llena, nadie se mueve, colapsa. Con backpressure, el portero dice ‘lleno, vuelve en 10’ a algunos, y los de adentro la pasan bien. Rechazar a unos salva la noche de los demás.”

Aplicar (capa 3)

Revisa si las llamadas al LLM y a Postgres en scholar-rag tienen timeout. Si no, agrégalos. Luego describe (aunque no lo implementes aún) qué debería pasar con la petición que excede tu capacidad medida en el Concepto 3: esa respuesta es tu política de backpressure, y va al design doc.


Cierre del módulo

Ya no ves scholar-rag como “un RAG”, sino como un sistema con un event loop que puede congelarse, un pool que puede ser cuello, y un límite tras el cual hay que degradar a propósito. Sabes qué se rompe primero bajo carga, porque lo mediste.

Lo que puedes responder ahora: “¿cuántas queries concurrentes aguanta antes de degradarse y qué se rompe primero?” — con el benchmark de carga.


Siguiente: M6 — Latencia, costo y operación, donde cierras el ciclo: perfilas dónde se va el tiempo y el dinero, y dejas el sistema observable.