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
PipelineBuilderclass insemantica/pipeline/pipeline_builder.pyprovides the fluent DSL entry point for declaring data-processing workflows. - Steps are defined as
PipelineStepdataclasses within an immutablePipelineobject, ensuring structural integrity. PipelineValidatorautomatically checks for circular dependencies and missing references whenbuild()is called.- The
ExecutionEngineorchestrates runtime behavior, respecting dependency chains while leveraging theParallelismManagerfor concurrent step execution. - Advanced features include custom handler registration,
parallel_safe=Truehints for concurrency, anddelta_modefor 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →