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

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) (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.

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.


# 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.

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) (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.

engine = ExecutionEngine()

For delta processing workflows, inject additional dependencies:

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.

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:


# 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).

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.

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).


# 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 for defining steps, dependencies, and parallelism settings.
  • ExecutionEngine in 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 and focuses on topology; the engine in 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.

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 →