Guzmán D. Darío Senior Python Developer English Contratar
Taller autoguiado

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

Paso 22 de 23
Paso 22 de 23

Bonus: cientos de streams con GStreamer, y el supervisor

Terminaste el bonus de RTSP con un techo: cada cámara IP consume un hilo decodificando H.264 sin parar, y con cuarenta cámaras tu máquina se queda sin CPU. Este capítulo es sobre qué se hace cuando llegas ahí.

La respuesta incomoda un poco, y por eso vale la pena decirla: a esa escala, el video sale de Python.

GStreamer ya resolvió lo que acabas de escribir

GStreamer es un framework de medios con veinte años encima. Está pensado exactamente para lo que tú acabas de construir a mano, y en la ruta del video lo hace mejor: decodificación por hardware, buffers sin copias, sincronización por timestamps, y las rarezas de RTSP de cada fabricante ya resueltas.

Lo incómodo es lo parecido que resulta a tu pipeline:

lo que escribiste GStreamer
asyncio.Queue(maxsize=32) queue max-size-buffers=32
backpressure al llenarse pads bloqueantes, por defecto
put_or_drop queue leaky=downstream
ocho trabajadores de inferencia la frontera de un queue es frontera de hilo
el monitor de longitud de colas pad probes, gst-shark
las cuatro cámaras entrando a una cola funnel, nvstreammux

Cada idea del taller, disponible como configuración.

Ojo

Entonces, ¿perdiste el tiempo construyéndolo a mano? No, y la razón es la columna de la izquierda: ahora entiendes qué significa cada opción de la derecha. Poner leaky=downstream sin haber visto crecer una cola sin límite es copiar una línea de un foro. Después del paso 10 y del paso 11, es una decisión de política que sabes justificar.

Lo que se pierde al usar GStreamer es ver tu propia concurrencia, que es justo para lo que lo escribimos hoy.

Cómo se ve a esa escala

Un pipeline de producción para cientos de cámaras se parece a esto:

el pipeline de medios text
uridecodebin x25 ─┐
                  ├─> nvstreammux ─> nvinfer ─> nvtracker ─> metadata
uridecodebin x25 ─┘   batch-size=25   batch-size=25

Tres cosas que aprender de ese dibujo:

El techo es la decodificación, no la inferencia. Al revés de lo que asumiría cualquiera. Por eso la primera optimización real suele ser drop-frame-interval=5: analizar uno de cada cinco frames. Para detectar personas o carros, cinco veces por segundo suele bastar, y es el 5x más barato que vas a conseguir.

Los frames viajan en lotes. nvstreammux junta 25 cámaras en un solo lote y el modelo lo procesa de una vez. Es la misma idea que tus ocho trabajadores concurrentes, llevada al hardware.

Veinticinco streams por proceso, no trescientos. Y esa decisión no es técnica, es de fallos: un stream corrupto que tumbe el proceso se lleva por delante a veinticuatro compañeros, no a doscientos noventa y nueve.

Y ahí es donde vuelve tu Python

Si el video se procesa en veinte procesos de GStreamer, alguien tiene que repartir las cámaras entre ellos, vigilarlos, reiniciar el que se muera y sacarlos de servicio ordenadamente cuando despliegues.

Ese alguien es un supervisor. Y es exactamente el mismo asyncio que llevas todo el taller escribiendo:

supervisor.py python
async def supervise(shard_id, cameras):
    backoff = 1.0
    while not stopping.is_set():
        proc = await asyncio.create_subprocess_exec(
            "python", "shard.py", "--id", str(shard_id), *cameras)
        code = await proc.wait()
        if stopping.is_set():
            return
        print(f"shard {shard_id} exited {code}, restarting in {backoff:.0f}s")
        await asyncio.sleep(backoff)
        backoff = min(backoff * 2, 30)


async def main():
    shards = chunk(await discover_cameras(), SHARD_SIZE)
    async with asyncio.TaskGroup() as tg:
        for i, cams in enumerate(shards):
            tg.create_task(supervise(i, cams))

Reconoce las piezas, porque son todas tuyas:

Cuatro shards de 25 streams bajo un supervisor asyncio. Cuando uno muere, los demás siguen entregando, y el supervisor lo reinicia con backoff creciente hasta que vuelve.
Cuatro shards de 25 streams bajo un supervisor asyncio. Cuando uno muere, los demás siguen entregando, y el supervisor lo reinicia con backoff creciente hasta que vuelve.

Analogía · El jefe de sala

En la cocina industrial ya nadie corta las verduras a mano: hay máquinas que lo hacen mil veces más rápido. Pero sigue haciendo falta alguien que reparta el trabajo entre las máquinas, que se dé cuenta cuando una se para, que la reinicie y que decida qué hacer mientras tanto.

Ese jefe de sala no toca un cuchillo en todo el día, y sin él la cocina se detiene igual.

Los píxeles salen de Python. El control, no.

Esa es la frase que resume este capítulo, y quizá el taller entero.

Repartir cámaras entre shards, reiniciar lo que se murió, aplicar backoff, exponer endpoints de salud, vaciar ordenadamente al desplegar: todo eso sigue siendo tuyo, y todo se escribe con las mismas cinco herramientas que ya conoces. TaskGroup, cancelación, timeouts, backoff y colas acotadas, una capa más arriba.

De hecho ya lo has visto funcionando: el services.py que llevas todo el taller dejando abierto en la terminal 1 es una versión pequeña de exactamente eso, un proceso Python que levanta y vigila otros dos.

Checkpoint

No hay nada que ejecutar en este capítulo, y es a propósito: montar DeepStream o una GPU no cabe en un taller. Lo que sí te llevas es el mapa.

Si algún día tu proyecto pasa de diez cámaras a trescientas, ya sabes dónde está la frontera: el video se va a un framework de medios, y tu asyncio se queda con el plano de control.