# How to Build and Execute a Custom Pipeline with PipelineBuilder and ExecutionEngine in Semantica

> Build and execute custom pipelines in Semantica using PipelineBuilder and ExecutionEngine. Define steps, manage dependencies, and leverage parallel processing for efficient runtime orchestration.

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

---

**Use `PipelineBuilder` to declaratively assemble steps and dependencies into a `Pipeline` dataclass, then pass it to `ExecutionEngine.execute_pipeline()` to orchestrate runtime execution with parallel processing, retries, and progress tracking.**

The **semantica-agi/semantica** repository provides a robust pipeline subsystem for constructing and running complex workflows. To build and execute a custom pipeline with **PipelineBuilder** and **ExecutionEngine**, you define your workflow topology declaratively and let the engine handle scheduling, dependency resolution, and concurrency.

## Building a Pipeline with PipelineBuilder

The **`PipelineBuilder`** class in [[`semantica/pipeline/pipeline_builder.py`](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/pipeline_builder.py)](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/pipeline_builder.py) (lines 20‑84) provides a fluent API for assembling workflows.

### Step 1: Create the Builder and Register Handlers

Instantiate the builder and register callable handlers that process step inputs.

```python
from semantica.pipeline import PipelineBuilder, ExecutionEngine

builder = PipelineBuilder()

def ingest(data, **cfg):
    return {"text": "Hello world"}

def parse(data, **cfg):
    return {"tokens": data["text"].split()}

builder.register_step_handler("file_ingest", ingest)
builder.register_step_handler("text_parse", parse)

```

Handlers are stored in `self.step_registry` (see `register_step_handler` at lines 45‑53) and invoked during execution.

### Step 2: Add and Connect Steps

Add steps using `add_step()` and link dependencies with `connect_steps()`. The **`PipelineStep`** dataclass (lines 55‑71) stores configuration, dependencies, and execution flags.

```python

# Add steps with configuration

builder.add_step("ingest", "file_ingest", parallel_safe=True)
builder.add_step("parse", "text_parse")

# Explicitly connect ingest -> parse

builder.connect_steps("ingest", "parse")

```

The `connect_steps()` method adds the source step to the target’s `dependencies` list (lines 77‑84). Steps marked `parallel_safe=True` enable concurrent execution when inputs are dictionaries.

### Step 3: Configure Parallelism and Build

Set parallelism levels and validate the pipeline topology.

```python
builder.set_parallelism(2)  # Validated at lines 97-101

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

```

The `build()` method runs **`PipelineValidator`** (lines 29‑37) to check for circular dependencies and configuration errors. It returns a **`Pipeline`** dataclass (lines 73‑80) containing ordered `PipelineStep` objects.

## Executing Pipelines with ExecutionEngine

The **`ExecutionEngine`** in [[`semantica/pipeline/execution_engine.py`](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/execution_engine.py)](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/execution_engine.py) (class starts at line 83) orchestrates runtime behavior.

### Step 1: Instantiate the Engine

The constructor (lines 95‑102) wires together a `FailureHandler`, `ParallelismManager`, and `ResourceScheduler`.

```python
engine = ExecutionEngine()

```

For delta processing workflows, inject additional dependencies:

```python
engine = ExecutionEngine(
    version_manager=my_version_manager,
    triplet_store=my_triplet_store,
)

```

### Step 2: Run the Pipeline

Invoke `execute_pipeline()` (lines 23‑38) to start execution.

```python
result = engine.execute_pipeline(pipeline, data=my_input)

```

The method initializes progress tracking (lines 41‑47), sets pipeline status to `RUNNING` (lines 88‑90), and calls `_execute_steps` to process steps. Parallel-safe layers run concurrently via `_execute_parallel_group` (lines 34‑67), while the engine deep-copies inputs for isolation (line 44).

### Step 3: Monitor Execution Status

Query runtime state using native engine methods:

```python

# Get enum: RUNNING, COMPLETED, FAILED, PAUSED

status = engine.get_pipeline_status(pipeline.name)

# Report total steps, completed count, and percentage (lines 82-99)

progress = engine.get_progress(pipeline.name)

```

