How sdata's DependencyGraph Tracks Data Relationships and Lineage

sdata tracks data lineage through a bipartite graph structure implemented in ProcessGraph, which maps inputs and outputs via directed edges and applies Kahn's topological sort algorithm to produce deterministic execution order.

sdata is an open-source framework for managing scientific data pipelines. Its dependency tracking system constructs a directed acyclic graph (DAG) connecting data objects to processing nodes, enabling both relationship mapping and lineage tracing across complex workflows.

Constructing the Bipartite Relationship Graph

Rather than exposing a class literally named DependencyGraph, sdata embeds dependency tracking inside the ProcessGraph class, which treats data objects and process nodes as distinct vertex types.

Mapping Data and Process Nodes

In sdata/sclass/process_graph.py, the constructor iterates over each ProcessNode's declared input_classes and output_classes to build the edge list. This creates a bipartite graph where edges flow from data nodes to process nodes (consumption) and from process nodes to data nodes (production).

for input_name, input_class in input_classes.items():
    self.data_nodes.add(input_name)
    self.class_map[input_name] = input_class
    self.edges.append((input_name, proc_name))

for output_name, output_class in output_classes.items():
    self.data_nodes.add(output_name)
    self.class_map[output_name] = output_class
    self.edges.append((proc_name, output_name))

Source: sdata/sclass/process_graph.py

Aggregating Per-Process Dependencies

Each ProcessNode exposes its own dependencies via get_dependencies(). The ProcessGraph aggregates these into a single dictionary that maps every node to its immediate prerequisites.

def get_dependencies(self):
    dependencies = {}
    for process in self.processes:
        d = process.get_dependencies()
        dependencies.update(d)
    return dependencies

Source: sdata/sclass/process_graph.py

Topological Sorting for Lineage Resolution

Once the graph edges are established, sdata determines execution order and lineage depth using a generic topological sort utility.

Kahn's Algorithm Implementation

The toposort function in sdata/sclass/dependency_graph.py implements Kahn's algorithm. It iteratively yields sets of nodes with no remaining dependencies, representing parallelizable layers of the pipeline.

def toposort(graph: Dict[T, Iterable[T]]) -> Iterable[Set[T]]:
    # Normalize to sets and drop self-references

    graph = {node: set(dep for dep in deps if dep != node) 
             for node, deps in graph.items()}
    
    # Ensure nodes appearing only as dependencies are tracked

    extra_nodes = {dep for deps in graph.values() for dep in deps} - set(graph)
    graph.update({node: set() for node in extra_nodes})

    while True:
        independent = {n for n, deps in graph.items() if not deps}
        if not independent:
            break
        yield independent
        graph = {n: (deps - independent) for n, deps in graph.items()
                 if n not in independent}
    if graph:
        raise CircularDependencyError(graph)

Source: sdata/sclass/dependency_graph.py

Generating a Flat Execution Order

For linear pipeline execution, toposort_flatten consumes the generator and optionally sorts each layer to produce a deterministic list.

def toposort_flatten(graph: Dict[T, Iterable[T]], sort: bool = True) -> List[T]:
    result = []
    for level in toposort(graph):
        result.extend(sorted(level) if sort else list(level))
    return result

Source: sdata/sclass/dependency_graph.py

Detecting Circular Dependencies

The topological sort acts as a validation mechanism. If any nodes remain with unresolved dependencies after the algorithm exhausts all independent sets, the function raises a CircularDependencyError, preventing infinite loops in pipeline definitions.

Practical Pipeline Construction

Developers interact with the dependency graph through ProcessGraph methods that wrap the sorting utilities.

from sdata.sclass.process_graph import ProcessGraph

# Build graph from ProcessNode subclasses

g = ProcessGraph(processes=[ExtractNode, TransformNode, LoadNode])

# Retrieve lineage as a flat execution order

execution_order = g.get_dag(flatten=True)

# Inspect raw relationships

print(g.edges)  # List of (data, process) and (process, data) tuples

The resulting execution_order list represents the complete data lineage from source inputs to final outputs, while g.edges serves as a lineage map for auditing or visualization via to_graphviz.

Summary

  • ProcessGraph constructs a bipartite DAG linking data nodes to process nodes via directed edges.
  • Relationship mapping is captured in the edges list and get_dependencies() aggregation.
  • Lineage resolution is performed by toposort in dependency_graph.py, which implements Kahn's algorithm to yield parallelizable execution layers.
  • Circular dependencies are detected and raised as CircularDependencyError when the graph contains cycles.
  • Deterministic ordering is provided by toposort_flatten, enabling reliable pipeline execution and full data provenance tracking.

Frequently Asked Questions

What is the difference between ProcessGraph and DependencyGraph in sdata?

sdata does not expose a class named DependencyGraph. Instead, ProcessGraph acts as the dependency graph engine, managing node relationships and execution order, while sdata/sclass/dependency_graph.py provides the generic topological sorting algorithms that power lineage resolution.

How does sdata detect circular dependencies in data relationships?

The toposort function validates the DAG by tracking remaining dependencies after each iteration. If unresolvable nodes remain after exhausting all independent layers, it raises a CircularDependencyError containing the offending cycle, preventing invalid pipeline configurations from executing.

Can sdata's dependency tracking handle parallel execution?

Yes. The toposort generator yields sets of independent nodes at each level, representing tasks that can execute concurrently. Callers can process these sets in parallel or flatten them into a serial order using toposort_flatten with deterministic sorting enabled.

Which method returns the complete lineage order for a pipeline?

Calling ProcessGraph.get_dag(flatten=True) returns the topologically sorted list representing the full data lineage from upstream inputs to downstream outputs. Alternatively, invoking toposort_flatten directly on the dependency dictionary produces the same deterministic execution sequence.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →