←

Pipeline multi-threaded con queue.Queue en Python

Contexto

Necesitaba un pipeline de audio en tiempo real: capturar audio del sistema, transcribir con Whisper, traducir, y mostrar en terminal. Cada paso tiene latencia diferente (la transcripcion es la mas lenta). Si todo corre en un solo thread, el audio se pierde mientras Whisper procesa.

Lo que aprendi

queue.Queue con maxsize conecta threads como tuberias. Cada componente consume de una queue y produce a otra:

import queue
import threading

audio_queue = queue.Queue(maxsize=10)
text_queue = queue.Queue(maxsize=50)
logger_queue = queue.Queue(maxsize=50)
display_queue = queue.Queue(maxsize=50)

capture = AudioCapture(audio_queue)
transcriber = Transcriber(audio_queue, text_queue)
translator = Translator(text_queue, [logger_queue, display_queue])
logger = TranscriptLogger(logger_queue)
display = Display(display_queue)

El patron clave es que el Translator hace fan-out: empuja el mismo resultado a multiples queues de salida. Y si una queue esta llena, lo descarta sin bloquear:

for q in self.output_queues:
    try:
        q.put(result, timeout=1.0)
    except queue.Full:
        pass  # backpressure: no bloquear el traductor

Cada thread usa threading.Event() para saber cuando parar:

class Transcriber:
    def __init__(self, audio_queue, text_queue):
        self._stop_event = threading.Event()

    def run(self):
        while not self._stop_event.is_set():
            try:
                chunk = self.audio_queue.get(timeout=1.0)
            except queue.Empty:
                continue
            # procesar...

    def stop(self):
        self._stop_event.set()

Los workers arrancan como daemon threads con un stagger de 100ms para evitar race conditions al inicio:

workers = [
    threading.Thread(target=capture.run, daemon=True),
    threading.Thread(target=transcriber.run, daemon=True),
    threading.Thread(target=translator.run, daemon=True),
    threading.Thread(target=logger.run, daemon=True),
]
for t in workers:
    t.start()
    time.sleep(0.1)

Por que importa

Este patron escala a cualquier pipeline de procesamiento donde los pasos tienen latencias diferentes. El maxsize en las queues actua como backpressure natural: si Whisper es lento, la queue de audio se llena y el capture espera en vez de acumular memoria infinita.