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

Apagado limpio

Todo el taller hemos parado los programas con Ctrl-C sin pensarlo. Ahora vamos a mirar qué pasa realmente cuando haces eso, porque en un pipeline con colas llenas la respuesta es fea.

Ctrl-C es una señal, no una excepción

Cuando pulsas Ctrl-C, el sistema operativo manda una señal SIGINT a tu proceso, y Python la convierte en un KeyboardInterrupt en el punto exacto donde estuviera ejecutando. Puede ser en mitad de un put, mientras un trabajador tiene un frame en la mano, o entre dos líneas de tu limpieza.

En un script pequeño da igual. En un pipeline significa perder lo que estuviera en vuelo y dejar sockets a medio cerrar.

La solución es atrapar la señal tú, y convertirla en una petición educada:

main.py python
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
    loop.add_signal_handler(sig, stop.set)

asyncio.Event es una banderita que las corrutinas pueden esperar. Cuando llega la señal, en vez de reventar por donde sea, simplemente se levanta la bandera y tu programa decide qué hacer y en qué orden.

Sin esto, Ctrl-C te corta por donde sea. Con esto, tú eliges por dónde.

Parar en orden: primero la entrada, al final la salida

Con la bandera levantada, el apagado sigue el sentido del flujo:

main.py python
for task in ingest:
    task.cancel()
await asyncio.gather(*ingest, return_exceptions=True)

async with asyncio.timeout(DRAIN_DEADLINE):
    await drain(infer_q, publish_q, infer_workers, publish_workers)

Primero se corta la entrada: se cancelan las tareas de ingesta, así que dejan de entrar frames nuevos. Después se deja que lo que ya está dentro termine su recorrido.

El orden importa y es el del flujo: cámaras, inferencia y, de última, publicación. Si apagaras primero la publicación, todo lo que quedaba en la segunda cola se perdería.

Y ese asyncio.timeout(DRAIN_DEADLINE) es la dosis de realismo: le das un plazo al vaciado y, si no termina a tiempo, dejas de ser educado y cortas. Un apagado sin plazo es un programa que no se apaga.

El centinela, etapa por etapa

¿Y cómo le dices a un trabajador que ya no vengan más frames? No cancelándolo: mandándole un mensaje especial por la misma cola. Eso es el centinela, y ya lo viste en pipeline.py sin que le prestáramos atención:

pipeline.py y main.py, lado a lado python
# lo que el worker ya esperaba, desde el paso 12
SENTINEL = None

while True:
    msg = await in_q.get()
    if msg is SENTINEL:
        break
    ...

# lo que el apagado manda
for _ in infer_workers:
    await infer_q.put(SENTINEL)
await asyncio.gather(*infer_workers)

for _ in publish_workers:
    await publish_q.put(SENTINEL)
await asyncio.gather(*publish_workers)

Dos detalles que valen oro:

Uno por trabajador. Cada worker saca un centinela y se va. Si mandaras uno solo con ocho trabajadores, siete se quedarían dormidos para siempre en el get.

¿Y por qué no cancelarlos, como hiciste con la ingesta? Porque cancelar bota el frame que el trabajador tenga en la mano, y ese frame quizá ya costó una petición al modelo. La ingesta se cancela porque lo que produce es reemplazable: vienen más frames. Los trabajadores se despiden con centinela porque lo suyo ya está pagado.

Algo va a fallar... · Prueba a apagarlo

Arranca el pipeline y déjalo procesar unos segundos:

python main.py

Ahora pulsa Ctrl-C una sola vez y observa la salida sin tocar nada más. Cuenta cuánto tarda en volver la línea de comandos y mira los números finales de stats.

Cómo resolverlo · Ese par de segundos son el vaciado

Lo que viste no es lentitud: es el apagado haciendo su trabajo. Entre que pulsaste Ctrl-C y el programa terminó, el pipeline canceló la ingesta, mandó los centinelas, esperó a que los ocho trabajadores terminaran lo que tenían en vuelo y cerró la sesión HTTP.

Fíjate en el stats final: los números de inferred y published cuadran. Nada quedó a medias.

Si lo comparas con matar el proceso a lo bruto, la diferencia es que ahí perderías todo lo que estuviera en las dos colas, y el servicio del otro lado se quedaría con conexiones abiertas hasta que expiren.

Dos descuidos que se pagan caro

Uno. Cancelar una tarea hace que lance CancelledError. Si el gather no lleva return_exceptions=True, ese error sube de inmediato y aborta el resto del apagado: nunca se vacían las colas, y las otras tareas siguen corriendo sin nadie que las espere.

Dos. Si capturas CancelledError y no lo relanzas, la tarea se niega a morir y el apagado espera por ella hasta que venza el plazo. Es el mismo error del paso 14, y aparece aquí porque es donde duele.

Cierra los recursos en finally: es lo único que corre igual si la tarea termina bien o si la cancelan.

Checkpoint

Tu pipeline ya se apaga sin perder trabajo en vuelo y sin dejar conexiones colgando. Con eso deja de ser un experimento y empieza a parecerse a algo que puedes dejar corriendo.