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 14 de 23
Paso 14 de 23

El montaje: el pipeline completo

Tienes todas las piezas. La ingesta transmite por sockets, la inferencia vive detrás de un servicio, las colas desacoplan y ya elegiste la política de cada una. Toca armarlo.

Tres etapas, todas esperando sockets

Abre pipeline.py. Las tres etapas del pipeline son tres corrutinas, y ninguna pasa de diez líneas:

pipeline.py python
async def ingest_source(source, out_q, stats):
    async for jpeg in source.frames():
        msg = FrameMsg(cam_id=source.cam_id, seq=stats["read"], jpeg=jpeg)
        stats["read"] += 1
        if not put_or_drop(out_q, msg):
            stats["dropped"] += 1

async def inference_worker(in_q, out_q, detector, stats):
    while True:
        msg = await in_q.get()
        if msg is SENTINEL:
            break
        result = await detector.detect(msg.jpeg)
        msg.detections = result["detections"]
        msg.jpeg = b""
        await out_q.put(msg)
        stats["inferred"] += 1

async def publish_worker(in_q, publisher, stats):
    while True:
        msg = await in_q.get()
        if msg is SENTINEL:
            break
        await publisher.publish(msg)
        stats["published"] += 1

Fíjate en tres detalles que se pueden pasar por alto:

Las políticas están donde las decidiste. La ingesta usa put_or_drop y cuenta lo que bota. La inferencia usa await out_q.put(msg) y bloquea, porque ese resultado ya costó una petición al modelo.

msg.jpeg = b"" después de inferir. Ahí estás soltando los píxeles: una vez que tienes las detecciones, la imagen no le hace falta a nadie más. Lo que viaja a la segunda cola es metadata de unos pocos bytes. Ese detalle es lo que permite aguantar 64 cámaras sin comerse la memoria.

if msg is SENTINEL: break es la puerta de salida ordenada, y la usaremos en el paso 15. Por ahora, ignórala.

Quién arranca a quién

Ahora main.py, donde todo cobra vida:

main.py python
infer_q = asyncio.Queue(maxsize=QUEUE_SIZE)
publish_q = asyncio.Queue(maxsize=QUEUE_SIZE)

infer_workers = [
    asyncio.create_task(inference_worker(infer_q, publish_q, detector, stats))
    for _ in range(INFER_WORKERS)]        # 8
publish_workers = [
    asyncio.create_task(publish_worker(publish_q, publisher, stats))
    for _ in range(PUBLISH_WORKERS)]      # 2
ingest = [
    asyncio.create_task(ingest_source(MjpegSource(cam, session), infer_q, stats))
    for cam in cameras]                   # 4, one per camera

Lee el orden con atención, porque no es casual: primero los consumidores, después la ingesta. Si arrancaras al revés, las cámaras empezarían a llenar la cola sin que hubiera nadie del otro lado sacando frames.

Cuenta las tareas: ocho de inferencia, dos de publicación, cuatro de ingesta. Catorce tareas, dos colas, un solo hilo.

Y esos ocho trabajadores de inferencia no son ocho hilos. Son ocho corrutinas que, la mayor parte del tiempo, están paradas en un await esperando a que el servicio del modelo conteste. Ocho peticiones en vuelo a la vez, y al event loop le cuestan ocho objetos en memoria.

Las cuatro cámaras entran a la misma cola, tres trabajadores sacan frames y los mandan a detectar, y los resultados salen por la segunda cola hacia la publicación. Ninguna etapa espera a la otra.
Las cuatro cámaras entran a la misma cola, tres trabajadores sacan frames y los mandan a detectar, y los resultados salen por la segunda cola hacia la publicación. Ninguna etapa espera a la otra.

Ejecútalo

terminal 2 bash
python main.py

Y ahora vete al navegador, al dashboard que tienes abierto desde el paso 2.

Checkpoint · El panel de eventos vivo

Las cuatro cámaras siguen transmitiendo, pero ahora aparecen cajas verdes dibujadas sobre el video y el panel de abajo se llena de eventos de las cuatro a la vez, mezclados.

Compáralo con el paso 3, cuando los eventos llegaban de una cámara y luego de la siguiente. Esa mezcla es la prueba visual de que las cuatro se están procesando concurrentemente.

El número

Deja que corra sus diez segundos y mira la salida:

salida text
4 cameras, 8 inference workers

ran 10.0s
stats: {'read': 972, 'dropped': 0, 'inferred': 972, 'published': 972}
throughput: 96.8 frames/s inferred

Checkpoint · Compara con tu baseline

Saca el número que anotaste en el paso 3. Aquí está la comparación completa:

baseline: 12.1 frames/s 160 frames, 0 dropped pipeline: 96.8 frames/s 972 frames, 0 dropped

Ocho veces el rendimiento, y ni un frame perdido. Mismo hardware, mismo modelo, mismo trabajo. Lo único que cambió es que tu programa dejó de estar parado esperando.

Y hay un dato más que merece atención: cuatro cámaras a 25 frames por segundo piden 100 por segundo. Estás procesando 96.8. Vas al día con la realidad, que era exactamente el objetivo del taller.

Lo que hizo la diferencia

No fue el paralelismo: sigues teniendo un solo hilo. Fue esto:

Tómate un momento aquí. El resto del taller es hacer que este pipeline sobreviva al mundo real: a la escala, a las cámaras que fallan, a los apagados y a no saber dónde está el cuello de botella.