Paso 1 El modelo mental: PHP-FPM vs event loop
Cómo lo piensas en PHP// cada request tiene su proceso; bloquear es gratis para los demás
$rating = $client->get("ratings/$id"); // 300 ms, este proceso espera
$fotos = $client->get("fotos/$id"); // 300 ms
$vecinos = $client->get("vecinos/$id"); // 300 ms
// ≈ 900 ms; el resto de requests siguen en SUS procesos (hasta agotar los workers)
Cómo lo piensas en m-b-core# un proceso, muchas peticiones; mientras una espera, el loop atiende otras
rating, fotos, vecinos = await asyncio.gather(
fetch_rating(client, id), # las tres salen a la vez
fetch_fotos(client, id),
fetch_vecinos(client, id),
)
# ≈ 300 ms; y NADIE puede bloquear el loop (time.sleep, requests, CPU pesado)
La regla del loop. Una función que hace I/O es async def y se llama con await. Todo lo que no cede el control (una llamada bloqueante, un bucle de CPU) congela el proceso entero para todas las peticiones concurrentes. En PHP ese error cuesta un worker; aquí cuesta el servicio.
En el widget de experiencias hay un buen ejemplo PHP de "hacer dos llamadas a la vez": ExperienceController.php:48-75 (WF-03) usa getAsync + Promise\settle con un solo cliente Guzzle. Es la misma idea; en Python es el modo por defecto, no una excepción.
Paso 2 Un cliente, con el timeout de la casa
Crea el módulo services/api/app/enrichment/ (con __init__.py y tests/__init__.py). Primero las constantes y la forma del resultado, luego el cliente:
services/api/app/enrichment/models.py"""Constantes y formas de datos del módulo enrichment (guideline 1: sin números mágicos)."""
from __future__ import annotations
from dataclasses import dataclass, field
# Timeout total por llamada saliente (segundos). En m-b-core el equivalente es FETCH_TIMEOUT_S.
HTTP_TIMEOUT_S = 5.0
# Máximo de llamadas concurrentes al upstream por request (no saturar al legacy).
MAX_CONCURRENCY = 3
# Radio de "cerca" para contar locales vecinos (km).
NEARBY_RADIUS_KM = 1.5
@dataclass(frozen=True)
class LocalEnriquecido:
"""Resultado de enriquecer un local con tres fuentes independientes."""
local_id: int
nombre: str
rating: float | None
fotos: int
vecinos_cerca: int
fuentes_ok: tuple[str, ...] = field(default_factory=tuple)
services/api/app/enrichment/client.py"""Un solo cliente HTTP saliente para todo el módulo, con timeout de la casa."""
from __future__ import annotations
import httpx
from .models import HTTP_TIMEOUT_S
def make_client(base_url: str) -> httpx.AsyncClient:
"""Cliente async compartido. Se crea una vez (lifespan de la app o fixture) y se inyecta."""
return httpx.AsyncClient(
base_url=base_url,
timeout=httpx.Timeout(HTTP_TIMEOUT_S),
headers={"User-Agent": "mesa-core-practicas/enrichment"},
)
Por qué en el constructor. Si el timeout va en cada llamada, alguien lo olvidará (en m-b-core hay cuatro AsyncClient() sin timeout que sobreviven porque fetch() lo pasa por llamada). En el constructor, nadie puede olvidarlo. Y un solo cliente = un pool de conexiones, no uno por método.
Paso 3 Las tres fuentes: solo I/O
services/api/app/enrichment/queries.py"""I/O: llamadas a los tres servicios que enriquecen un local. Sin lógica de negocio.
`r.json()` devuelve `Any`; con `warn_return_any` (config de m-b-core) mypy exige que
digamos qué esperamos. Anotar la variable es la forma barata de hacerlo explícito.
"""
from __future__ import annotations
from typing import Any
import httpx
async def fetch_rating(client: httpx.AsyncClient, local_id: int) -> dict[str, Any]:
r = await client.get(f"/ratings/{local_id}")
r.raise_for_status()
data: dict[str, Any] = r.json()
return data
async def fetch_fotos(client: httpx.AsyncClient, local_id: int) -> list[dict[str, Any]]:
r = await client.get(f"/fotos/{local_id}")
r.raise_for_status()
data: list[dict[str, Any]] = r.json()
return data
async def fetch_vecinos(client: httpx.AsyncClient, local_id: int) -> list[dict[str, Any]]:
r = await client.get(f"/vecinos/{local_id}")
r.raise_for_status()
data: list[dict[str, Any]] = r.json()
return data
Y la función pura que después va a correr fuera del loop (misma fórmula que cities/utils.py en m-b-core):
services/api/app/enrichment/utils.py"""Funciones puras: cálculo sin I/O. Se testean con parametrize, sin mocks."""
from __future__ import annotations
import math
from typing import Any
from .models import NEARBY_RADIUS_KM
_EARTH_KM = 6371.0
def haversine_km(lat1: float, lng1: float, lat2: float, lng2: float) -> float:
"""Distancia en km entre dos coordenadas."""
p1, p2 = math.radians(lat1), math.radians(lat2)
dphi = math.radians(lat2 - lat1)
dl = math.radians(lng2 - lng1)
a = math.sin(dphi / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dl / 2) ** 2
return 2 * _EARTH_KM * math.asin(math.sqrt(a))
def contar_vecinos_cerca(
origen: tuple[float, float], vecinos: list[dict[str, Any]], radio_km: float = NEARBY_RADIUS_KM
) -> int:
"""Cuenta vecinos dentro del radio. CPU pura: en el servicio corre en to_thread."""
lat0, lng0 = origen
return sum(
1
for v in vecinos
if v.get("lat") is not None
and v.get("lng") is not None
and haversine_km(lat0, lng0, float(v["lat"]), float(v["lng"])) <= radio_km
)
Fíjate en v.get("lat"). Con filas sucias del legacy, v["lat"] lanzaría KeyError (trampa 5 de la guía). Aquí una fila sin coordenadas simplemente no cuenta, y lo decimos con un test.
Paso 4 Orquestar: gather + Semaphore + to_thread, degradando por fuente
services/api/app/enrichment/service.py"""Orquestador: tres fuentes EN PARALELO, con límite de concurrencia, degradando por fuente.
Patrón de m-b-core: services/api/app/availability_engine/legacy_queries.py:136
results = await asyncio.gather(*[_query(q, params, enabled) for q in QUERIES])
y el Semaphore del workbench (availability_engine/workbench/parity.py:920).
"""
from __future__ import annotations
import asyncio
import logging
from typing import Any
import httpx
from . import queries
from .models import MAX_CONCURRENCY, LocalEnriquecido
from .utils import contar_vecinos_cerca
log = logging.getLogger(__name__)
async def _guardado(sem: asyncio.Semaphore, nombre: str, coro: Any) -> tuple[str, Any | None]:
"""Ejecuta una fuente bajo el semáforo; un fallo de red degrada a None y se registra.
Solo capturamos errores de red/HTTP (lo que sabemos manejar); cualquier otra
excepción sube.
"""
async with sem:
try:
return nombre, await coro
except (httpx.HTTPError, TimeoutError) as e:
log.warning("enrichment.fuente_fallo", extra={"fuente": nombre, "error": type(e).__name__})
return nombre, None
async def enriquecer_local(
client: httpx.AsyncClient,
local_id: int,
nombre: str,
coords: tuple[float, float],
) -> LocalEnriquecido:
sem = asyncio.Semaphore(MAX_CONCURRENCY)
resultados = await asyncio.gather(
_guardado(sem, "rating", queries.fetch_rating(client, local_id)),
_guardado(sem, "fotos", queries.fetch_fotos(client, local_id)),
_guardado(sem, "vecinos", queries.fetch_vecinos(client, local_id)),
)
por_fuente: dict[str, Any] = dict(resultados)
ok = tuple(k for k, v in por_fuente.items() if v is not None)
rating = por_fuente["rating"]["rating"] if por_fuente["rating"] else None
fotos = len(por_fuente["fotos"] or [])
# CPU pura sobre una lista potencialmente grande: fuera del loop.
vecinos = await asyncio.to_thread(contar_vecinos_cerca, coords, por_fuente["vecinos"] or [])
return LocalEnriquecido(
local_id=local_id,
nombre=nombre,
rating=rating,
fotos=fotos,
vecinos_cerca=vecinos,
fuentes_ok=ok,
)
Corre ruff. Te va a marcar algo real:
Deberías ver (si escribiste asyncio.TimeoutError, como es natural viniendo de la doc vieja)services/api/app/enrichment/service.py:34:16: UP041 [*] Replace aliased errors with `TimeoutError`
|
34 | except (httpx.HTTPError, asyncio.TimeoutError) as e:
= help: Replace with builtin `TimeoutError`
Qué te acaba de enseñar ruff. Desde Python 3.11 asyncio.TimeoutError es un alias del TimeoutError nativo. La regla UP (pyupgrade) de la config de m-b-core mantiene el código en el Python que corres, no en el de los tutoriales de 2019. Déjalo en TimeoutError y vuelve a correr: All checks passed!
Tres decisiones del orquestador. (1) gather lanza las tres corrutinas a la vez; el Semaphore(3) aquí no limita nada, pero es el patrón: cuando sean 40 locales, limita a 3 en vuelo. (2) _guardado captura solo errores de red y HTTP: una fuente caída degrada a None y se loguea con extra={}; un KeyError o un bug tuyo suben, como debe ser (vicio 12: el catch que traga). (3) to_thread saca el bucle de CPU del loop: con 2.000 vecinos no importa, con 200.000 sí, y el hábito se adquiere en el caso chico.
Paso 5 Tests: funciones puras con parametrize, orquestador con respx
Sigue el árbol de decisión de TESTING.md: la función pura se parametriza sin mocks; el orquestador se prueba mockeando la red (el nivel más bajo), no fetch_rating.
services/api/app/enrichment/tests/conftest.pyimport pytest
from services.api.app.enrichment.client import make_client
BASE = "https://upstream.test"
@pytest.fixture()
async def client():
"""Un cliente por test, cerrado al terminar (async fixture: asyncio_mode = auto)."""
c = make_client(BASE)
yield c
await c.aclose()
@pytest.fixture()
def coords_lima() -> tuple[float, float]:
return (-12.1211, -77.0300)
services/api/app/enrichment/tests/test_utils.py"""Funciones puras → parametrizar bordes, sin mocks (TESTING.md §1)."""
import pytest
from services.api.app.enrichment.utils import contar_vecinos_cerca, haversine_km
LIMA = (-12.1211, -77.0300)
@pytest.mark.parametrize(
"a, b, esperado_km",
[
(LIMA, LIMA, 0.0),
(LIMA, (-12.1211, -77.0400), 1.09), # ~1 km al oeste
(LIMA, (-33.4489, -70.6693), 2458.0), # Santiago
],
)
def test_haversine(a, b, esperado_km):
assert haversine_km(*a, *b) == pytest.approx(esperado_km, abs=0.5)
@pytest.mark.parametrize(
"vecinos, esperado",
[
([], 0),
([{"lat": -12.1211, "lng": -77.0300}], 1), # el mismo punto
([{"lat": -12.1211, "lng": -77.0400}, {"lat": -33.4, "lng": -70.6}], 1), # 1 km sí, Santiago no
([{"lat": None, "lng": -77.0}, {"lng": -77.0}], 0), # filas sucias no cuentan (dict.get)
],
)
def test_contar_vecinos_cerca(vecinos, esperado):
assert contar_vecinos_cerca(LIMA, vecinos) == esperado
services/api/app/enrichment/tests/test_service.py"""Orquestador → red simulada con respx (mock en el borde de red, no en la función bajo prueba)."""
import httpx
import respx
from services.api.app.enrichment.service import enriquecer_local
from .conftest import BASE
@respx.mock(base_url=BASE)
async def test_tres_fuentes_ok(respx_mock, client, coords_lima):
respx_mock.get("/ratings/11").mock(return_value=httpx.Response(200, json={"rating": 4.6}))
respx_mock.get("/fotos/11").mock(return_value=httpx.Response(200, json=[{"id": 1}, {"id": 2}]))
respx_mock.get("/vecinos/11").mock(
return_value=httpx.Response(200, json=[{"lat": -12.1211, "lng": -77.0310}])
)
r = await enriquecer_local(client, 11, "Maido", coords_lima)
assert r.rating == 4.6
assert r.fotos == 2
assert r.vecinos_cerca == 1
assert r.fuentes_ok == ("rating", "fotos", "vecinos")
@respx.mock(base_url=BASE)
async def test_una_fuente_caida_degrada_no_rompe(respx_mock, client, coords_lima):
respx_mock.get("/ratings/11").mock(side_effect=httpx.ConnectError("boom")) # servidor caído
respx_mock.get("/fotos/11").mock(return_value=httpx.Response(200, json=[]))
respx_mock.get("/vecinos/11").mock(return_value=httpx.Response(200, json=[]))
r = await enriquecer_local(client, 11, "Maido", coords_lima)
assert r.rating is None
assert r.fuentes_ok == ("fotos", "vecinos")
@respx.mock(base_url=BASE)
async def test_500_del_upstream_es_fuente_caida(respx_mock, client, coords_lima):
respx_mock.get("/ratings/11").mock(return_value=httpx.Response(500))
respx_mock.get("/fotos/11").mock(return_value=httpx.Response(200, json=[]))
respx_mock.get("/vecinos/11").mock(return_value=httpx.Response(200, json=[]))
r = await enriquecer_local(client, 11, "Maido", coords_lima)
assert "rating" not in r.fuentes_ok
@respx.mock(base_url=BASE)
async def test_timeout_no_cuelga_el_request(respx_mock, client, coords_lima):
"""respx no simula el paso del tiempo: el timeout se simula lanzando la excepción que
lanzaría el transporte real (httpx.ReadTimeout). Lo que probamos es que el orquestador
la trata como fuente caída y no la deja subir."""
respx_mock.get("/ratings/11").mock(side_effect=httpx.ReadTimeout("timed out"))
respx_mock.get("/fotos/11").mock(return_value=httpx.Response(200, json=[]))
respx_mock.get("/vecinos/11").mock(return_value=httpx.Response(200, json=[]))
r = await enriquecer_local(client, 11, "Maido", coords_lima)
assert r.rating is None
assert r.fuentes_ok == ("fotos", "vecinos")
async def test_el_cliente_de_la_casa_tiene_timeout(client):
"""Guardarraíl: si alguien quita el timeout del cliente, este test lo pilla."""
assert client.timeout.connect == pytest.approx(5.0)
assert client.timeout.read == pytest.approx(5.0)
(test_service.py necesita también import pytest arriba, por el approx; ruff te lo ordena.)
Deberías verservices/api/app/enrichment/tests/test_service.py::test_tres_fuentes_ok PASSED [ 8%]
services/api/app/enrichment/tests/test_service.py::test_una_fuente_caida_degrada_no_rompe PASSED [ 16%]
services/api/app/enrichment/tests/test_service.py::test_500_del_upstream_es_fuente_caida PASSED [ 25%]
services/api/app/enrichment/tests/test_service.py::test_timeout_no_cuelga_el_request PASSED [ 33%]
services/api/app/enrichment/tests/test_service.py::test_el_cliente_de_la_casa_tiene_timeout PASSED [ 41%]
services/api/app/enrichment/tests/test_utils.py::test_haversine[a0-b0-0.0] PASSED [ 50%]
…
services/api/app/enrichment/tests/test_utils.py::test_contar_vecinos_cerca[vecinos3-0] PASSED [100%]
============================== 12 passed in 0.04s ==============================
Dos cosas que aprendí escribiendo estos tests (te van a pasar). (1) Mi primera versión del test de timeout hacía await asyncio.sleep(0.5) dentro del mock con un cliente de 50 ms de timeout, y el test pasaba con rating=5: respx intercepta el transporte, así que no hay reloj de red que expire. El timeout se simula con side_effect=httpx.ReadTimeout(...). (2) Puse "Santiago está a 2455,6 km" de memoria; la fórmula dijo 2458,0. El test parametrizado te obliga a medir, no a recordar.
Paso 6 Medir: serie vs paralelo con un upstream falso
Un servidor FastAPI de pruebas cuyos tres endpoints tardan 300 ms cada uno (simula el legacy lento):
scripts/upstream_fake.py"""Servidor de pruebas: tres endpoints que tardan 300 ms cada uno (simulan el legacy)."""
import asyncio
from fastapi import FastAPI
app = FastAPI()
LATENCIA_S = 0.3
@app.get("/ratings/{local_id}")
async def ratings(local_id: int):
await asyncio.sleep(LATENCIA_S)
return {"rating": 4.6}
@app.get("/fotos/{local_id}")
async def fotos(local_id: int):
await asyncio.sleep(LATENCIA_S)
return [{"id": i} for i in range(3)]
@app.get("/vecinos/{local_id}")
async def vecinos(local_id: int):
await asyncio.sleep(LATENCIA_S)
return [{"lat": -12.1211 + i * 0.001, "lng": -77.03} for i in range(2000)]
scripts/serie_vs_paralelo.py"""Mide: tres llamadas en serie (como en PHP) vs gather (como en m-b-core)."""
import asyncio
import time
from services.api.app.enrichment import queries
from services.api.app.enrichment.client import make_client
from services.api.app.enrichment.service import enriquecer_local
BASE = "http://127.0.0.1:8099"
async def en_serie(client, local_id):
await queries.fetch_rating(client, local_id)
await queries.fetch_fotos(client, local_id)
await queries.fetch_vecinos(client, local_id)
async def main():
async with make_client(BASE) as client:
t0 = time.perf_counter()
await en_serie(client, 11)
t_serie = time.perf_counter() - t0
t0 = time.perf_counter()
r = await enriquecer_local(client, 11, "Maido", (-12.1211, -77.03))
t_par = time.perf_counter() - t0
print(f"en serie : {t_serie:.2f} s (3 x 300 ms + red)")
print(f"con gather: {t_par:.2f} s -> {r}")
asyncio.run(main())
# terminal 1
PYTHONPATH=. uvicorn scripts.upstream_fake:app --port 8099 --log-level warning
# terminal 2
PYTHONPATH=. python scripts/serie_vs_paralelo.py
Deberías veren serie : 0.93 s (3 x 300 ms + red)
con gather: 0.32 s -> LocalEnriquecido(local_id=11, nombre='Maido', rating=4.6, fotos=3, vecinos_cerca=14, fuentes_ok=('rating', 'fotos', 'vecinos'))
0,93 → 0,32. Tres esperas de 300 ms que se solapan. Es la razón por la que el motor de disponibilidad de m-b-core lanza todas sus queries al legacy con un gather (legacy_queries.py:136) en vez de una tras otra. Y fíjate en vecinos_cerca=14: 2.000 vecinos contados en un hilo aparte mientras el loop seguía libre.
Paso 7 Los tres errores clásicos, provocados a propósito
Escribe los tres scripts y córrelos (el upstream falso sigue arriba para el tercero). El objetivo es que reconozcas la huella de cada error cuando aparezca en m-b-core.
scripts/error_1_sin_await.py"""Error clásico 1: olvidar el await. No da error: devuelve una corrutina sin ejecutar."""
import asyncio
async def buscar_local(local_id: int) -> dict:
await asyncio.sleep(0.01)
return {"id": local_id, "nombre": "Maido"}
async def main():
local = buscar_local(11) # <- falta el await
print("tipo:", type(local).__name__)
print("nombre:", local["nombre"]) # TypeError: 'coroutine' object is not subscriptable
asyncio.run(main())
python -W default scripts/error_1_sin_await.py
Deberías vertipo: coroutine
Traceback (most recent call last):
…
print("nombre:", local["nombre"]) # TypeError: 'coroutine' object is not subscriptable
~~~~~^^^^^^^^^^
TypeError: 'coroutine' object is not subscriptable
sys:1: RuntimeWarning: coroutine 'buscar_local' was never awaited
La huella. Dos síntomas que siempre significan "olvidaste el await": type(x).__name__ == 'coroutine' y el RuntimeWarning: coroutine '…' was never awaited al final. Con mypy en el módulo (paso 8) el error aparece antes de correr: "Coroutine[…]" has no attribute "__getitem__".
scripts/error_2_time_sleep.py"""Error clásico 2: time.sleep dentro de async bloquea el loop; asyncio.sleep lo cede."""
import asyncio
import time
async def tarea(nombre: str, dormir):
await dormir(0.5)
return nombre
async def correr(dormir, etiqueta):
t0 = time.perf_counter()
await asyncio.gather(*[tarea(f"t{i}", dormir) for i in range(5)])
print(f"{etiqueta:<22} 5 tareas de 0.5 s -> {time.perf_counter() - t0:.2f} s")
async def dormir_bloqueante(s: float):
time.sleep(s) # <- bloquea el loop: nadie más corre mientras tanto
asyncio.run(correr(asyncio.sleep, "asyncio.sleep (cede):"))
asyncio.run(correr(dormir_bloqueante, "time.sleep (bloquea):"))
python scripts/error_2_time_sleep.py
Deberías verasyncio.sleep (cede): 5 tareas de 0.5 s -> 0.50 s
time.sleep (bloquea): 5 tareas de 0.5 s -> 2.51 s
La huella. "Usé gather y sigue tardando lo mismo que en serie". Es que dentro hay algo bloqueante: time.sleep, un driver síncrono, subprocess.run (m-b-core lo tiene en worker/app/modules/menus/router.py:105 con pdftoppm), o un bucle de CPU. La solución es la versión async o asyncio.to_thread.
scripts/error_3_requests_sync.py"""Error clásico 3: un cliente HTTP síncrono (requests) dentro de async def.
Funciona... y bloquea el loop en cada llamada. Comparado con httpx.AsyncClient."""
import asyncio
import time
import httpx
import requests
BASE = "http://127.0.0.1:8099"
async def con_requests(local_id: int):
return requests.get(f"{BASE}/ratings/{local_id}", timeout=5).json() # bloqueante
async def con_httpx(client: httpx.AsyncClient, local_id: int):
r = await client.get(f"/ratings/{local_id}")
return r.json()
async def main():
t0 = time.perf_counter()
await asyncio.gather(*[con_requests(i) for i in range(5)])
print(f"requests en async def: 5 llamadas -> {time.perf_counter() - t0:.2f} s (en serie, aunque uses gather)")
async with httpx.AsyncClient(base_url=BASE, timeout=5) as client:
t0 = time.perf_counter()
await asyncio.gather(*[con_httpx(client, i) for i in range(5)])
print(f"httpx.AsyncClient : 5 llamadas -> {time.perf_counter() - t0:.2f} s (de verdad en paralelo)")
asyncio.run(main())
python scripts/error_3_requests_sync.py
Deberías verrequests en async def: 5 llamadas -> 1.52 s (en serie, aunque uses gather)
httpx.AsyncClient : 5 llamadas -> 0.31 s (de verdad en paralelo)
La huella. No hay warning, no hay error: solo lentitud y un servicio que "se cae" con carga. Por eso m-b-core no tiene requests en requirements/base.txt y por eso el conector MySQL síncrono de shared/legacy_db lleva el comentario "don't use en async". Cuando termines la práctica, quita requests de requirements/local.txt: estaba solo para esto.
Paso 8 Lint, tipos y commit
make lint
PYTHONPATH=. mypy services/api/app/enrichment/ --explicit-package-bases
Deberías ver la primera vez (si escribiste return r.json() directo)services/api/app/enrichment/queries.py:19: error: Returning Any from function declared to return "list[dict[str, Any]]" [no-any-return]
services/api/app/enrichment/queries.py:25: error: Returning Any from function declared to return "list[dict[str, Any]]" [no-any-return]
Found 3 errors in 1 file (checked 10 source files)
Es warn_return_any = true, de la config de m-b-core: r.json() devuelve Any y tú prometiste una lista. Anota la variable (data: list[dict[str, Any]] = r.json()) como en el paso 3 y vuelve a correr:
Deberías verSuccess: no issues found in 10 source files
All checks passed!
17 files already formatted
git switch -c feat/enrichment-async
git add .
git commit -m "feat(enrichment): tres fuentes en paralelo con cliente único, timeout y degradación por fuente"