Execution¶
This page explains what happens when you call
compute().
Execution pipeline¶
Target resolution: identify which nodes are needed to produce the requested targets.
Topological sort: order the reachable nodes so that every node executes after its dependencies. Uses Kahn’s algorithm; cycles are detected and rejected.
Cache lookup: for each node, derive its result key from the node, resolved inputs, tool identity, environment identity, and selected upstream records. If
current.jsonselects a reusable record, load it and skip execution.Input resolution: for each row, resolve column bindings and constants into concrete values. Output path templates are resolved at this stage.
Tool execution:
DataFrameTool: call
merge_dataframes(if multiple upstreams), thentransform. Runs in the main process.ProcessingTool: call
process_rowfor each row (orprocess_batchfor all rows). Runs in the tool’s declared environment.
Result assembly: collect outputs into a DataFrame. For ProcessingTools, each
Outputsinstance becomes a row.Cache publication: publish the output DataFrame and owned assets as an immutable record, then select it through
current.json.
Published cache records are retained until an explicit storage maintenance operation prunes them.
Index alignment¶
When a ProcessingTool receives inputs from multiple upstream nodes, the engine aligns rows by index. Row 0 of node A matches row 0 of node B.
For explosion tools (one-to-many), child indices use :: separators:
Parent index: "0"
Child indices: "0::0", "0::1", "0::2"
Downstream tools receiving both parent and child data use the parent prefix to align rows correctly.
Engines and scheduling¶
The default engine is engine="wetlands" with execution="parallel".
It runs ProcessingTool methods in isolated Wetlands
worker environments. In that backend, max_workers and
per-environment overrides control worker-process pools and row dispatch.
See Parallelism for max_workers,
ResourceSpec, and per-environment overrides.
Use execution="sequential" for single-node-at-a-time debugging.
Storage layout¶
{storage_path}/
├── cache/
│ └── v1/
│ └── results/
│ └── {shard}/{shard}/{result_key}/
│ ├── current.json
│ ├── attempts/{attempt_id}/staging/
│ └── records/{record_id}/
│ ├── manifest.json
│ ├── dataframe.parquet
│ └── assets/
├── views/
│ ├── runs/
│ └── latest/
└── outputs/
├── runs/
└── latest/
cache/v1 is the canonical machine-readable cache root. views/runs/
and views/latest/ are portable JSON provenance views over selected
records and are not used to decide cache hits. outputs/ contains optional
materialized files for human browsing when an output view is enabled.
Result keys¶
The result key is the public cache identity exposed by planning, progress, invalidation, and run-view APIs. It includes the node’s logical inputs and selected upstream record references when upstream records are available. The full material is documented in Output and Cache Storage Specification.
Progress events¶
The engine emits ProgressEvent objects via the
on_progress callback:
started: node execution beginsrow_complete: a single row finished (withrowandtotal_rows)completed: node execution finishedcached: node result loaded from cache (no execution)
Cache-related events expose result_key / record_id values.
Diagnostic signatures are separate debug values and are not cache keys.
See also¶
Caching and Provenance — the result-key/current-record model,
plan(),invalidate(), andNodePlanStatusvalues.