Guzmán D. Darío Senior Python Developer Español Hire me
Taller autoguiado

De bloqueante a tiempo real: un pipeline de video multi-cámara con asyncio

Paso 03 de 23
Paso 3 de 23

El mundo exterior: cámaras y modelo

Antes de construir tu pipeline necesitas algo a lo que conectarlo: cámaras que transmitan video y un modelo que detecte objetos.

En un proyecto real esas dos cosas existen sin que tú las escribas: las cámaras cuelgan de una pared y el modelo corre en un servidor. Aquí las vamos a simular en tu máquina, y por eso este paso es el único del taller donde vas a copiar código sin analizarlo línea por línea.

Copia estos cuatro archivos tal cual y sigue adelante. No son la lección: son el escenario donde ocurre la lección.

Los videos que harán de cámaras

Primero necesitas material que transmitir. Crea make_samples.py:

make_samples.py python
import cv2
import numpy as np
from pathlib import Path

Path("videos").mkdir(exist_ok=True)

def make_video(path, color, n_frames=150, size=(640, 480), fps=25):
    w, h = size
    writer = cv2.VideoWriter(path, cv2.VideoWriter_fourcc(*"mp4v"), fps, size)
    for i in range(n_frames):
        frame = np.full((h, w, 3), 30, dtype=np.uint8)
        x = int((i / n_frames) * (w - 80))
        cv2.rectangle(frame, (x, h // 2 - 40), (x + 80, h // 2 + 40), color, -1)
        cv2.putText(frame, f"{Path(path).stem}  frame {i}", (10, 30),
                    cv2.FONT_HERSHEY_SIMPLEX, 0.7, (255, 255, 255), 2)
        writer.write(frame)
    writer.release()
    print("wrote", path)

for idx, color in enumerate([(0, 0, 200), (0, 200, 0), (200, 0, 0), (0, 200, 200)]):
    make_video(f"videos/cam{idx}.mp4", color)
terminal bash
python make_samples.py

Eso crea cuatro clips de seis segundos con un rectángulo de color moviéndose. Si prefieres usar videos propios, cópialos a la carpeta videos/ en vez de ejecutar esto: cualquier .mp4 corto sirve, y si sale gente o carros el detector real tendrá algo que encontrar.

El servicio de cámaras

Este es el que convierte esos archivos en cámaras. Sirve cada clip como un stream MJPEG infinito por HTTP, recibe los eventos que tu pipeline publique, y aloja el dashboard.

Crea svc_cameras.py:

svc_cameras.py python
import asyncio
import os
from collections import deque
from pathlib import Path

import cv2
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from fastapi.staticfiles import StaticFiles

N_CAMERAS = int(os.environ.get("N_CAMERAS", "8"))
FPS = float(os.environ.get("CAMERA_FPS", "25"))
BOUNDARY = "frame"

app = FastAPI()
app.state.clips = []
app.state.events = deque(maxlen=200)
app.state.subscribers = set()


def load_clips():
    frames_per_clip = []
    for path in sorted(Path("videos").glob("*.mp4")):
        cap = cv2.VideoCapture(str(path))
        jpegs = []
        while True:
            ok, frame = cap.read()
            if not ok:
                break
            ok, buf = cv2.imencode(".jpg", frame, [cv2.IMWRITE_JPEG_QUALITY, 80])
            if ok:
                jpegs.append(buf.tobytes())
        cap.release()
        if jpegs:
            frames_per_clip.append(jpegs)
            print(f"loaded {path.name}: {len(jpegs)} frames")
    return frames_per_clip


@app.on_event("startup")
async def startup():
    app.state.clips = load_clips()
    if not app.state.clips:
        raise SystemExit("no videos found in videos/ (run: python make_samples.py)")
    print(f"serving {N_CAMERAS} camera streams at {FPS} fps")


@app.get("/cameras")
async def cameras():
    return {"cameras": [f"cam{i}" for i in range(N_CAMERAS)], "fps": FPS}


async def mjpeg_stream(clip, offset):
    delay = 1.0 / FPS
    i = offset
    while True:
        jpeg = clip[i % len(clip)]
        yield (b"--" + BOUNDARY.encode() + b"\r\n"
               b"Content-Type: image/jpeg\r\n"
               b"Content-Length: " + str(len(jpeg)).encode() + b"\r\n\r\n"
               + jpeg + b"\r\n")
        i += 1
        await asyncio.sleep(delay)


@app.get("/stream/{cam_id}")
async def stream(cam_id: str):
    index = int(cam_id.removeprefix("cam"))
    clip = app.state.clips[index % len(app.state.clips)]
    offset = (index * 37) % len(clip)
    return StreamingResponse(
        mjpeg_stream(clip, offset),
        media_type=f"multipart/x-mixed-replace; boundary={BOUNDARY}")


@app.post("/events")
async def publish(request: Request):
    event = await request.json()
    app.state.events.append(event)
    dead = set()
    for queue in app.state.subscribers:
        try:
            queue.put_nowait(event)
        except asyncio.QueueFull:
            dead.add(queue)
    app.state.subscribers -= dead
    return {"ok": True}


@app.get("/events/stream")
async def event_stream():
    queue = asyncio.Queue(maxsize=100)
    app.state.subscribers.add(queue)

    async def gen():
        try:
            while True:
                event = await queue.get()
                yield f"data: {event}\n\n".replace("'", '"')
        finally:
            app.state.subscribers.discard(queue)

    return StreamingResponse(gen(), media_type="text/event-stream")


app.mount("/", StaticFiles(directory="static", html=True), name="static")

Ojo

Fíjate en una línea que te va a servir más adelante: offset = (index * 37) % len(clip). Cada cámara arranca su clip en un punto distinto, así que aunque tengas un solo video puedes servir sesenta y cuatro cámaras y todas entregan frames diferentes. Eso es lo que hará posible el paso de escala sin que tengas que grabar nada.

El servicio del modelo

Recibe una imagen por HTTP y devuelve las detecciones en JSON. Trae dos implementaciones: la real con YOLO y una simulada que duerme 30 milisegundos, para que el taller funcione sin PyTorch.

Crea svc_model.py:

svc_model.py python
import asyncio
import os
import random
from concurrent.futures import ThreadPoolExecutor

from fastapi import FastAPI, Request

FAKE = os.environ.get("FAKE_MODEL") == "1"
LATENCY = float(os.environ.get("FAKE_LATENCY", "0.03"))
POOL_SIZE = int(os.environ.get("MODEL_THREADS", "4"))

app = FastAPI()


class RealModel:
    def __init__(self):
        import cv2
        import numpy as np
        from ultralytics import YOLO
        self.cv2 = cv2
        self.np = np
        self.model = YOLO("yolo11n.pt")

    def infer(self, jpeg: bytes):
        arr = self.np.frombuffer(jpeg, dtype=self.np.uint8)
        frame = self.cv2.imdecode(arr, self.cv2.IMREAD_COLOR)
        if frame is None:
            return []
        results = self.model(frame, conf=0.35, verbose=False)[0]
        h, w = frame.shape[:2]
        out = []
        for box in results.boxes:
            x1, y1, x2, y2 = box.xyxy[0].tolist()
            out.append({
                "label": results.names[int(box.cls)],
                "conf": round(float(box.conf), 3),
                "box": [x1 / w, y1 / h, x2 / w, y2 / h],
            })
        return out


class FakeModel:
    def infer(self, jpeg: bytes):
        import time
        time.sleep(LATENCY)
        return [{"label": "person", "conf": 0.9, "box": [0.3, 0.3, 0.6, 0.9]}
                for _ in range(random.randint(0, 3))]


@app.on_event("startup")
async def startup():
    app.state.model = FakeModel() if FAKE else RealModel()
    app.state.pool = ThreadPoolExecutor(max_workers=POOL_SIZE)
    app.state.served = 0
    print(f"model service ready (fake={FAKE}, threads={POOL_SIZE})")


@app.post("/detect")
async def detect(request: Request):
    jpeg = await request.body()
    loop = asyncio.get_running_loop()
    detections = await loop.run_in_executor(
        app.state.pool, app.state.model.infer, jpeg)
    app.state.served += 1
    return {"detections": detections, "count": len(detections)}


@app.get("/health")
async def health():
    return {"ok": True, "fake": FAKE, "served": app.state.served}

Ese servicio hace algo que vas a entender del todo en el paso 9, y que conviene que notes ya: mete el cálculo en un ThreadPoolExecutor. El modelo es cómputo puro, así que el servicio lo aparta a un pool de hilos para no bloquear su propio event loop.

Es la misma trampa que vas a sufrir en carne propia dentro de unos pasos, resuelta aquí porque este servicio no es la lección.

El lanzador

En vez de abrir dos terminales para los dos servicios, un pequeño lanzador los levanta juntos y los apaga juntos. Crea services.py:

services.py python
import os
import signal
import subprocess
import sys
import time

PROCS = [
    ("cameras", ["uvicorn", "svc_cameras:app", "--port", "8001", "--log-level", "warning"]),
    ("model", ["uvicorn", "svc_model:app", "--port", "8002", "--log-level", "warning"]),
]


def main():
    env = os.environ.copy()
    if "--fake" in sys.argv:
        env["FAKE_MODEL"] = "1"
    if "--cameras" in sys.argv:
        env["N_CAMERAS"] = sys.argv[sys.argv.index("--cameras") + 1]

    running = []
    for name, cmd in PROCS:
        print(f"starting {name}: {' '.join(cmd)}")
        running.append((name, subprocess.Popen(cmd, env=env)))

    def shutdown(*_):
        for name, proc in running:
            proc.send_signal(signal.SIGINT)
        for name, proc in running:
            proc.wait()
        sys.exit(0)

    signal.signal(signal.SIGINT, shutdown)
    signal.signal(signal.SIGTERM, shutdown)

    print("\ncameras   http://127.0.0.1:8001")
    print("model     http://127.0.0.1:8002/health")
    print("dashboard http://127.0.0.1:8001\n")

    while True:
        for name, proc in running:
            if proc.poll() is not None:
                print(f"{name} exited with {proc.returncode}")
                shutdown()
        time.sleep(0.5)


if __name__ == "__main__":
    main()

El dashboard

Falta la cara visible. Crea la carpeta static/ y dentro el archivo index.html:

static/index.html html
<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<title>pipeline dashboard</title>
<style>
  body { background:#0d1117; color:#c9d1d9; font:14px ui-monospace,monospace; margin:0; padding:16px; }
  h1 { font-size:15px; font-weight:600; margin:0 0 14px; color:#79c0ff; }
  #grid { display:grid; grid-template-columns:repeat(auto-fill,minmax(240px,1fr)); gap:10px; }
  .cam { position:relative; background:#000; border:1px solid #21262d; border-radius:4px; overflow:hidden; }
  .cam img { width:100%; display:block; }
  .cam canvas { position:absolute; inset:0; width:100%; height:100%; }
  .cam span { position:absolute; top:4px; left:6px; font-size:11px; background:#000a; padding:1px 5px; border-radius:2px; }
  #feed { margin-top:16px; border-top:1px solid #21262d; padding-top:10px; max-height:200px; overflow-y:auto; font-size:12px; color:#8b949e; }
  #feed div { padding:1px 0; }
</style>
</head>
<body>
<h1>detections</h1>
<div id="grid"></div>
<div id="feed"></div>
<script>
const grid = document.getElementById("grid");
const feed = document.getElementById("feed");
const canvases = {};

fetch("/cameras").then(r => r.json()).then(({cameras}) => {
  for (const cam of cameras) {
    const box = document.createElement("div");
    box.className = "cam";
    box.innerHTML = `<img src="/stream/${cam}"><canvas></canvas><span>${cam}</span>`;
    grid.appendChild(box);
    canvases[cam] = box.querySelector("canvas");
  }
});

function draw(cam, detections) {
  const c = canvases[cam];
  if (!c) return;
  c.width = c.clientWidth; c.height = c.clientHeight;
  const ctx = c.getContext("2d");
  ctx.clearRect(0, 0, c.width, c.height);
  ctx.strokeStyle = "#3fb950"; ctx.lineWidth = 2;
  ctx.fillStyle = "#3fb950"; ctx.font = "11px monospace";
  for (const d of detections || []) {
    const [x1, y1, x2, y2] = d.box;
    const x = x1 * c.width, y = y1 * c.height;
    ctx.strokeRect(x, y, (x2 - x1) * c.width, (y2 - y1) * c.height);
    ctx.fillText(`${d.label} ${d.conf}`, x, Math.max(10, y - 3));
  }
}

new EventSource("/events/stream").onmessage = (e) => {
  const ev = JSON.parse(e.data);
  draw(ev.cam_id, ev.detections);
  const line = document.createElement("div");
  line.textContent = `${ev.cam_id}  seq=${ev.seq}  count=${ev.count}`;
  feed.prepend(line);
  while (feed.childNodes.length > 60) feed.lastChild.remove();
};
</script>
</body>
</html>

Vale la pena entender qué hace esa página, porque explica algo del diseño: el video no pasa por tu pipeline. El <img src="/stream/cam0"> va directo del servicio de cámaras al navegador, y un MJPEG dentro de un img es simplemente una imagen que nunca termina de cargar.

Lo único que aporta tu pipeline son las cajas, que llegan por EventSource y se dibujan encima en un canvas. Dos caminos independientes que se juntan en la pantalla.

Arráncalo

Con los cuatro archivos creados y los videos generados:

terminal 1 bash
python services.py --fake --cameras 4

Deja esa terminal abierta el resto del taller. Vas a trabajar siempre con dos: esta con los servicios, y otra donde escribirás y ejecutarás tu pipeline.

Ahora abre http://127.0.0.1:8001 en el navegador.

Checkpoint · Cuatro cámaras transmitiendo

Deberías ver cuatro recuadros de video reproduciéndose a la vez, etiquetados cam0 a cam3, y debajo un panel de eventos vacío.

Está vacío a propósito: las cámaras transmiten y el modelo espera, pero todavía no hay nada que los conecte. Ese puente es lo que vas a construir en el resto del taller.

Algo va a fallar... · Comprueba el modelo

Antes de seguir, asegúrate de que el segundo servicio también responde. En una terminal nueva:

curl http://127.0.0.1:8002/health

Si ves algo como {"ok":true,"fake":true,"served":0}, todo está en orden. Si en cambio la terminal se queda colgada o dice que no se pudo conectar, mira la terminal 1.

Cómo resolverlo · Los dos errores típicos aquí

no videos found in videos/: el servicio de cámaras arrancó antes de que existieran los clips. Ejecuta python make_samples.py y vuelve a levantar los servicios.

ModuleNotFoundError: No module named 'uvicorn' o similar: el entorno virtual no está activo en esa terminal, o faltó el pip install. Fíjate si tienes el (venv) delante y repite la instalación del paso 1.

Y si curl no existe en tu Windows, abre http://127.0.0.1:8002/health en el navegador: es lo mismo.