Rendre effectif le délai de sécurité de 2 h des workers parallèles
This commit is contained in:
@ -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:
|
||||
|
||||
Reference in New Issue
Block a user