Replace plain markdown lists, tables, and numbered steps with interactive Mintlify v3 MDX components across all 50+ documentation files: - Tabs: provider/parser/method selection guides, citation formats, component details - Steps: setup flows, pipeline stages, connection initialization - CardGroup/Card: feature overviews, "what you get" sections, navigation footers - AccordionGroup: FAQ entries - Check/Warning/Tip/Note/Info: callouts replacing plain bold text and inline notes Files improved span the full docs surface: reference modules (context, llms, kg, reasoning, embeddings, deduplication, provenance, parse, ontology, core, semantic_extract), integrations (agno, docling, snowflake), graph/vector store backends (apache_age, pgvector), and top-level guides (contributing, governance, glossary, citation, community-projects, learning-more).
21 KiB
title, description, icon
| title | description | icon |
|---|---|---|
| Pipeline Module | Pipeline DSL with parallel workers, retry policies, failure handling, and progress tracking. | gear |
semantica.pipeline lets you chain Semantica components into reproducible, fault-tolerant workflows:
- Per-step failure strategies:
skip,retry,abort, orfallback - Parallel workers via
ParallelismManager— thread or process pool PipelineValidatorcatches cycles, missing handlers, and config errors before running- Pre-built templates:
"document_processing","rag_pipeline","kg_construction","ontology_generation" - Pipelines are serializable to YAML — save and reload in any environment
Exported Classes
| Class | Role |
|---|---|
PipelineBuilder |
DSL for wiring steps: add_step, connect_steps, set_parallel, build |
ExecutionEngine |
Runs a built pipeline: execute_pipeline(pipeline, data) → ExecutionResult |
ExecutionResult |
{success, output, metadata, metrics, errors} — full run summary |
FailureHandler |
Per-step strategy: skip, retry, abort, or fallback on failure |
ParallelismManager |
Thread or process pool for concurrent step execution with configurable workers |
PipelineValidator |
Catches dependency cycles, missing handlers, and config errors before running |
PipelineTemplateManager |
Pre-built templates: "document_processing", "rag_pipeline", "kg_construction", "ontology_generation" |
Why Use a Pipeline?
You could wire Semantica modules together with plain Python code. Pipelines add:
A single bad document doesn't crash a 10,000-document run. Run extraction across multiple workers with one parameter. tqdm console bar or WebSocket streaming to Explorer. Save the exact pipeline configuration to YAML and replay on any machine. On re-runs, only process documents that changed since the last run. Catch misconfigured steps and dependency cycles before they fail mid-run. Use plain module calls for quick scripts and notebooks. Use pipelines for anything you run repeatedly, at scale, or in production.<img src="/assets/img/diagrams/pipeline-flow.svg" alt="Pipeline step sequence: Ingest → Parse → Normalize → Extract → Build KG → QA → Store → Deliver" style={{ width: '100%', borderRadius: '10px', margin: '0 0 24px' }} />
Quick Start
```python from semantica.pipeline import PipelineBuilder from semantica.ingest import FileIngestor from semantica.parse import DocumentParser from semantica.semantic_extract import NERExtractor from semantica.kg import GraphBuilderingestor = FileIngestor()
parser = DocumentParser()
extractor = NERExtractor(method="ml")
kg_builder = GraphBuilder(merge_entities=True)
builder = PipelineBuilder()
builder.add_step("ingest", "file_ingest", handler=ingestor.ingest_file)
builder.add_step("parse", "document_parse", handler=parser.parse)
builder.add_step("extract", "ner_extract", handler=extractor.extract)
builder.add_step("build_kg", "graph_build", handler=kg_builder.build)
builder.connect_steps("ingest", "parse")
builder.connect_steps("parse", "extract")
builder.connect_steps("extract","build_kg")
pipeline = builder.build("my_pipeline")
```
validator = PipelineValidator()
result = validator.validate_pipeline(pipeline)
if not result.valid:
for error in result.errors: # errors is List[str]
print(f"Error: {error}")
for warning in result.warnings:
print(f"Warning: {warning}")
```
engine = ExecutionEngine()
result = engine.execute_pipeline(pipeline, data="data/")
kg = result.output
print(f"Success: {result.success}")
print(f"Steps executed: {result.metrics['steps_executed']}")
print(f"Steps failed: {result.metrics['steps_failed']}")
print(f"Duration: {result.metrics['execution_time']:.1f}s")
```
Parallel Processing
Set parallelism on the builder and pass max_workers to ExecutionEngine:
from semantica.pipeline import PipelineBuilder, ExecutionEngine
builder = PipelineBuilder()
builder.add_step("ingest", "file_ingest", handler=ingestor.ingest_file)
builder.add_step("parse", "document_parse", handler=parser.parse)
builder.add_step("extract", "ner_extract", handler=extractor.extract)
builder.add_step("build", "graph_build", handler=kg_builder.build)
builder.set_parallelism(4)
pipeline = builder.build("parallel_pipeline")
engine = ExecutionEngine(max_workers=4)
result = engine.execute_pipeline(pipeline, data="data/")
Retry and Error Handling
```python from semantica.pipeline import RetryPolicy, RetryStrategy, FailureHandler, ExecutionEnginepolicy = RetryPolicy(
max_retries=3,
strategy=RetryStrategy.EXPONENTIAL,
initial_delay=1.0, # 1s → 2s → 4s
backoff_factor=2.0
)
handler = FailureHandler()
handler.set_retry_policy("ner_extract", policy) # keyed by step_type
engine = ExecutionEngine(default_max_retries=3, default_backoff_factor=2.0)
result = engine.execute_pipeline(pipeline, data="data/")
```
Best for transient API errors and rate limits — waits longer with each retry, giving upstream services time to recover.
policy = RetryPolicy(
max_retries=3,
strategy=RetryStrategy.LINEAR,
initial_delay=2.0 # 2s → 4s → 6s
)
```
Use when the delay between retries should grow predictably — e.g., waiting for a database lock to release.
policy = RetryPolicy(
max_retries=5,
strategy=RetryStrategy.FIXED,
initial_delay=1.0 # 1s every attempt
)
```
Use when retrying against a service with a fixed cooldown window.
Failure Strategies
| Strategy | Behaviour | When to Use |
|---|---|---|
"skip" |
Log failure, continue to next document | Production — one bad doc shouldn't stop 10k |
"stop" |
Raise exception immediately | Development — surface errors fast |
"retry" |
Retry via RetryPolicy, then skip |
When failures are likely transient |
Progress Tracking
```python from semantica.pipeline import ExecutionEngineengine = ExecutionEngine()
result = engine.execute_pipeline(pipeline, data="data/")
# The progress tracker outputs tqdm bars to the console during execution
```
Displays a live progress bar in the terminal via Semantica's built-in progress tracker. Best for scripts and CLI tools.
engine = ExecutionEngine()
# Run in a background thread, poll progress from main thread
def run():
engine.execute_pipeline(pipeline, data="data/")
t = threading.Thread(target=run, daemon=True)
t.start()
while t.is_alive():
progress = engine.get_progress(pipeline.name)
if progress:
print(f" {progress['completed_steps']}/{progress['total_steps']} steps — {progress['status']}")
time.sleep(2)
```
Poll `get_progress()` for live status during execution.
Pipeline DSL
PipelineBuilder uses add_step(name, type, **config) and connect_steps(from, to) to define a DAG:
from semantica.pipeline import PipelineBuilder, ExecutionEngine
builder = PipelineBuilder()
# Add steps — step_type is a string label, handler is the callable invoked at runtime
builder.add_step("ingest", "file_ingest", handler=ingestor.ingest_file)
builder.add_step("parse", "document_parse", handler=parser.parse)
builder.add_step("normalize", "text_normalize", handler=normalizer.normalize)
builder.add_step("extract", "ner_extract", handler=extractor.extract)
builder.add_step("rel_extract", "rel_extract", handler=rel_extractor.extract)
builder.add_step("build_kg", "graph_build", handler=kg_builder.build)
builder.add_step("deduplicate", "dedup", handler=deduplicator.deduplicate)
builder.add_step("export", "rdf_export", handler=exporter.export, format="turtle", path="output.ttl")
# Wire the data flow
builder.connect_steps("ingest", "parse")
builder.connect_steps("parse", "normalize")
builder.connect_steps("normalize", "extract")
builder.connect_steps("extract", "rel_extract")
builder.connect_steps("rel_extract", "build_kg")
builder.connect_steps("build_kg", "deduplicate")
builder.connect_steps("deduplicate", "export")
pipeline = builder.build("full_pipeline")
result = ExecutionEngine().execute_pipeline(pipeline, data="data/")
Serialize and Restore Pipelines
PipelineSerializer converts a pipeline to JSON or dict for storage and reloads it later:
from semantica.pipeline import PipelineSerializer
serializer = PipelineSerializer()
# Serialize to JSON string
json_str = serializer.serialize_pipeline(pipeline, format="json")
# Save to file
with open("pipeline_config.json", "w") as f:
f.write(json_str)
# Restore on any machine and execute
with open("pipeline_config.json") as f:
restored = serializer.deserialize_pipeline(f.read())
result = ExecutionEngine().execute_pipeline(restored, data="data/")
Pre-Built Templates
PipelineTemplateManager wires common workflows with the correct step order — no manual wiring required:
from semantica.pipeline import PipelineTemplateManager
manager = PipelineTemplateManager()
The create_pipeline_from_template(name) method returns a configured PipelineBuilder. Call .build(pipeline_name) on it to produce a runnable Pipeline.
Complete document processing from ingestion to knowledge graph.
```python
builder = manager.create_pipeline_from_template("document_processing")
pipeline = builder.build("doc_pipeline")
```
RAG pipeline for question answering — builds a vector-indexed store.
```python
builder = manager.create_pipeline_from_template("rag_pipeline")
pipeline = builder.build("rag_pipeline")
```
Knowledge graph construction from multiple sources.
```python
builder = manager.create_pipeline_from_template("kg_construction")
pipeline = builder.build("kg_pipeline")
```
Ontology generation from extracted data.
```python
builder = manager.create_pipeline_from_template("ontology_generation")
pipeline = builder.build("ontology_pipeline")
```
ExecutionEngine
Fine-grained control over pipeline execution — pause, resume, cancel, and inspect live progress:
from semantica.pipeline import ExecutionEngine
engine = ExecutionEngine(max_workers=4)
# pipeline.name is the pipeline ID used for all control operations
result = engine.execute_pipeline(pipeline, data="data/")
pipeline_id = pipeline.name # e.g. "my_pipeline"
# Pause after the current step finishes
engine.pause_pipeline(pipeline_id)
progress = engine.get_progress(pipeline_id)
print(f"Completed: {progress['completed_steps']}/{progress['total_steps']}")
print(f"Status: {progress['status']}")
engine.resume_pipeline(pipeline_id)
engine.stop_pipeline(pipeline_id)
| Method | Returns | Description |
|---|---|---|
execute_pipeline(pipeline, data) |
ExecutionResult |
Execute pipeline from start to finish |
get_pipeline_status(pipeline_id) |
PipelineStatus |
Current state (RUNNING, PAUSED, STOPPED) |
get_progress(pipeline_id) |
Dict |
completed_steps, total_steps, progress_percentage, status |
pause_pipeline(pipeline_id) |
None |
Suspend after current step completes |
resume_pipeline(pipeline_id) |
None |
Resume from paused state |
stop_pipeline(pipeline_id) |
None |
Cancel and clean up immediately |
PipelineValidator
Catches problems before they surface as mid-run failures:
from semantica.pipeline import PipelineValidator
validator = PipelineValidator()
result = validator.validate_pipeline(pipeline)
if result.valid:
print("Pipeline is valid — safe to run")
else:
for error in result.errors: # errors is List[str]
print(f"Error: {error}")
for warning in result.warnings: # warnings is List[str]
print(f"Warning: {warning}")
Checks performed:
- Dependency cycle detection — A depends on B, B depends on A
- Step type validation — each step type must be registered
- Connection integrity — referenced step names must exist
- Configuration completeness — required parameters must be present
ParallelismManager
```python from semantica.pipeline import ParallelismManager, Task# use_processes=False (default) → thread pool for I/O-bound tasks
manager = ParallelismManager(max_workers=8, use_processes=False)
tasks = [Task(task_id=f"t{i}", handler=ner.extract, args=(text,)) for i, text in enumerate(texts)]
results = manager.execute_parallel(tasks)
# returns List[ParallelExecutionResult]
successes = [r for r in results if r.success]
failures = [r for r in results if not r.success]
```
Use thread pools for **I/O-bound** steps: web fetching, database queries, API calls.
# use_processes=True → process pool, bypasses Python GIL
manager = ParallelismManager(max_workers=4, use_processes=True)
tasks = [Task(task_id=f"t{i}", handler=embedder.generate_embeddings, args=(chunk,)) for i, chunk in enumerate(chunks)]
results = manager.execute_parallel(tasks)
```
Use process pools for **CPU-bound** steps: embedding, OCR, large NER batches.
ResourceScheduler
Prevents memory oversubscription on large runs:
from semantica.pipeline import ResourceScheduler, ExecutionEngine
scheduler = ResourceScheduler()
engine = ExecutionEngine()
resources = scheduler.allocate_resources(pipeline)
try:
result = engine.execute_pipeline(pipeline, data="data/")
finally:
scheduler.release_resources(resources)
Delta Mode
Re-process only data that has changed since the last run:
from semantica.pipeline import PipelineBuilder, ExecutionEngine
builder = PipelineBuilder()
# delta_mode=True tells ExecutionEngine to compute the diff between two snapshots
# and pass only changed triples to this step's handler
builder.add_step(
"ingest", "file_ingest",
handler=ingestor.ingest_file,
delta_mode=True, base_version_id="v1", target_version_id="v2"
)
builder.add_step(
"extract", "ner_extract",
handler=extractor.extract,
delta_mode=True, base_version_id="v1", target_version_id="v2"
)
builder.add_step(
"build", "graph_build",
handler=kg_builder.build,
delta_mode=False # always rebuild the merged graph
)
builder.connect_steps("ingest", "extract")
builder.connect_steps("extract", "build")
pipeline = builder.build("delta_pipeline")
engine = ExecutionEngine()
result = engine.execute_pipeline(
pipeline,
data="data/",
version_manager=version_manager, # required for delta mode
triplet_store=triplet_store # required for delta mode
)
Schemas
@dataclass
class ExecutionResult:
success: bool # True if all steps completed without failure
output: Any # output from the final pipeline step
metadata: Dict[str, Any] # {"pipeline_id": "...", "execution_time": 1.23}
metrics: Dict[str, Any] # {"steps_executed": 4, "steps_failed": 0, "execution_time": 1.23}
errors: List[str] # error messages from failed steps (empty on full success)
# Access pattern
result.success # bool
result.output # final step output
result.metadata["pipeline_id"] # pipeline name used as ID
result.metadata["execution_time"] # total wall-clock seconds
result.metrics["steps_executed"] # count of successfully completed steps
result.metrics["steps_failed"] # count of failed steps
result.errors # List[str] of error messages
@dataclass
class PipelineStep:
name: str
step_type: str
config: Dict[str, Any]
dependencies: List[str] # names of steps this step waits for
handler: Optional[Callable]
status: StepStatus
result: Any
error: Optional[Exception]
delta_mode: bool # True = process only changed data
base_version_id: Optional[str] # snapshot ID to diff against
target_version_id: Optional[str] # snapshot ID being produced
from semantica.pipeline import StepStatus
StepStatus.PENDING # Not yet started
StepStatus.RUNNING # Currently executing
StepStatus.COMPLETED # Finished successfully
StepStatus.FAILED # Error occurred — check step.error
StepStatus.SKIPPED # Skipped due to FailureHandler "skip" strategy