Ingesta asíncrona: transmitir, no leer
La inferencia ya es I/O de red. Falta la entrada: la ingesta sigue leyendo con urllib, que es bloqueante, y mientras espera un frame de cam0 nadie atiende a las otras tres cámaras.
Qué es exactamente un stream MJPEG
Antes del código conviene entender qué está mandando el servicio de cámaras, porque el formato explica la forma del código.
MJPEG sobre HTTP es una respuesta que nunca termina. El servidor abre la conexión y va enviando imágenes JPEG una detrás de otra, separadas por una marca, cada una con una cabecera que dice cuántos bytes ocupa:
--frame
Content-Type: image/jpeg
Content-Length: 24601
<24601 bytes de JPEG>
--frame
Content-Type: image/jpeg
Content-Length: 24384
<24384 bytes de JPEG>
...
Por eso no puedes hacer un response.read() y esperar el final: no hay final. Tienes que ir leyendo por partes, y ahí es donde aiohttp encaja perfecto. Crea sources.py:
async def frames(self):
async with self.session.get(self.url) as response:
reader = response.content
while True:
header = await reader.readuntil(b"\r\n\r\n")
length = content_length(header)
jpeg = await reader.readexactly(length)
await reader.readexactly(2)
yield jpeg
Lee ese bucle línea por línea, porque es el corazón de la ingesta:
await reader.readuntil(b"\r\n\r\n")lee hasta el final de la cabecera. Cede el control mientras los bytes llegan.content_length(header)saca de esa cabecera cuántos bytes mide el JPEG que viene.await reader.readexactly(length)lee exactamente esa cantidad. Otro punto donde cede.await reader.readexactly(2)se traga el salto de línea que separa una parte de la siguiente.yield jpegentrega el frame a quien lo esté consumiendo.
Ese yield dentro de un async def lo convierte en un generador asíncrono: algo que puedes recorrer con async for y que cede el control en cada vuelta. Es la forma natural de modelar una cámara, que es justo eso, una fuente infinita de frames que llegan cuando llegan.
Y lo más importante de todo: no hay ni un to_thread ni un executor en toda la clase. Cada frame llega por un await sobre un socket.
Cuatro streams, un hilo
Para comprobarlo, crea test_source.py:
python test_source.py
cameras: ['cam0', 'cam1', 'cam2', 'cam3']
4 streams read at once: 100 frames in 0.99s
one stream alone would have taken about 1.0s for 25 frames at 25 fps
Checkpoint · Las otras tres salieron gratis
Cien frames de cuatro cámaras en un segundo. Y la línea de abajo es la que hay que leer dos veces: una sola cámara habría tardado ese mismo segundo en dar sus veinticinco frames, porque transmite a 25 por segundo y no puede ir más rápido.
Las otras tres cámaras no costaron tiempo adicional. Mientras el programa espera un frame de cam0, los frames de cam1, cam2 y cam3 van llegando a sus propios sockets.
Analogía · Cuatro ollas en el mismo fuego
Cocinar arroz toma veinte minutos, y ese tiempo no lo puedes acortar. Pero si pones cuatro ollas a la vez, las cuatro estarán listas en esos mismos veinte minutos, no en ochenta.
El tiempo de espera no se suma cuando esperas varias cosas a la vez. Y eso vale igual para el arroz que para los sockets.
Dónde estamos
Ya tienes las dos puntas del pipeline en su forma definitiva:
- La ingesta transmite MJPEG con
aiohttp, cediendo el control en cada lectura. - La inferencia manda una petición HTTP al servicio del modelo y espera.
Las dos son I/O puro sobre sockets, y las dos pueden ocurrir cientos de veces a la vez en un solo hilo. Lo que falta es conectarlas, y hacerlo de forma que la ingesta no quede a merced de lo lenta que sea la inferencia. Ese es el trabajo de las colas.
Ojo
Detente un momento aquí, porque es el punto donde mucha gente se lanza a conectar las dos etapas con una llamada directa: leer un frame y mandarlo a detectar en el mismo bucle. Eso funciona, y vuelve a atarte a la etapa más lenta. En el siguiente paso vas a ver por qué la cola es la pieza que faltaba.