How to Use the Declarative Pipeline DSL in Semantica for Data Processing

The declarative pipeline DSL in Semantica provides a fluent, method-chaining interface through the PipelineBuilder class that lets you define data-processing workflows as named steps with explicit dependencies, validation, and parallel execution capabilities.

The Semantica AGI repository offers a powerful declarative pipeline DSL for constructing complex data-processing workflows without imperative boilerplate. This fluent interface allows developers to define pipeline topology, configure step-specific parameters, and orchestrate parallel execution through immutable definitions that serialize to JSON for reproducibility.

Core Components of the Declarative Pipeline DSL

The architecture centers on immutable data classes and specialized orchestrators defined across the semantica/pipeline/ module.

PipelineBuilder – The DSL Entry Point

Located in semantica/pipeline/pipeline_builder.py, the PipelineBuilder class serves as the primary interface for constructing pipelines. It collects PipelineStep definitions, supports method chaining for fluency, and exposes serialize() and deserialize() methods for JSON persistence. When you invoke build(), the builder automatically delegates structural validation to PipelineValidator before returning an immutable Pipeline instance.

PipelineStep and Immutable Pipeline Definitions

Each step in your workflow is represented by a PipelineStep dataclass that stores the step name, type string, configuration dictionary, dependency list, optional handler reference, and runtime metadata. The Pipeline dataclass acts as an immutable container comprising a list of these steps plus optional metadata, ensuring that once defined, a pipeline’s structure cannot be accidentally mutated before execution.

Validation and Execution Infrastructure

The PipelineValidator class in semantica/pipeline/pipeline_validator.py enforces structural integrity by checking for unique step names, missing dependencies, and circular references. At runtime, the ExecutionEngine (semantica/pipeline/execution_engine.py) orchestrates step execution, respecting declared dependencies while leveraging the ParallelismManager to identify steps that can run concurrently. Optional resource allocation is handled by the ResourceScheduler, and failures are managed through the centralized FailureHandler.

Constructing Data Pipelines with PipelineBuilder

You define workflows by instantiating PipelineBuilder and chaining method calls to register steps, configure parallelism, and finalize the definition.

from semantica.pipeline import PipelineBuilder, ExecutionEngine

# 1. Instantiate the builder

builder = PipelineBuilder()

# 2. Declare steps with dependencies

builder.add_step("ingest", "file_ingest", path="data/input.csv")
builder.add_step("parse", "document_parse", dependencies=["ingest"])
builder.add_step("embed", "semantic_embedding", dependencies=["parse"], parallel_safe=True)

# 3. Configure global parallelism

builder.set_parallelism(level=2)

# 4. Build the immutable pipeline (triggers validation automatically)

pipeline = builder.build(name="my_data_pipeline")

# 5. (Optional) Serialize for storage or remote dispatch

pipeline_json = builder.serialize(pipeline)

# 6. Execute

engine = ExecutionEngine()
result = engine.execute_pipeline(pipeline, data=None)
print(result.success, result.output)

The add_step() method accepts a name, type identifier, and arbitrary keyword arguments for configuration. If you add steps out of order, use connect_steps(from_step, to_step) to wire explicit dependencies manually. The set_parallelism(level) method hints to the execution engine how many workers to allocate for concurrent step processing.

Advanced DSL Patterns

Beyond basic step chaining, the DSL supports custom logic injection, concurrent execution hints, and versioned data processing.

Registering Custom Step Handlers

For steps requiring arbitrary Python logic, register a callable handler and reference it by type:

def normalize_text(step_config, previous_output):
    # previous_output contains results from dependency steps

    return {"normalized": previous_output["text"].lower().strip()}

builder = PipelineBuilder()

# Register handler under a custom type key

builder.step_registry["text_normalize"] = normalize_text

builder.add_step("load", "file_ingest", path="article.txt")
builder.add_step(
    "clean",
    "text_normalize",
    dependencies=["load"],
    handler=normalize_text
)

pipeline = builder.build()

Enabling Parallel Execution

Mark independent steps with parallel_safe=True to allow the ParallelismManager to schedule them concurrently:

builder.add_step(
    "embed",
    "semantic_embedding",
    dependencies=["parse"],
    parallel_safe=True,  # Permits concurrent execution

    batch_size=100
)

The execution engine enforces that parallel-safe steps return dict-like results to ensure mergeable outputs.

Delta Mode for Versioned Datasets

When processing versioned data, enable delta mode to compute changes between versions rather than full recomputation:

builder.add_step(
    "update_index",
    "vector_index_update",
    delta_mode=True,
    base_version_id="v1.2",
    target_version_id="v1.3"
)

Validating and Executing Pipelines

Calling build() implicitly invokes PipelineValidator.validate(), which raises ValidationError for topological issues like circular dependencies or missing step references. This prevents invalid pipelines from reaching the execution phase.

Once validated, pass the Pipeline instance to an ExecutionEngine via execute_pipeline(pipeline, data). The engine consults the ParallelismManager to determine the execution graph, allocates resources through the ResourceScheduler, and tracks progress through the global ProgressTracker. The method returns an ExecutionResult object containing success status, output data, execution metrics, and error details if failures occurred.

Summary

  • The PipelineBuilder class in semantica/pipeline/pipeline_builder.py provides the fluent DSL entry point for declaring data-processing workflows.
  • Steps are defined as PipelineStep dataclasses within an immutable Pipeline object, ensuring structural integrity.
  • PipelineValidator automatically checks for circular dependencies and missing references when build() is called.
  • The ExecutionEngine orchestrates runtime behavior, respecting dependency chains while leveraging the ParallelismManager for concurrent step execution.
  • Advanced features include custom handler registration, parallel_safe=True hints for concurrency, and delta_mode for efficient versioned data processing.

Frequently Asked Questions

What is the declarative pipeline DSL in Semantica?

The declarative pipeline DSL in Semantica is a fluent, method-chaining interface exposed through the PipelineBuilder class that allows developers to define data-processing workflows as a series of named, typed steps with explicit dependencies. It abstracts away execution details, letting you focus on what data transformations occur rather than how they are orchestrated.

How does the PipelineBuilder validate pipeline topology?

When you call build(), the PipelineBuilder automatically invokes PipelineValidator.validate() to check for structural errors such as duplicate step names, missing dependencies, or circular reference chains. Any validation failures raise a ValidationError immediately, preventing invalid pipeline definitions from proceeding to execution.

Can I execute pipeline steps in parallel using the DSL?

Yes. Set parallel_safe=True when calling add_step() to mark steps that have no side effects or resource conflicts with other concurrent steps. The ExecutionEngine uses the ParallelismManager to identify which parallel-safe steps share no overlapping dependencies and schedules them concurrently up to the limit specified by set_parallelism(level).

How do I persist a pipeline definition for later execution?

Call builder.serialize(pipeline) to convert the immutable Pipeline object into a JSON string suitable for storage, version control, or remote dispatch. You can later reconstruct the pipeline using builder.deserialize(json_str) without redefining the step configuration, ensuring reproducible execution across different environments or scheduled runs.

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 →