From bb086e308a8a6a4d590c0f108e1ed495ae640bc6 Mon Sep 17 00:00:00 2001 From: Antoine Jacquin Date: Fri, 18 Sep 2026 21:08:25 +0200 Subject: [PATCH] =?UTF-8?q?Rendre=20effectif=20le=20d=C3=A9lai=20de=20s?= =?UTF-8?q?=C3=A9curit=C3=A9=20de=202=20h=20des=20workers=20parall=C3=A8le?= =?UTF-8?q?s?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- lidar_pipeline/pipeline.py | 49 +++++++++++++++++++++----------------- 1 file changed, 27 insertions(+), 22 deletions(-) diff --git a/lidar_pipeline/pipeline.py b/lidar_pipeline/pipeline.py index 533be08..510deb5 100644 --- a/lidar_pipeline/pipeline.py +++ b/lidar_pipeline/pipeline.py @@ -11,7 +11,7 @@ import logging import multiprocessing import shutil import time -from concurrent.futures import ProcessPoolExecutor, as_completed +from concurrent.futures import ProcessPoolExecutor, as_completed, TimeoutError as FuturesTimeoutError from datetime import datetime from pathlib import Path import subprocess @@ -589,27 +589,32 @@ class LidarArchaeoPipeline: done = 0 t_deadline = time.time() + 7200 try: - for future in as_completed(future_to_file): - if time.time() > t_deadline: - logger.error("Délai dépassé (2h) — annulation des workers restants") - for f in future_to_file: - f.cancel() - break - laz_file = future_to_file[future] - done += 1 - try: - success = future.result() - results[laz_file.name] = success - status = "✓" if success else "✗" - logger.info(f" [{done}/{len(files)}] {status} {laz_file.name}") - if success and self.incremental_index: - self._rebuild_index_incremental() - except Exception as e: - logger.error(f" [{done}/{len(files)}] ✗ {laz_file.name}: {e}") - logger.debug(f" Traceback:", exc_info=True) - report_event(self.output_dir, _file_basename(laz_file), - "tile", "fail", detail=str(e)) - results[laz_file.name] = False + # timeout= : sans lui, as_completed bloque entre deux + # complétions et le délai de 2 h n'est jamais évalué si + # aucun worker ne rend la main (run figé pour toujours). + try: + futures_iter = as_completed(future_to_file, + timeout=max(1.0, t_deadline - time.time())) + for future in futures_iter: + laz_file = future_to_file[future] + done += 1 + try: + success = future.result() + results[laz_file.name] = success + status = "✓" if success else "✗" + logger.info(f" [{done}/{len(files)}] {status} {laz_file.name}") + if success and self.incremental_index: + self._rebuild_index_incremental() + except Exception as e: + logger.error(f" [{done}/{len(files)}] ✗ {laz_file.name}: {e}") + logger.debug(f" Traceback:", exc_info=True) + report_event(self.output_dir, _file_basename(laz_file), + "tile", "fail", detail=str(e)) + results[laz_file.name] = False + except FuturesTimeoutError: + logger.error("Délai dépassé (2h) — annulation des workers restants") + for f in future_to_file: + f.cancel() except KeyboardInterrupt: logger.info("Interruption — annulation des travaux en cours...") for f in future_to_file: