How to Use the sdata ProcessGraph Module for Complex Workflow Management
The sdata ProcessGraph module enables you to define, validate, and visualize complex data-processing pipelines as directed acyclic graphs using typed data containers, declarative processing nodes, and automatic dependency resolution.
The lepy/sdata repository provides a lightweight framework for structuring scientific data workflows. The sdata ProcessGraph module sits at the core of this system, offering a Pythonic API to chain processing steps, enforce type safety through metadata, and render execution graphs without external orchestration tools.
Core Architecture of the ProcessGraph Module
The module implements a three-layer architecture that separates data definition, processing logic, and workflow orchestration.
ProcessData: Typed Containers with Metadata
ProcessData is the foundational class for all inputs and outputs in the system. Defined in sdata/sclass/process.py (line 6), it enforces a strict schema via a metadata dictionary that tracks required attributes, data types, units, and descriptions. Subclasses override the constructor to pre-populate attribute specifications, ensuring that every data object carries self-describing metadata compliant with the workflow’s contract.
ProcessNode and CompositeProcess: Declarative Processing Steps
Individual computations are encapsulated in ProcessNode (line 33) or grouped hierarchically in CompositeProcess (line 77), both located in sdata/sclass/process.py. Each node explicitly declares required_inputs and required_outputs as dictionaries mapping port names to ProcessData subclasses. The base implementation provides validate_inputs() to check type compatibility at runtime, while concrete subclasses implement run() or execute() to perform the actual computation. The helper create_process_class() (line 32) automates the boilerplate of subclassing ProcessNode, filling in the input/output class declarations and generating the validation logic.
ProcessGraph: DAG Construction and Validation
The ProcessGraph class in sdata/sclass/process_graph.py (line 26) aggregates multiple nodes into a directed acyclic graph (DAG). During construction, it analyzes each node’s required_inputs and required_outputs to build an internal edge list representing data dependencies. The static method toposort() (line 68) performs a topological sort using a set-based algorithm; if any nodes remain unprocessed after the sort, it raises CircularDependencyError, preventing invalid cyclic workflows. For visualization, to_graphviz() (line 110) generates a Graphviz Digraph object that can be rendered to PNG or SVG.
Building a Workflow with the ProcessGraph Module
Implementing a pipeline follows a four-step pattern: define data types, create node classes, assemble the graph, and optionally execute manually.
1. Define Typed Data Classes
Subclass ProcessData to create strongly-typed inputs and outputs. The metadata container (self.md) is inherited from sdata/base.py.
from sdata.sclass.process import ProcessData
class Temperature(ProcessData):
"""Temperature in °C."""
def __init__(self, **kwargs):
attrs = [{"name": "value", "dtype": float, "unit": "°C",
"description": "Measured temperature", "required": True}]
kwargs.setdefault('attributes', []).extend(attrs)
super().__init__(**kwargs)
class Pressure(ProcessData):
"""Pressure in kPa."""
def __init__(self, **kwargs):
attrs = [{"name": "value", "dtype": float, "unit": "kPa",
"description": "Measured pressure", "required": True}]
kwargs.setdefault('attributes', []).extend(attrs)
super().__init__(**kwargs)
2. Create Processing Nodes
Use create_process_class() to generate concrete ProcessNode subclasses, then attach a custom run method that implements the business logic.
from sdata.sclass.process import create_process_class
# Node A: converts Celsius to Kelvin
TempToK = create_process_class(
process_name='TempToK',
input_classes={'temp_c': Temperature},
output_classes={'temp_k': Temperature},
)
def run_temp_to_k(self, inputs):
temp_c = inputs['temp_c'].md.get('value').value
out = Temperature(name='temp_k')
out.md.add('value', temp_c + 273.15)
return {'temp_k': out}
setattr(TempToK, 'run', run_temp_to_k)
# Node B: pressure-adjusted temperature calculation
PressAdj = create_process_class(
process_name='PressAdj',
input_classes={'temp_k': Temperature, 'press': Pressure},
output_classes={'adj_temp': Temperature},
)
def run_press_adj(self, inputs):
t_k = inputs['temp_k'].md.get('value').value
p = inputs['press'].md.get('value').value
out = Temperature(name='adj_temp')
out.md.add('value', t_k * (1 + p / 1000.0))
return {'adj_temp': out}
setattr(PressAdj, 'run', run_press_adj)
3. Assemble and Visualize the Graph
Pass a list of node classes to get_process_graph() in sdata/sclass/process_graph.py (line 215). This instantiates ProcessGraph, registers each node via add_process() (lines 70–94), and returns a Graphviz object.
from sdata.sclass.process_graph import get_process_graph
# Build the DAG
dot = get_process_graph([TempToK, PressAdj])
# Render to PNG (requires Graphviz installed)
dot.render('workflow', format='png', cleanup=True)
4. Execute the Pipeline Manually
While ProcessGraph handles structure, execution is typically invoked manually on instantiated nodes.
# Prepare concrete data
temp_c = Temperature(name='temp_c')
temp_c.md.add('value', 25.0)
press = Pressure(name='press')
press.md.md.add('value', 101.3)
# Run Node A
node_a = TempToK(inputs={'temp_c': temp_c})
out_a = node_a.run(node_a.inputs)
# Run Node B with upstream output
node_b = PressAdj(inputs={'temp_k': out_a['temp_k'], 'press': press})
out_b = node_b.run(node_b.inputs)
print('Adjusted temperature (K):', out_b['adj_temp'].md.get('value').value)
Validating Dependencies and Visualizing Pipelines
Before execution, validate the graph topology using ProcessGraph.toposort(). This method converts the dependency map into a set representation, iteratively extracts nodes with no unmet dependencies, and raises CircularDependencyError if cycles exist.
For visualization, to_graphviz() generates a standard Graphviz digraph, while to_graphviz_cluster() (also in process_graph.py) supports subgraph clustering for CompositeProcess hierarchies. Both methods respect the port names defined in input_classes and output_classes, labeling edges to show data flow between specific attributes.
Summary
- ProcessData in
sdata/sclass/process.pyprovides the metadata-rich container system required for type-safe workflow inputs and outputs. - ProcessNode and CompositeProcess define discrete processing steps with declarative input/output contracts and built-in validation.
- ProcessGraph constructs a DAG from node declarations, detects circular dependencies via
toposort(), and exports to Graphviz for visualization. - The
get_process_graph()convenience wrapper automates graph instantiation and registration, returning a renderable Graphviz object immediately.
Frequently Asked Questions
How does the sdata ProcessGraph module detect circular dependencies?
The ProcessGraph.toposort() static method (line 68 in sdata/sclass/process_graph.py) implements a set-based topological sort. It repeatedly extracts nodes whose dependencies are already satisfied; if any nodes remain after this process completes, the method raises CircularDependencyError, indicating an unresolvable cycle in the workflow graph.
What is the difference between ProcessNode and CompositeProcess?
ProcessNode represents a single atomic processing step with explicit required_inputs and required_outputs. CompositeProcess (line 77 in sdata/sclass/process.py) inherits from ProcessNode but acts as a container for nested sub-graphs, allowing hierarchical workflow construction where a single node logically encapsulates multiple child processes.
Can I execute a ProcessGraph automatically or only manually?
The ProcessGraph class itself focuses on structure, validation, and visualization rather than execution orchestration. While it provides the dependency order via toposort(), running the pipeline requires manual instantiation of each ProcessNode and invocation of its run() method, as shown in the execution example above.
How do I customize the visual appearance of the workflow graph?
The to_graphviz() and to_graphviz_cluster() methods return a standard Graphviz Digraph object. You can modify this object directly—adding colors, shapes, or cluster labels—before calling render(). The source code in sdata/sclass/process_graph.py (line 110) shows that these methods accept optional formatting parameters that propagate to the generated DOT syntax.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →