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

> Master Semantica's declarative pipeline DSL to build efficient data processing workflows. Define steps, manage dependencies, and enable parallel execution with ease.

- Repository: [Semantica /semantica](https://github.com/semantica-agi/semantica)
- Tags: how-to-guide
- Published: 2026-09-12

---

**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`](https://github.com/semantica-agi/semantica/blob/main/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`](https://github.com/semantica-agi/semantica/blob/main/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`](https://github.com/semantica-agi/semantica/blob/main/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.

```python
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:

```python
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:

```python
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:

```python
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`](https://github.com/semantica-agi/semantica/blob/main/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.