Processing LiDAR Tiles with Dask Distributed
TL;DR: Build one image with PDAL, Python and dask[distributed]; run dask scheduler on one host and dask worker tcp://<scheduler>:8786 --nworkers 8 --nthreads 1 --memory-limit 6GiB on each worker host; in the client, Client("tcp://<scheduler>:8786") then client.map(process_tile, tile_keys, retries=2, pure=False) and gather small results with as_completed. Tiles stream from S3 via /vsis3/, outputs go straight back to S3.
# Context and Motivation
This guide is part of Dask Distributed Processing for LiDAR. The topic explains the model; this page is the concrete first deployment — two or three cloud VMs or on-premises servers, the same container everywhere, and a client script run from a laptop or a CI job. It deliberately avoids Kubernetes so every piece is visible: a scheduler process, worker processes, and the network between them.
Once this works, moving to an operator-managed cluster changes how processes are started, not the code.
# Prerequisites and Assumptions
- Two or more Linux hosts with Docker, able to reach each other on TCP ports 8786 (scheduler) and 8787 (dashboard), plus the worker-to-worker ports Dask assigns.
- An S3 bucket readable and writable through instance roles (no credential files in images).
- A per-tile PDAL function that already works locally.
# Step-by-Step Implementation
# Step 1 — Build one image
Start from a conda-forge PDAL environment, add dask[distributed], and copy in the module containing process_tile. Client, scheduler and workers all run this image.
# Step 2 — Start the scheduler
docker run --network host lidar-dask:1.4 dask scheduler on the scheduler host.
# Step 3 — Start workers
On each worker host: docker run --network host lidar-dask:1.4 dask worker tcp://10.0.1.10:8786 --nworkers 8 --nthreads 1 --memory-limit 6GiB.
# Step 4 — Submit from the client
Connect with Client, map process_tile over the tile keys, and consume as_completed.
# Step 5 — Watch and collect
Open the dashboard, then write the manifest of results and failures when the loop ends.
# Complete Working Example
Dockerfile:
FROM mambaorg/micromamba:1.5-jammy
COPY env.lock /tmp/env.lock
RUN micromamba install -y -n base -f /tmp/env.lock && micromamba clean -a -y
COPY lidar_tasks.py /app/lidar_tasks.py
ENV PYTHONPATH=/app OMP_NUM_THREADS=1 GDAL_DISABLE_READDIR_ON_OPEN=EMPTY_DIRlidar_tasks.py — imported by workers, so the function is importable by name:
"""Per-tile task executed on Dask workers."""
from __future__ import annotations
import json
import time
from pathlib import PurePosixPath
import pdal
OUT = "/vsis3/lidar-out/dtm"
def process_tile(key: str) -> dict:
stem = PurePosixPath(key).stem
spec = {"pipeline": [
f"/vsis3/lidar-in/{key}",
{"type": "filters.range", "limits": "Classification![7:7],Classification![18:18]"},
{"type": "filters.smrf", "slope": 0.15, "window": 18, "threshold": 0.5},
{"type": "filters.range", "limits": "Classification[2:2]"},
{"type": "writers.gdal", "filename": f"{OUT}/{stem}.tif", "resolution": 1.0,
"output_type": "idw", "window_size": 6, "data_type": "float32",
"gdalopts": "COMPRESS=DEFLATE,TILED=YES"},
]}
t0 = time.perf_counter()
n = pdal.Pipeline(json.dumps(spec)).execute()
return {"tile": stem, "ground_points": n, "seconds": round(time.perf_counter() - t0, 1)}submit.py — run anywhere that can reach the scheduler:
"""Submit a tile list to a remote Dask scheduler and write a manifest."""
import json
import sys
from pathlib import Path
from dask.distributed import Client, as_completed
from lidar_tasks import process_tile
client = Client(sys.argv[1]) # e.g. tcp://10.0.1.10:8786
print("dashboard:", client.dashboard_link)
keys = [k.strip() for k in Path("tile_list.txt").read_text().splitlines() if k.strip()]
futures = client.map(process_tile, keys, retries=2, pure=False,
key=[f"tile-{Path(k).stem}" for k in keys])
done, failed = [], []
for i, f in enumerate(as_completed(futures), 1):
(failed.append({"key": f.key, "error": repr(f.exception())}) if f.status == "error"
else done.append(f.result()))
if i % 50 == 0:
print(f"{i}/{len(futures)} done, {len(failed)} failed", flush=True)
Path("manifest.json").write_text(json.dumps({"done": done, "failed": failed}, indent=2))
client.close()# Sizing the First Cluster
The first real run is also the best moment to measure what the batch needs. Start with two worker hosts, run a hundred representative tiles, and read three numbers from the performance report: median task time, peak worker memory, and the share of time workers spent idle. Median task time times the number of tiles, divided by the number of workers, gives a first estimate of wall time for the full list. Peak memory tells you whether --memory-limit can be lowered to fit more workers per host, or must be raised because the largest tiles came close to the pause threshold. Idle time above a few percent usually points to I/O — workers waiting on S3 — which argues for hosts in the bucket’s region or larger instance network bandwidth rather than more workers.
With those numbers, choose the host count for the real run from the deadline and the budget. Cloud VMs are billed per second or minute, so forty hosts for two hours costs roughly the same as ten hosts for eight; the only reason to run fewer is a quota or a shared bucket’s request limits.
When the batch ends, stop the workers and scheduler explicitly. Idle Dask workers on on-demand instances cost exactly as much as busy ones, and a forgotten fleet over a weekend is the most expensive mistake in cloud LiDAR processing.
# Key Parameter Table
| Setting | Where | Value | Why |
|---|---|---|---|
--nworkers |
dask worker |
cores or memory ÷ peak | Processes per host |
--nthreads |
dask worker |
1 | PDAL is CPU-bound, single-threaded |
--memory-limit |
dask worker |
e.g. 6GiB | Pause and restart thresholds |
retries |
client.map |
2 | Absorb transient failures |
key |
client.map |
tile-<name> |
Readable dashboard and errors |
OMP_NUM_THREADS |
image env | 1 | No oversubscription |
# Verification
- One-tile smoke test. Submit a single task first and check its output appears in S3.
- Worker count.
len(client.scheduler_info()["workers"])equals hosts ×--nworkers. - Manifest complete. Every key appears once in
doneorfailed; retry the failures as a new, small batch.
# Gotchas and Edge Cases
Functions defined in __main__. Workers must be able to import the task function. Defining it in the submitting script works only because Dask pickles it by value; putting it in a module inside the image is more robust and avoids version mismatches.
Firewalls. Workers connect to each other as well as to the scheduler. Open the worker port range or run all hosts in one security group.
Mismatched images. A worker started from an older image may accept tasks and fail them. Tag images immutably and log pdal.__version__ in task results.
Dashboard exposure. Port 8787 shows task names and logs. Keep it on a private network or behind an authenticating proxy.
# Frequently Asked Questions
How do I start a Dask cluster across several machines?
Run dask scheduler on one host and dask worker with the scheduler’s address on each other host, all from the same container image. Clients then connect to the scheduler address.
Where should the task function live?
In a module installed in the image that workers use, so workers can import it by name. This avoids pickling surprises and guarantees workers run the same code as the client expects.
How do workers read tiles from S3?
Through GDAL’s virtual file system: PDAL readers accept /vsis3/bucket/key paths and use the instance role or service account credentials available on each worker.
What should the task return?
A small summary such as point counts, timings and output paths. Outputs themselves should be written to storage from the worker, never returned through Dask.
# Related
- Dask Distributed Processing for LiDAR — the model and tuning
- Dask vs ProcessPoolExecutor for PDAL — whether you need Dask yet
- Tracking Tile Progress and Failures in Dask — manifests and reruns
- Building a Slim PDAL Docker Image — the image all nodes share
- Streaming LAZ from S3 with PDAL — the I/O pattern workers use