Compare commits
2 Commits
translate-
...
batch-time
| Author | SHA1 | Date | |
|---|---|---|---|
| d0dc8d90e9 | |||
| 6e8580138c |
@ -84,7 +84,7 @@
|
||||
- `_d8_accumulate_numba` uses numba with `argsort` top-down sweep. Python fallback exists.
|
||||
- Ray-tracing (SVF, openness): processes one direction at a time to limit VRAM. Auto-falls back to CPU on OOM via `_ray_trace_horizons`.
|
||||
- Multi-resolution: 0.5 m has no filename suffix whatever its position; every other resolution (including the 0.2 m default) uses a `_r0p2` style suffix. Ground classification done once, shared across resolutions.
|
||||
- `ProcessPoolExecutor` has a 2-hour wall-clock safety timeout (prevents indefinite hang from stuck workers).
|
||||
- `ProcessPoolExecutor` batch runs are unlimited by default; set `LIDAR_BATCH_TIMEOUT` (seconds) to cap a batch's wall clock and cancel the remaining workers past it.
|
||||
|
||||
### Numba usage pattern
|
||||
- Defined at function scope with `@njit(cache=True)` — first call compiles (~2-3s), subsequent calls hit disk cache.
|
||||
|
||||
@ -38,6 +38,8 @@ services:
|
||||
# Generations started from a remote map use the GPU
|
||||
- LIDAR_GPU=1
|
||||
- LIDAR_WORKERS=auto
|
||||
# Wall-clock cap for one batch of tiles, in seconds (unset or 0 = unlimited)
|
||||
- LIDAR_BATCH_TIMEOUT=21600
|
||||
# Workers per GPU capped by the free VRAM when the run starts
|
||||
# ((free - reserve) / per-worker peak); the excess runs on the CPU.
|
||||
# Peak estimated at 2048 MiB: adjust after measuring (nvidia-smi during a run).
|
||||
|
||||
@ -47,6 +47,22 @@ def resolve_workers(value):
|
||||
return 1
|
||||
|
||||
|
||||
def _batch_timeout_s():
|
||||
"""Wall-clock budget for one batch, from LIDAR_BATCH_TIMEOUT (seconds).
|
||||
|
||||
Unset, empty or 0 means unlimited (default). A negative or non-numeric
|
||||
value is ignored with a warning.
|
||||
"""
|
||||
raw = (os.environ.get("LIDAR_BATCH_TIMEOUT") or "").strip()
|
||||
if not raw:
|
||||
return 0.0
|
||||
try:
|
||||
return max(0.0, float(raw))
|
||||
except ValueError:
|
||||
logger.warning(f"LIDAR_BATCH_TIMEOUT={raw!r} is not a number — ignored (unlimited)")
|
||||
return 0.0
|
||||
|
||||
|
||||
def _file_basename(path):
|
||||
"""Extract base name from a LAZ/LAS file, removing all known extensions.
|
||||
|
||||
@ -852,15 +868,17 @@ class LidarArchaeoPipeline:
|
||||
for laz_file in files
|
||||
}
|
||||
done = 0
|
||||
t_deadline = time.time() + 7200
|
||||
budget_s = _batch_timeout_s()
|
||||
deadline = time.time() + budget_s if budget_s > 0 else None
|
||||
try:
|
||||
# timeout=: without it, as_completed blocks between two
|
||||
# completions and the 2 h deadline is never evaluated if
|
||||
# no worker returns (run stuck forever).
|
||||
# timeout= lets the deadline be checked between two
|
||||
# completions; without it a stuck worker would block
|
||||
# as_completed forever. No deadline (LIDAR_BATCH_TIMEOUT
|
||||
# unset or 0) → wait indefinitely.
|
||||
as_completed_kwargs = ({"timeout": max(1.0, deadline - time.time())}
|
||||
if deadline is not None else {})
|
||||
try:
|
||||
futures_iter = as_completed(future_to_file,
|
||||
timeout=max(1.0, t_deadline - time.time()))
|
||||
for future in futures_iter:
|
||||
for future in as_completed(future_to_file, **as_completed_kwargs):
|
||||
laz_file = future_to_file[future]
|
||||
done += 1
|
||||
try:
|
||||
@ -877,7 +895,8 @@ class LidarArchaeoPipeline:
|
||||
"tile", "fail", detail=str(e))
|
||||
results[laz_file.name] = False
|
||||
except FuturesTimeoutError:
|
||||
logger.error("Timeout exceeded (2 h) — cancelling remaining workers")
|
||||
logger.error(f"Batch exceeded LIDAR_BATCH_TIMEOUT "
|
||||
f"({budget_s:g} s) — cancelling remaining workers")
|
||||
for f in future_to_file:
|
||||
f.cancel()
|
||||
except KeyboardInterrupt:
|
||||
|
||||
@ -440,3 +440,21 @@ class TestResolveWorkers:
|
||||
for _ in range(4): # 4th call: empty queue
|
||||
pipeline._init_worker_slot(q)
|
||||
assert calls == [("gpu", 1), ("cpu",)]
|
||||
|
||||
|
||||
class TestBatchTimeout:
|
||||
def test_env_var_parsing(self, monkeypatch):
|
||||
"""LIDAR_BATCH_TIMEOUT: unset/0/invalid = unlimited, N seconds = N."""
|
||||
from lidar_pipeline.pipeline import _batch_timeout_s
|
||||
monkeypatch.delenv("LIDAR_BATCH_TIMEOUT", raising=False)
|
||||
assert _batch_timeout_s() == 0.0
|
||||
monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "")
|
||||
assert _batch_timeout_s() == 0.0
|
||||
monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "0")
|
||||
assert _batch_timeout_s() == 0.0
|
||||
monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "3600")
|
||||
assert _batch_timeout_s() == 3600.0
|
||||
monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "-5")
|
||||
assert _batch_timeout_s() == 0.0
|
||||
monkeypatch.setenv("LIDAR_BATCH_TIMEOUT", "abc")
|
||||
assert _batch_timeout_s() == 0.0
|
||||
|
||||
Reference in New Issue
Block a user