Control the lifecycle with `pause_pipeline()`, `resume_pipeline()`, or `stop_pipeline()` (lines 55‑77).

## Advanced Pipeline Patterns

### Parallel-Safe Layer Execution

Steps without interdependencies and `parallel_safe=True` execute concurrently. The engine groups them via `_can_run_layer_in_parallel` (lines 11‑22) and merges results with `_merge_parallel_results` (lines 90‑115).

```python
def add_one(data, **cfg):
    return {"a": data.get("a", 0) + 1}

def add_two(data, **cfg):
    return {"b": data.get("b", 0) + 2}

builder = PipelineBuilder()
builder.register_step_handler("inc_one", add_one)
builder.register_step_handler("inc_two", add_two)

builder.add_step("step1", "inc_one", parallel_safe=True)
builder.add_step("step2", "inc_two", parallel_safe=True)
builder.set_parallelism(2)

```

### Delta Mode for Incremental Processing

Enable `delta_mode=True` to compute graph deltas before invoking handlers (see `_execute_step` lines 46‑71). This requires `base_version_id` and `target_version_id` configuration.

```python
builder.add_step(
    "kg_update",
    "kg_enrich",
    delta_mode=True,
    base_version_id="v001",
    target_version_id="v002",
    parallel_safe=False,
)

```

## Pipeline Serialization

Persist pipeline definitions for reproducibility using **`PipelineBuilder.serialize()`** (lines 60‑70) and **`PipelineSerializer`** (lines 21‑55).

```python

# Export to JSON

json_def = builder.serialize(format="json")

# Later reconstruction

from semantica.pipeline import PipelineSerializer
pipeline = PipelineSerializer.deserialize_pipeline(json_def)

```

## Summary

- **PipelineBuilder** provides a declarative DSL in [`semantica/pipeline/pipeline_builder.py`](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/pipeline_builder.py) for defining steps, dependencies, and parallelism settings.
- **ExecutionEngine** in [`semantica/pipeline/execution_engine.py`](https://github.com/semantica-agi/semantica/blob/main/semantica/pipeline/execution_engine.py) handles orchestration, concurrent layer execution via `_execute_parallel_group`, and retry logic via `_execute_step_with_retries`.
- Steps marked `parallel_safe=True` with dictionary inputs execute concurrently when grouped in the same dependency layer.
- `delta_mode=True` enables incremental graph processing by computing deltas before handler invocation.
- The engine exposes lifecycle controls (pause/resume/stop) and progress tracking through `get_pipeline_status()` and `get_progress()`.

## Frequently Asked Questions

### What is the difference between PipelineBuilder and ExecutionEngine?

**PipelineBuilder** is a compile-time tool that constructs and validates a static `Pipeline` dataclass, while **ExecutionEngine** is a runtime orchestrator that schedules steps, manages resources, and handles failures. The builder lives in [`pipeline_builder.py`](https://github.com/semantica-agi/semantica/blob/main/pipeline_builder.py) and focuses on topology; the engine in [`execution_engine.py`](https://github.com/semantica-agi/semantica/blob/main/execution_engine.py) focuses on state management and concurrency.

### How does parallel execution work in Semantica pipelines?

The engine analyzes the dependency graph to identify layers where all steps are `parallel_safe=True` and inputs are dictionaries. When `_can_run_layer_in_parallel` (lines 11‑22) returns true, `_execute_parallel_group` (lines 34‑67) distributes work across workers up to the configured parallelism level, then merges results while preserving key isolation.

### Can I serialize and reuse pipeline definitions?

Yes. Call `builder.serialize(format="json")` (lines 60‑70) to export the pipeline definition, then use `PipelineSerializer.deserialize_pipeline()` to reconstruct it later. This supports reproducible experiments, version control, and CI/CD pipelines.

### How do I handle failures and retries in pipeline execution?

The engine wraps each step in `_execute_step_with_retries` (lines 37‑55), which applies the configured retry policy before marking a step as failed. After retries exhaust, the `FailureHandler` determines whether to fail the entire pipeline or continue, based on step configuration and global engine settings.