live

ThreadPoolExecutor y FuturesTimeoutError: el bloqueo que no esperabas

concurrent.futures.as_completed(timeout=N) lanza FuturesTimeoutError cuando algun futuro no termina a tiempo -- pero el bloque with ThreadPoolExecutor sigue bloqueando en el __exit__ porque shutdown(wait=True) es el default. Este es el patron correcto para evitarlo.

El problema

Tienes un pool de threads con un timeout global. Cuando alguno no responde a tiempo, capturas la excepcion y marcas los pendientes como fallidos:

with ThreadPoolExecutor(max_workers=4) as executor:
    futures = {executor.submit(tarea_lenta, item): item for item in items}
    try:
        for future in as_completed(futures, timeout=20):
            results.append(future.result())
    except FuturesTimeoutError:
        for future, item in futures.items():
            if not future.done():
                results.append({"error": "timeout"})
        # aqui ya tienes tus resultados

Parece correcto. El problema esta en lo que pasa despues del bloque except: al salir del with, Python llama executor.__exit__, que internamente llama executor.shutdown(wait=True).

Resultado: aunque hayas capturado el FuturesTimeoutError y procesado los pendientes, el proceso espera a que todos los threads activos terminen antes de continuar.

Por que importa

Si cada tarea hace llamadas bloqueantes secuenciales con subprocess.run(timeout=10), un thread pendiente puede tardar hasta 50 segundos en terminar. Con varios workers y mala suerte, son varios minutos de espera extra despues del timeout que creias que acotaba el tiempo total.

El caso real: infra_report_job en Dagster lanzaba 4 threads SSH a nodos remotos con un timeout de 20s via as_completed. Si un nodo estaba caido, el job tardaba 20s (timeout declarado) + hasta 50s (shutdown bloqueante silencioso) = ~70s. El log mostraba el FuturesTimeoutError capturado correctamente, pero el proceso no continuaba.

La solucion

Gestionar el executor manualmente en lugar de usar el context manager, y llamar a shutdown con los parametros correctos:

executor = ThreadPoolExecutor(max_workers=4)
try:
    futures = {executor.submit(tarea_lenta, item): item for item in items}
    try:
        for future in as_completed(futures, timeout=20):
            results.append(future.result())
    except FuturesTimeoutError:
        for future, item in futures.items():
            if not future.done():
                results.append({"error": "timeout"})
finally:
    executor.shutdown(wait=False, cancel_futures=True)

cancel_futures=True (Python 3.9+) cancela los futuros en cola que aun no han empezado. Los threads ya en ejecucion no se pueden cancelar -- Python no tiene ese mecanismo -- pero wait=False hace que shutdown retorne inmediatamente y los threads terminen en background sin bloquear el proceso principal.

Los limites del fix

wait=False no mata los threads, solo deja de esperar. Si las tareas internas hacen I/O bloqueante sin timeout propio, siguen vivos en background potencialmente para siempre.

La solucion completa tiene dos capas:

  1. shutdown(wait=False, cancel_futures=True) -- no bloquea el proceso principal
  2. Cada tarea interna tiene su propio timeout -- acota el peor caso de los threads en background

En el caso de infra_report_job, las llamadas SSH usan subprocess.run(..., timeout=10), lo que garantiza que ningun thread pendiente vive mas de ~60 segundos aunque el proceso principal ya haya continuado.

Por que el default es wait=True

Es la opcion segura para el caso general: asegura que todos los recursos se liberan antes de continuar. El problema es que "seguro" y "correcto para mi caso de uso" no son lo mismo cuando has disenado explicitamente un timeout global.

La documentacion de Python lo menciona, pero es facil de pasar por alto: el context manager de ThreadPoolExecutor siempre espera a que todos los futures terminen, independientemente de si has capturado excepciones dentro del bloque. El timeout de as_completed acota cuanto esperas para iterar, no cuanto tarda el with en cerrarse.

aqui cualquier cosa mientras cuadramos el logo