376 lines
13 KiB
Python
376 lines
13 KiB
Python
"""GPU acceleration helpers for LiDAR pipeline.
|
||
|
||
Provides CuPy/numpy abstraction layer. If CuPy is available and a CUDA GPU
|
||
is detected, array operations are accelerated on the GPU. Otherwise, all
|
||
operations fall back to numpy/scipy on CPU.
|
||
|
||
GPU errors (e.g. in forked subprocesses) are caught gracefully and
|
||
cause an automatic fallback to CPU for the current operation.
|
||
|
||
Multi-GPU support: each worker process sets CUDA_VISIBLE_DEVICES before
|
||
CuPy is imported, so CuPy only sees its assigned GPU. This avoids kernel
|
||
cache incompatibilities that occur with Device.use() switching.
|
||
"""
|
||
|
||
import logging
|
||
import os
|
||
import numpy as np
|
||
from scipy import ndimage
|
||
|
||
logger = logging.getLogger("lidar")
|
||
|
||
# Detect total GPU count via nvidia-smi (no CUDA context created).
|
||
# This must happen before any CUDA_VISIBLE_DEVICES manipulation.
|
||
_NUM_GPUS = 0
|
||
HAS_GPU = False
|
||
_gpu_name = None
|
||
_gpu_mem_gb = 0
|
||
# System-level GPU IDs that are currently visible (after restrict_gpus).
|
||
# Populated by restrict_gpus() or auto-detected at import time.
|
||
_available_gpu_ids: list[int] = []
|
||
|
||
try:
|
||
import subprocess
|
||
_result = subprocess.run(
|
||
['nvidia-smi', '--query-gpu=count,name,memory.total', '--format=csv,noheader,nounits'],
|
||
capture_output=True, text=True, timeout=5
|
||
)
|
||
if _result.returncode == 0:
|
||
_lines = _result.stdout.strip().split('\n')
|
||
_NUM_GPUS = len(_lines)
|
||
# Parse first GPU info for logging
|
||
_parts = _lines[0].split(',')
|
||
if len(_parts) >= 3:
|
||
_gpu_name = _parts[1].strip()
|
||
try:
|
||
_gpu_mem_gb = int(float(_parts[2].strip())) // 1024
|
||
except (ValueError, IndexError):
|
||
pass
|
||
HAS_GPU = True
|
||
# All GPUs are visible by default
|
||
_available_gpu_ids = list(range(_NUM_GPUS))
|
||
except (FileNotFoundError, subprocess.TimeoutExpired, Exception):
|
||
pass
|
||
|
||
# Lazy CuPy initialization — imported only when first needed.
|
||
# This allows CUDA_VISIBLE_DEVICES to be set before CuPy creates
|
||
# a CUDA context, enabling per-process GPU assignment.
|
||
_xp = np # Default: CPU
|
||
_cp = None # cupy module (or None)
|
||
_cp_ndimage = None # cupyx.scipy.ndimage (or None)
|
||
_gpu_initialized = False
|
||
|
||
|
||
def _init_gpu():
|
||
"""Lazily initialize CuPy on first GPU use.
|
||
|
||
Import CuPy only when needed, so CUDA_VISIBLE_DEVICES can be
|
||
set before the CUDA context is created.
|
||
|
||
With CUPY_CUDA_COMPILE_WITH_CACHE=1, JIT-compiled kernels are cached
|
||
on disk and shared across processes. We warm up on ALL visible GPUs
|
||
so workers that later restrict to a single GPU find pre-compiled kernels.
|
||
|
||
Per-GPU warm-up errors are caught individually: if GPU 1 fails to
|
||
compile a kernel (e.g. CUDA_ERROR_NO_BINARY_FOR_GPU on sm_89), we
|
||
record it and continue warming up GPU 0. Workers assigned to a
|
||
failed GPU fall back to CPU automatically.
|
||
"""
|
||
global _xp, _cp, _cp_ndimage, _gpu_initialized, HAS_GPU, _gpu_failed_ids
|
||
if _gpu_initialized:
|
||
return
|
||
_gpu_initialized = True
|
||
_gpu_failed_ids = set()
|
||
|
||
try:
|
||
import cupy as _real_cupy
|
||
import cupyx.scipy.ndimage as _real_cupy_ndimage
|
||
|
||
n_devs = _real_cupy.cuda.runtime.getDeviceCount()
|
||
all_failed = False
|
||
|
||
# Warm up each GPU independently — one failure doesn't kill the others.
|
||
for dev_id in range(n_devs):
|
||
try:
|
||
with _real_cupy.cuda.Device(dev_id):
|
||
props = _real_cupy.cuda.runtime.getDeviceProperties(dev_id)
|
||
name = props['name'].decode() if isinstance(props['name'], bytes) else props['name']
|
||
_test = _real_cupy.array([1.0, 2.0, 3.0], dtype=_real_cupy.float32)
|
||
_result = _real_cupy.sum(_test * _test)
|
||
_ = _result.get()
|
||
del _test, _result
|
||
logger.info(f" GPU {dev_id}: {name} — kernels pré-compilés")
|
||
except Exception as dev_err:
|
||
_gpu_failed_ids.add(dev_id)
|
||
logger.warning(
|
||
f" GPU {dev_id}: échec warm-up ({dev_err.__class__.__name__}: "
|
||
f"{dev_err}) — workers sur ce GPU passeront en CPU"
|
||
)
|
||
|
||
# If ALL visible GPUs failed, disable GPU entirely.
|
||
# This is critical for worker subprocesses that restricted
|
||
# CUDA_VISIBLE_DEVICES to a single GPU before importing CuPy.
|
||
if len(_gpu_failed_ids) >= n_devs:
|
||
all_failed = True
|
||
logger.warning("Tous les GPUs échoués au warm-up — passage en mode CPU")
|
||
|
||
if not all_failed:
|
||
_xp = _real_cupy
|
||
_cp = _real_cupy
|
||
_cp_ndimage = _real_cupy_ndimage
|
||
|
||
# Limit GPU memory pool per worker to avoid OOM when multiple
|
||
# workers share one GPU. Each worker gets at most 3.5 GB (or
|
||
# 50 % of total VRAM on smaller cards).
|
||
try:
|
||
props = _real_cupy.cuda.runtime.getDeviceProperties(0)
|
||
total_mem = props['totalGlobalMem']
|
||
max_worker_mem = min(int(total_mem * 0.5), 3.5 * 1024**3)
|
||
_real_cupy.cuda.set_memory_pool(0, max_worker_mem)
|
||
logger.info(
|
||
f" Pool mémoire GPU limité à "
|
||
f"{max_worker_mem // (1024**3) * 1000 // 1024} MB"
|
||
)
|
||
except Exception:
|
||
pass # pool config is best-effort
|
||
else:
|
||
_xp = np
|
||
_cp = None
|
||
_cp_ndimage = None
|
||
HAS_GPU = False
|
||
|
||
except (ImportError, Exception) as e:
|
||
logger.warning(f"GPU non disponible — mode CPU: {e}")
|
||
_xp = np
|
||
_cp = None
|
||
_cp_ndimage = None
|
||
HAS_GPU = False
|
||
|
||
|
||
def restrict_gpus(gpu_ids: list[int], set_env_var: bool = False):
|
||
"""Restrict which GPUs are visible to the process.
|
||
|
||
By default (set_env_var=False), only records the GPU IDs for
|
||
num_gpus() and logging. Does NOT touch CUDA_VISIBLE_DEVICES
|
||
in the main process because CuPy 13.x JIT compilation (sm_89)
|
||
needs ALL GPUs visible during warm-up to pre-compile kernels.
|
||
|
||
Workers call this with set_env_var=True (via assign_gpu_to_worker)
|
||
after forking, when CuPy hasn't been imported yet in the child.
|
||
|
||
Args:
|
||
gpu_ids: List of system-level GPU indices to make visible.
|
||
set_env_var: If True, actually set CUDA_VISIBLE_DEVICES.
|
||
Default False (safe for main process).
|
||
"""
|
||
global _NUM_GPUS, HAS_GPU, _available_gpu_ids
|
||
if not gpu_ids or not HAS_GPU:
|
||
return
|
||
|
||
# Validate IDs against total GPU count from nvidia-smi
|
||
total_count = _NUM_GPUS or 1
|
||
valid_ids = [gid % total_count for gid in gpu_ids]
|
||
_available_gpu_ids = valid_ids
|
||
_NUM_GPUS = len(valid_ids)
|
||
|
||
if set_env_var:
|
||
os.environ['CUDA_VISIBLE_DEVICES'] = ','.join(str(g) for g in valid_ids)
|
||
logger.info(f"GPU visibles: {_available_gpu_ids}")
|
||
|
||
|
||
def num_gpus():
|
||
"""Return the number of visible (available) GPUs."""
|
||
return _NUM_GPUS
|
||
|
||
|
||
def set_active_gpu(gpu_id):
|
||
"""Set the active GPU for the current process via CUDA_VISIBLE_DEVICES.
|
||
|
||
gpu_id is an index into the currently visible GPU list
|
||
(_available_gpu_ids), not a system-level ID.
|
||
|
||
MUST be called before any GPU operation (to_gpu, etc.) to ensure
|
||
CuPy creates its CUDA context on the correct device. With lazy
|
||
initialization, CuPy is imported AFTER this call, so it only
|
||
sees the assigned GPU.
|
||
|
||
If the assigned GPU failed warm-up (e.g. NO_BINARY_FOR_GPU), this
|
||
function disables GPU for this worker so all operations fall back
|
||
to CPU without crashing.
|
||
|
||
Args:
|
||
gpu_id: 0-based index into the visible GPU list.
|
||
"""
|
||
if not HAS_GPU or _NUM_GPUS <= 1:
|
||
return # Nothing to do for single GPU or no GPU
|
||
|
||
gpu_id = gpu_id % _NUM_GPUS
|
||
|
||
# Map visible-GPU index back to the real system GPU ID
|
||
system_gpu_id = _available_gpu_ids[gpu_id]
|
||
|
||
# If this GPU failed warm-up, disable GPU for this worker
|
||
if system_gpu_id in _gpu_failed_ids:
|
||
logger.warning(
|
||
f" GPU {system_gpu_id} échouée au warm-up — "
|
||
f"worker passe en mode CPU"
|
||
)
|
||
disable_gpu()
|
||
return
|
||
|
||
# Set CUDA_VISIBLE_DEVICES to isolate this worker to one GPU.
|
||
# This MUST happen before CuPy is imported (lazy init).
|
||
# The JIT kernels were already pre-compiled by the main process
|
||
# on all GPUs, so the worker finds them in the cache.
|
||
os.environ['CUDA_VISIBLE_DEVICES'] = str(system_gpu_id)
|
||
|
||
logger.info(f" GPU {system_gpu_id} sélectionnée pour ce worker")
|
||
|
||
|
||
def _gpu_available():
|
||
"""Check if GPU is usable right now (may fail in forked subprocesses)."""
|
||
if not HAS_GPU:
|
||
return False
|
||
try:
|
||
_init_gpu()
|
||
_cp.cuda.runtime.getDevice()
|
||
return True
|
||
except Exception:
|
||
return False
|
||
|
||
|
||
def log_gpu_status():
|
||
"""Log GPU detection result. Called after logging is configured."""
|
||
if _gpu_available():
|
||
# Get actual device name from CuPy (after init)
|
||
try:
|
||
dev = _cp.cuda.Device()
|
||
name = _cp.cuda.runtime.getDeviceProperties(0)['name']
|
||
if isinstance(name, bytes):
|
||
name = name.decode()
|
||
mem_gb = _cp.cuda.runtime.getDeviceProperties(0)['totalGlobalMem'] // (1024 ** 3)
|
||
gpu_info = f"GPU: {name} ({mem_gb} Go VRAM)"
|
||
except Exception:
|
||
gpu_info = f"GPU: {_gpu_name} ({_gpu_mem_gb} Go VRAM)"
|
||
if _NUM_GPUS > 1:
|
||
gpu_info += f" × {_NUM_GPUS}"
|
||
logger.info(gpu_info)
|
||
else:
|
||
logger.info("Pas de GPU — mode CPU uniquement")
|
||
|
||
|
||
def to_gpu(arr):
|
||
"""Send array to GPU if available, otherwise return as float32 numpy.
|
||
|
||
Uses float32 to reduce GPU memory usage. Falls back to CPU if GPU
|
||
is unavailable (e.g. in forked subprocess).
|
||
"""
|
||
if _gpu_available():
|
||
try:
|
||
return _cp.asarray(arr.astype(np.float32))
|
||
except Exception:
|
||
pass # Fall back to CPU
|
||
return arr.astype(np.float32)
|
||
|
||
|
||
def to_cpu(arr):
|
||
"""Bring array back to CPU (numpy). No-op if already on CPU."""
|
||
if _cp is not None and isinstance(arr, _cp.ndarray):
|
||
try:
|
||
return _cp.asnumpy(arr)
|
||
except Exception:
|
||
pass # Already on CPU or GPU error
|
||
return arr
|
||
|
||
|
||
def xp_gaussian_filter(arr, sigma):
|
||
"""Gaussian filter — uses GPU if array is on GPU, CPU otherwise."""
|
||
if _cp is not None and isinstance(arr, _cp.ndarray):
|
||
try:
|
||
return _cp_ndimage.gaussian_filter(arr, sigma)
|
||
except Exception:
|
||
arr = to_cpu(arr)
|
||
return ndimage.gaussian_filter(arr, sigma)
|
||
|
||
|
||
def xp_uniform_filter(arr, size):
|
||
"""Uniform filter — uses GPU if array is on GPU, CPU otherwise."""
|
||
if _cp is not None and isinstance(arr, _cp.ndarray):
|
||
try:
|
||
return _cp_ndimage.uniform_filter(arr, size)
|
||
except Exception:
|
||
arr = to_cpu(arr)
|
||
return ndimage.uniform_filter(arr, size)
|
||
|
||
|
||
def xp_minimum_filter(arr, footprint=None, size=None):
|
||
"""Minimum filter — uses GPU if array is on GPU, CPU otherwise."""
|
||
if _cp is not None and isinstance(arr, _cp.ndarray):
|
||
try:
|
||
return _cp_ndimage.minimum_filter(arr, footprint=footprint, size=size)
|
||
except Exception:
|
||
arr = to_cpu(arr)
|
||
return ndimage.minimum_filter(arr, footprint=footprint, size=size)
|
||
|
||
|
||
def xp_maximum_filter(arr, footprint=None, size=None):
|
||
"""Maximum filter — uses GPU if array is on GPU, CPU otherwise."""
|
||
if _cp is not None and isinstance(arr, _cp.ndarray):
|
||
try:
|
||
return _cp_ndimage.maximum_filter(arr, footprint=footprint, size=size)
|
||
except Exception:
|
||
arr = to_cpu(arr)
|
||
return ndimage.maximum_filter(arr, footprint=footprint, size=size)
|
||
|
||
|
||
def gpu_cleanup():
|
||
"""Free GPU memory. Call between visualizations to prevent OOM."""
|
||
if _cp is not None:
|
||
try:
|
||
_cp.get_default_memory_pool().free_all_blocks()
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def disable_gpu():
|
||
"""Disable GPU acceleration for the rest of this process.
|
||
|
||
Called when a CUDA error indicates the GPU is unusable (e.g.
|
||
CUDA_ERROR_NO_BINARY_FOR_GPU). Falls back to numpy for all
|
||
subsequent operations.
|
||
"""
|
||
global HAS_GPU, _xp, _cp, _cp_ndimage
|
||
if not HAS_GPU:
|
||
return # Already disabled
|
||
logger.warning("GPU désactivé — passage en mode CPU pour la suite du processus")
|
||
HAS_GPU = False
|
||
_xp = np
|
||
_cp = None
|
||
_cp_ndimage = None
|
||
|
||
|
||
def is_gpu_active():
|
||
"""Check if GPU acceleration is currently active.
|
||
|
||
Unlike the HAS_GPU module-level variable (which can go stale if
|
||
imported directly), this always reflects the current runtime state.
|
||
Use this in logging tags and conditional GPU paths.
|
||
"""
|
||
return HAS_GPU
|
||
|
||
|
||
def safe_gpu_call(func, *args, **kwargs):
|
||
"""Call a function with GPU arrays, retrying on CPU if GPU fails.
|
||
|
||
Usage:
|
||
result = safe_gpu_call(generate_svf, dem_file, basename, vis_dir, resolution, shared=shared)
|
||
"""
|
||
try:
|
||
return func(*args, **kwargs)
|
||
except Exception as e:
|
||
err_msg = str(e)
|
||
if _cp is not None and ('CUDA' in err_msg or 'cuda' in err_msg or 'GPU' in err_msg):
|
||
logger.warning(f"Erreur GPU ({e.__class__.__name__}), retry en CPU...")
|
||
disable_gpu()
|
||
return func(*args, **kwargs)
|
||
raise |