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:
shutdown(wait=False, cancel_futures=True)-- no bloquea el proceso principal- 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.