Skip to main content

Crate swale

Crate swale 

Source
Expand description

A scheduled, dependency-ordered orchestrator of asset graphs.

swale runs on the taquba durable execution crates. Every task instance runs as a single workflow run, and the record of its outcome and the event for the scheduler commit in one transaction.

§Definition file

One TOML file per graph. An asset node declares the asset it produces and the assets it consumes, and a task node declares its upstream nodes with after.

[graph]
name = "orders_daily"
schedule = "0 2 * * *"
partition = "daily"

[[node]]
name = "extract"
produces = "orders_raw"
operator = "subprocess"
[node.params]
argv = ["python", "tasks/extract.py"]

[[node]]
name = "load"
produces = "orders_warehouse"
consumes = ["orders_raw"]
operator = "subprocess"
[node.params]
argv = ["python", "tasks/load.py", "{{ upstream.extract.rows_key }}"]

The definition module documents every field, and the loader reports every fault of a definition at once. From Rust, load_path returns the checked graph:

use std::path::Path;
use swale::OperatorSet;

let graph = swale::load_path(
    Path::new("examples/orders_daily.toml"),
    &OperatorSet::builtin(),
)?;
println!("{}: {} nodes", graph.name(), graph.nodes().len());

§Command

The swale binary is enabled by the default cli feature. A program that embeds the library can set default-features = false.

swale validate loads a definition and prints each fault, one per line.

swale run runs a graph for one partition on a store, prints each node’s record as it is written and waits for the run to settle. The store is a directory or an object store URL (s3://bucket/prefix, gs://bucket/prefix, az://container/prefix). A cloud scheme needs the matching cargo feature (aws, gcp or azure) and reads the provider’s environment variables for its credentials. A second run for the same partition resumes the existing graph run, and a process interrupted mid-run resumes without repeating a completed node.

$ swale validate examples/orders_daily.toml
orders_daily: 5 nodes, 3 edges
$ swale run examples/local.toml --store ./swale-store
local/none: started, 1 root node(s) submitted
  first: succeeded (local-none-first-r0)
  second: succeeded (local-none-second-r0)
local/none: complete

The exit status is 0 for a valid definition or a complete run, 1 for a fault, a failed run or an error, and 2 for a usage error.

Re-exports§

pub use definition::load_path;
pub use definition::load_str;
pub use graph::Graph;
pub use graph::GraphSpec;
pub use graph::Node;
pub use graph::NodeKind;
pub use graph::NodeSpec;
pub use graph::Partitioning;
pub use graph::Problem;
pub use graph::TriggerRule;
pub use hook::EVENTS_QUEUE;
pub use hook::Event;
pub use hook::RecordHook;
pub use operator::Operator;
pub use operator::OperatorSet;
pub use operator::Outcome;
pub use operator::Task;
pub use partition::Partition;
pub use records::GraphRunRecord;
pub use records::GraphRunState;
pub use records::NodeRecord;
pub use records::RecordStatus;
pub use scheduler::DefinitionStore;
pub use scheduler::Pools;
pub use scheduler::Scheduler;
pub use scheduler::SchedulerOptions;
pub use scheduler::StartOutcome;
pub use task::TaskIdentity;
pub use template::Template;

Modules§

definition
The definition file: one TOML document per graph.
dispatch
The step runner of every pool: it reads the task input from the payload and the identity from the headers, renders the parameters and runs the operator.
duration
Durations in the definition file: an integer and a unit, such as 30s, 5m, 6h or 7d.
graph
The internal graph model: nodes, their edges and the checks a graph passes before it runs. Every front-end builds a Graph from a GraphSpec through Graph::build.
hook
The terminal hook of every pool: it writes the node’s record and enqueues the event for the scheduler, and both commit with the notification’s acknowledgement.
input
The input of a task instance: the operator, the node’s parameters before rendering and the records of its upstreams. The input is the run’s payload, so a definition edit that changes a node’s parameters while its run is active fails the resubmission with an input mismatch.
operator
The operators a definition can name: the check of their parameters at load time and their execution at run time.
partition
The partition key of an asset, in the charset of a run id.
records
The KV records of the orchestrator and their keys. Every key has the swale/ prefix, and every value is JSON written as an absolute value.
scheduler
The scheduler: it starts graph runs, submits a node when its upstreams satisfy its trigger rule, and settles the state of a graph run.
subprocess
The subprocess operator: runs a program with a JSON document on stdin and reads its output JSON from stdout.
task
The identity of a task instance: the graph, the partition, the node and the rerun count, which together form the run id and the swale. headers of the run.
template
The substitution language of a parameter string. A reference is written between {{ and }}, and the language is substitution only.

Enums§

Error
The failure of a definition load.