Structured Logging for Spatial Pipelines
This page makes a spatial pipeline’s logs queryable — emitting JSON events with structlog, binding the run identifier and the CRS once per run, putting the shard key and the counts in fields rather than in message text, sampling the per-tile chatter so a 4,096-shard job does not produce a million lines, and writing a run summary that outlives the log retention.
Why you hit this
When a district goes missing from a twin, the investigation starts in the logs, and what is usually there is a few thousand lines of Processing tile 120210233010... done interleaved from eight workers. Answering “was that shard processed, and how long did it take?” means grep and guesswork; answering “which shards did this run skip?” is impossible. Structured events fix both, and they cost nothing at write time. The wider instrumentation picture is in pipeline observability and monitoring.
Prerequisites
- Python 3.10+ with
structlog>=24.1. The standard library’sloggingworks too, with a JSON formatter, at the cost of more boilerplate. - Somewhere to send the output: a file per run is enough to start; a log system with field queries — Loki, OpenSearch, CloudWatch — is what makes it pay off.
- A stable spatial key for the unit of work: a quadkey, a tile index, a delivery tile name.
Step-by-Step
1. Configure once, JSON out
import logging
import sys
import structlog
def configure_logging(level="INFO", json_output=True):
shared = [
structlog.contextvars.merge_contextvars,
structlog.processors.add_log_level,
structlog.processors.TimeStamper(fmt="iso", utc=True),
structlog.processors.StackInfoRenderer(),
structlog.processors.format_exc_info,
]
renderer = (structlog.processors.JSONRenderer() if json_output
else structlog.dev.ConsoleRenderer(colors=True))
structlog.configure(
processors=shared + [renderer],
wrapper_class=structlog.make_filtering_bound_logger(getattr(logging, level)),
logger_factory=structlog.PrintLoggerFactory(file=sys.stdout),
cache_logger_on_first_use=True,
)
return structlog.get_logger()
log = configure_logging(json_output=not sys.stdout.isatty())
One switch decides the renderer: JSON when the output is a pipe or a file, human-readable colours when a developer is watching a terminal. That keeps the same call sites useful in both settings, which is what stops people adding print statements alongside the logger.
format_exc_info matters more than it looks. It turns an exception into structured fields rather than a multi-line traceback embedded in a JSON string, so a failure is searchable by exception type across runs.
2. Bind the run context, not every call
import os
import uuid
from datetime import datetime, timezone
def start_run(dataset, crs, delivery_date):
run_id = os.environ.get("CI_JOB_ID") or uuid.uuid4().hex[:12]
structlog.contextvars.clear_contextvars()
structlog.contextvars.bind_contextvars(
run_id=run_id,
dataset=dataset,
crs=crs,
delivery_date=delivery_date,
pipeline_version=os.environ.get("GIT_SHA", "dev"),
host=os.uname().nodename,
started=datetime.now(timezone.utc).isoformat(),
)
log.info("run.start")
return run_id
run_id = start_run("city", "EPSG:25832+7837", "2026-09-15")
Context variables attach these fields to every subsequent event in the same task or thread, so no call site has to remember them. The CRS is in there deliberately: half the confusing incidents in a spatial pipeline come down to coordinates in an unexpected system, and having it on every line means the question “what CRS was this run in?” is never open.
Using the CI job identifier when there is one links logs to the build that produced them without a separate correlation step.
3. Name events like data, not like sentences
def tile_shard(shard, features):
bound = log.bind(shard=shard, feature_count=len(features))
bound.info("shard.start")
try:
result = do_tiling(shard, features)
except Exception:
bound.exception("shard.failed") # exc_info captured as fields
raise
bound.info("shard.done", triangles=result.triangles, bytes=result.bytes,
duration_s=round(result.seconds, 2), geometric_error=result.geometric_error)
return result
An event name like shard.done is a stable key that can be counted, filtered and alerted on; Finished tiling shard 120210233010 in 4.2s is prose that cannot. The convention that pays off across a pipeline is noun.verb in the past tense, one name per meaningful outcome, and every variable part in a field.
log.bind returns a logger with extra fields, which is the right tool for a scope narrower than the run — a shard, a delivery file, a retry attempt.
4. Sample the chatter, keep the outcomes
import random
class Sampler:
"""Log every event for failures and slow work, a fraction of the routine ones."""
def __init__(self, rate=0.02, slow_s=30.0, seed=7):
self.rate, self.slow_s, self.rng = rate, slow_s, random.Random(seed)
def should_log(self, duration_s=0.0, failed=False):
return failed or duration_s >= self.slow_s or self.rng.random() < self.rate
sampler = Sampler()
for shard in shards:
res = tile_shard_quiet(shard)
if sampler.should_log(res.seconds, res.failed):
log.info("shard.done", shard=shard, duration_s=round(res.seconds, 2),
triangles=res.triangles, sampled=True)
totals.add(res)
log.info("stage.done", stage="tiling", shards=totals.count, failed=totals.failed,
triangles=totals.triangles, duration_s=round(totals.seconds, 1),
p95_shard_s=round(totals.p95, 2))
A per-shard line for 4,096 shards across twenty runs a month is a million events that nobody reads and somebody pays to store. Sampling keeps a representative few percent, always keeps failures and slow outliers, and relies on the stage summary for the totals — which is the line that actually gets queried. Seeding the sampler keeps a run reproducible, so re-running a build produces the same log volume.
5. Write a run summary that outlives the logs
import json
from pathlib import Path
def write_run_summary(path, run_id, stages, gate_results, outputs):
summary = {
"run_id": run_id,
"finished": datetime.now(timezone.utc).isoformat(),
"pipeline_version": os.environ.get("GIT_SHA", "dev"),
"dataset": "city",
"crs": "EPSG:25832+7837",
"stages": stages, # {"tiling": {"shards": 4096, "failed": 0, "seconds": 812}}
"gates": gate_results, # {"jsonld": "pass", "validator": "pass"}
"outputs": outputs, # {"tileset": "s3://…/tileset.json", "bytes": 1.2e9}
}
Path(path).write_text(json.dumps(summary, indent=2))
log.info("run.summary_written", path=str(path), stages=list(stages))
return summary
write_run_summary("build/run_summary.json", run_id,
stages={"tiling": {"shards": 4096, "failed": 0, "seconds": 812}},
gate_results={"validator": "pass", "wordcount": "n/a"},
outputs={"tileset": "s3://twin-tiles/city/v42/tileset.json"})
Log retention is typically weeks; a twin’s questions arrive after months. A summary JSON stored alongside the outputs is small, permanent and enough to answer most of them — which version of the pipeline produced this tileset, how many shards it had, whether the gates passed. It is also what makes the metrics in exporting Prometheus metrics from tiling jobs reconstructible if the metrics backend loses a day.
6. Query it
import json
from collections import Counter
def shards_missing(log_path, expected):
seen, failed = set(), Counter()
for line in Path(log_path).read_text().splitlines():
try:
rec = json.loads(line)
except json.JSONDecodeError:
continue
if rec.get("event") == "shard.done":
seen.add(rec["shard"])
elif rec.get("event") == "shard.failed":
failed[rec["shard"]] += 1
return sorted(set(expected) - seen - set(failed)), failed
missing, failed = shards_missing("build/run.log", expected_shards)
print(f"{len(missing)} shards with no outcome logged, {len(failed)} failed")
With sampling in place, shard.done is not emitted for every shard, so this query works against a run with sampling disabled — which is exactly what you turn on when investigating. Keeping a --log-every-shard flag in the pipeline, off by default, is the cheapest debugging affordance available.
Expected Output & Verification
{"event": "run.start", "level": "info", "timestamp": "2026-09-17T02:00:04Z", "run_id": "ci-88421", "dataset": "city", "crs": "EPSG:25832+7837", "delivery_date": "2026-09-15", "pipeline_version": "1f6e183", "host": "build-07"}
{"event": "shard.done", "level": "info", "timestamp": "2026-09-17T02:11:52Z", "run_id": "ci-88421", "dataset": "city", "crs": "EPSG:25832+7837", "shard": "120210233010", "feature_count": 311, "triangles": 12408, "bytes": 982144, "duration_s": 4.2, "sampled": true}
{"event": "stage.done", "level": "info", "timestamp": "2026-09-17T02:14:16Z", "run_id": "ci-88421", "stage": "tiling", "shards": 4096, "failed": 0, "triangles": 50104882, "duration_s": 812.4, "p95_shard_s": 6.8}
Verify that the logs answer the questions they exist for. Three assertions, run against a real log file in CI, are enough:
lines = [json.loads(l) for l in Path("build/run.log").read_text().splitlines() if l.startswith("{")]
events = Counter(r["event"] for r in lines)
assert events["run.start"] == 1, "exactly one run.start per run"
assert all({"run_id", "dataset", "crs"} <= set(r) for r in lines), "context missing from some events"
stage = next(r for r in lines if r["event"] == "stage.done" and r["stage"] == "tiling")
assert stage["shards"] == 4096 and stage["failed"] == 0
print("log contract satisfied:", dict(events))
Asserting the contract rather than the content is what keeps logging useful over time. A refactor that drops the context binding, or renames an event, breaks a test instead of quietly removing the field an alert depends on.
Performance Notes
- JSON rendering costs microseconds, which is irrelevant next to any spatial operation. Do not optimise logging; optimise log volume.
- Bind once, not per call. Context variables are cheap to read and the alternative — passing fields through every function — is what leads people back to prose messages.
- Write to stdout and let the platform ship it. A pipeline that writes its own log files, rotates them and uploads them is reinventing the part of the stack that already works.
- Sample per unit of work, never per stage. Losing a stage summary costs the query that matters most.
- Keep the summary JSON small — kilobytes — so it can live next to the outputs forever without anyone deciding to clean it up.
Common Errors
Every worker logs the same shard. Multiprocessing workers inherit the parent’s context variables at fork and then bind their own; if a shard key is bound before the fan-out, every child carries it. Bind inside the worker.
Context disappears in a thread pool. contextvars do not propagate into threads started before the binding. Bind inside the thread’s entry point, or pass the fields explicitly to the worker.
Logs are JSON but unsearchable. The fields are nested inside one message string because a formatter serialised the event dict into text. Check that the collector is parsing JSON rather than treating the line as a message.
Timestamps are in local time. Comparing runs across a daylight-saving change becomes guesswork. Use ISO 8601 in UTC, as the configuration above does.
Frequently Asked Questions
Is structlog necessary, or will logging do?
logging with python-json-logger and a LoggerAdapter produces the same output with more code. structlog’s context variables and bind are the parts worth having; if a project already standardises on logging, keep it and add a JSON formatter.
How long should logs be kept?
Long enough to investigate what users report, which is usually 30–90 days. The permanent record is the run summary, not the log stream.
Should logs carry geometry?
No. A bounding box as four numbers is useful; a WKT polygon in a log line is not, and a log line per feature will bury everything else. Geometry belongs in the outputs and in the change report.
What about logging from inside PDAL or GDAL?
Their native logs go to stderr and are not structured. Capture them per stage, attach the tail to a stage.failed event when something goes wrong, and do not attempt to parse them routinely.
Related Guides
- Pipeline Observability and Monitoring — how logs, metrics and traces divide the work
- Exporting Prometheus Metrics from Tiling Jobs — the numeric counterpart to the stage summary
- Tracing Pipeline Stages with OpenTelemetry — linking events across processes