Tracking Tile Progress and Failures in Dask
TL;DR: Give every task a readable key (tile-<name>), consume futures with as_completed, and append one JSON line per finished tile — success or failure with the exception type — to a manifest file on disk. On rerun, read the manifest, skip tiles already done (or whose output already exists), and submit only the rest. Group failures by exception type to tell bad tiles from bad infrastructure.
# Context and Motivation
This guide is part of Dask Distributed Processing for LiDAR. A batch of twenty thousand tiles will not finish cleanly on the first attempt. Some tiles are corrupt, some exceed memory, some hit a network timeout, and sometimes the client itself dies — a laptop sleeps, a CI job times out. Without a durable record, the only options are to rerun everything or to guess. The Dask dashboard shows what is happening now; a manifest shows what happened, survives the process that wrote it, and turns “rerun the failures” into a one-line filter.
# Prerequisites and Assumptions
- A Dask client and cluster as in processing LiDAR tiles with Dask distributed.
- A per-tile function with deterministic output names, so a completed tile can also be recognized by its output.
- Local or shared disk for the manifest.
# Step-by-Step Implementation
# Step 1 — Name every task
client.map(..., key=[f"tile-{stem}" for stem in stems]) makes keys readable and maps them back to tiles trivially.
# Step 2 — Append a line per completed tile
Inside the as_completed loop, write one JSON object per line with tile, status, duration and error type; flush after each write.
# Step 3 — Reconcile on start
Read existing manifest lines; the latest line per tile wins. Optionally also check whether the output object exists.
# Step 4 — Submit only what is needed
Tiles without a done line are submitted; others are skipped.
# Step 5 — Summarize failures
Count failures by exception type and list the tiles for each, to decide between fixing data, resizing workers or simply retrying.
# Complete Working Example
"""Resumable Dask batch with a JSON-lines manifest and failure summary."""
from __future__ import annotations
import collections
import json
import sys
import time
from pathlib import Path
from dask.distributed import Client, as_completed
from lidar_tasks import process_tile
MANIFEST = Path("manifest.jsonl")
def load_manifest() -> dict[str, dict]:
state: dict[str, dict] = {}
if MANIFEST.exists():
for line in MANIFEST.read_text().splitlines():
rec = json.loads(line)
state[rec["tile"]] = rec # latest record wins
return state
def record(fh, tile: str, status: str, **extra) -> None:
fh.write(json.dumps({"tile": tile, "status": status, "ts": time.time(), **extra}) + "\n")
fh.flush()
def run(address: str, keys: list[str]) -> None:
state = load_manifest()
todo = [k for k in keys if state.get(Path(k).stem, {}).get("status") != "done"]
print(f"{len(keys) - len(todo)} already done, submitting {len(todo)}")
if not todo:
return
client = Client(address)
futures = client.map(process_tile, todo, retries=1, pure=False,
key=[f"tile-{Path(k).stem}" for k in todo])
with MANIFEST.open("a") as fh:
for i, fut in enumerate(as_completed(futures), 1):
tile = fut.key.removeprefix("tile-")
if fut.status == "error":
exc = fut.exception()
record(fh, tile, "failed", error=type(exc).__name__, message=str(exc)[:300])
else:
r = fut.result()
record(fh, tile, "done", seconds=r.get("seconds"))
if i % 100 == 0:
print(f"{i}/{len(futures)}", flush=True)
client.close()
failures = [r for r in load_manifest().values() if r["status"] == "failed"]
by_type = collections.Counter(r["error"] for r in failures)
print("failures by type:", dict(by_type))
for err, n in by_type.most_common():
sample = [r["tile"] for r in failures if r["error"] == err][:5]
print(f" {err}: {n} tiles, e.g. {sample}")
if __name__ == "__main__":
tile_keys = [k.strip() for k in Path("tile_list.txt").read_text().splitlines() if k.strip()]
run(sys.argv[1], tile_keys)Running the script again after a crash or after fixing a problem submits only tiles that are not yet done.
# Reporting to People, Not Just Scripts
The manifest is primarily for machines, but a batch that runs for hours has human stakeholders too. A few derived views make it useful to them. A running count of done, failed and pending tiles, written to a small status file or posted to a chat channel every few hundred tiles, answers “how far along is it?” without anyone opening the dashboard. The median and 95th-percentile task time from the seconds field, recomputed as the batch progresses, gives an honest estimate of the remaining time. And at the end, a one-page summary — counts, total compute time, failures by type with example tiles — is the artefact that goes into the processing report for the delivery.
Because the manifest is plain JSON lines, all of this is a few lines of pandas: pd.read_json("manifest.jsonl", lines=True), keep the last record per tile, and group. Keeping these views derived from the manifest, rather than tracked separately, guarantees they agree with what actually happened.
# Key Parameter Table
| Element | Choice | Why |
|---|---|---|
| task key | tile-<stem> |
Readable, reversible to the tile |
| manifest format | JSON lines, append-only | Crash-safe, easy to parse |
| latest-wins | per tile | Reruns overwrite earlier failures |
retries |
1 | Transient errors retried before recording failure |
| failure grouping | by exception type | Distinguishes data from infrastructure |
| output check | optional existence test | Recovers state if the manifest is lost |
# Verification
- Kill and resume. Interrupt the client mid-batch, rerun, and confirm it submits only unfinished tiles.
- Counts. After completion,
doneplusfailedequals the tile list. - Output agreement. Every
donetile has an output object; a quick listing of the output prefix confirms it.
# Gotchas and Edge Cases
Tasks still running when the client dies. When the client disconnects, the scheduler cancels its futures by default, so tiles in progress stop and have no manifest line; they are resubmitted on rerun. Outputs they partly wrote are overwritten thanks to deterministic names.
Duplicate submissions. Two clients running the same batch against the same outputs duplicate work and race on the manifest. Use a lock file or run batches from one orchestrator.
Large messages. Recording full tracebacks for thousands of failures bloats the manifest. Store the exception type and a short message; keep full tracebacks in worker logs.
Retrying persistent failures. A tile that fails the same way twice will fail a third time. Stop retrying after one or two reruns and route those tiles to investigation.
# Frequently Asked Questions
How do I resume a Dask batch after the client crashes?
Record each finished tile in a durable manifest as results arrive. On restart, read the manifest and submit only tiles not marked done; deterministic output names make reruns of interrupted tiles safe.
Why use JSON lines for the manifest?
Each record is written and flushed as one line, so a crash leaves at most one partial line and everything before it is valid. It is easy to append to, easy to read back and easy to inspect with standard tools.
How do I find out why tiles failed?
Store the exception type and a short message per failure, then group by type. A few types usually account for all failures, and each points to a specific fix: memory, corrupt input, storage timeouts or metadata problems.
Should failed tiles be retried automatically?
Once or twice, for transient problems. Tiles that fail repeatedly with the same error need investigation, not more retries.
# Related
- Dask Distributed Processing for LiDAR — the overall pattern
- Processing LiDAR Tiles with Dask Distributed — the cluster this runs on
- Making Tile Outputs Idempotent — why reruns are safe
- Retrying Failed Tiles in Airflow — the same idea in Airflow
- Diagnosing PDAL Out-of-Memory Failures — the commonest failure type