--- title: "Pipeline Module" description: "Pipeline DSL with parallel workers, retry policies, failure handling, and progress tracking." icon: "gear" --- `semantica.pipeline` lets you chain Semantica components into reproducible, fault-tolerant workflows with parallel execution and configurable error handling. Pipelines are serializable — save them to YAML 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: `"full-qa"`, `"extract-only"`, `"kg-build"` | ## 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. Pipeline step sequence: Ingest → Parse → Normalize → Extract → Build KG → QA → Store → Deliver ## 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 GraphBuilder from semantica.llms import Groq import os llm = Groq(model="llama-3.3-70b-versatile", api_key=os.getenv("GROQ_API_KEY")) ingestor = FileIngestor() parser = DocumentParser() extractor = NERExtractor(method="llm", llm_provider=llm) 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") ``` ```python from semantica.pipeline import PipelineValidator 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}") ``` ```python from semantica.pipeline import ExecutionEngine 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`: ```python 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, ExecutionEngine policy = RetryPolicy( max_retries=3, strategy=RetryStrategy.EXPONENTIAL, initial_delay=1.0, # 1s → 2s → 4s backoff_factor=2.0 ) handler = FailureHandler() handler.retry_policies["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. ```python from semantica.pipeline import RetryPolicy, RetryStrategy 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. ```python from semantica.pipeline import RetryPolicy, RetryStrategy 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 | Always use `strategy="skip"` in production. A single malformed document shouldn't stop a pipeline processing thousands of documents. Inspect `result.errors` after the run to find and reprocess failures. ## Progress Tracking ```python from semantica.pipeline import ExecutionEngine engine = 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. ```python from semantica.pipeline import ExecutionEngine import threading, time 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: ```python 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: ```python 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/") ``` Serialized pipelines capture step names, types, and config — but not handler functions (callables can't be serialized). Re-register handlers on the restored steps before executing. ## Pre-Built Templates `PipelineTemplateManager` wires common workflows with the correct step order — no manual wiring required: ```python from semantica.pipeline import PipelineTemplateManager from semantica.llms import Groq import os llm = Groq(model="llama-3.3-70b-versatile", api_key=os.getenv("GROQ_API_KEY")) manager = PipelineTemplateManager() ``` **Ingest → Parse → Extract → Build KG** Standard knowledge base construction from documents. ```python pipeline = manager.get_template( "ingest-extract-build", llm_provider=llm ) ``` **Ingest → Parse → Embed → Index** Retrieval-augmented generation — builds a vector-indexed knowledge graph. ```python pipeline = manager.get_template( "graphrag", llm_provider=llm, vector_backend="faiss" ) ``` **Build KG → Analytics → Export Report** Graph analysis and reporting — centrality, community detection, HTML output. ```python pipeline = manager.get_template( "analytics", export_format="html" ) ``` **Ingest → Normalize → Extract → Dedup → Conflicts → Build** Production-quality KG with full data quality pipeline. ```python pipeline = manager.get_template( "full-qa", llm_provider=llm ) ``` ## ExecutionEngine Fine-grained control over pipeline execution — pause, resume, cancel, and inspect live progress: ```python 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: ```python 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 # use_processes=False (default) → thread pool for I/O-bound tasks manager = ParallelismManager(max_workers=8, use_processes=False) tasks = [{"fn": ner.extract, "args": [text]} for text in texts] results = manager.execute_parallel(tasks, timeout=60) # 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. ```python # use_processes=True → process pool, bypasses Python GIL manager = ParallelismManager(max_workers=4, use_processes=True) tasks = [{"fn": embedder.generate_embeddings, "args": [chunk]} for chunk in chunks] results = manager.execute_parallel(tasks, timeout=120) ``` Use process pools for **CPU-bound** steps: embedding, OCR, large NER batches. ## ResourceScheduler Prevents memory oversubscription on large runs: ```python 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: ```python 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 ) ``` Delta detection uses SHA-256 checksums on source content. Only sources whose checksum differs from `base_version_id` are passed to downstream steps. For pipelines that run hourly or daily against a growing corpus, delta mode eliminates redundant re-embedding and re-extraction. ## Schemas ```python @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 ``` ```python @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 ``` ```python 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 ``` ## Tips and Common Pitfalls **Use `PipelineValidator` before running in production.** It catches dependency cycles, missing step names, and misconfigured connections that would only surface as errors mid-run. Validation is instant; catching them after a 30-minute extraction job is not. **Set `workers=` based on workload type.** Thread workers for I/O-bound steps (web fetching, DB queries), process workers for CPU-bound steps (embedding, OCR, large NER batches). Mixing pool types on the wrong step type wastes resources without speed gains. **Use `failure_handler=FailureHandler(strategy="skip")` in production.** A single malformed document shouldn't stop a pipeline processing 10,000 documents. `skip` logs the failure and continues; inspect `result.errors` after the run to find and reprocess failed documents. **Use templates from `PipelineTemplateManager` for common patterns.** `get_template("full-qa")` wires up normalization, deduplication, conflict detection, and graph construction in the right order — saving you from common mistakes like deduplicating before normalizing. **Inspect `result.metrics` to find bottlenecks.** `result.metrics['steps_executed']` and `result.metrics['execution_time']` give a quick read on overall pipeline health. For per-step timing, check `step.result` on each `PipelineStep` after the run. First step in most pipelines. Core extraction step. Graph construction step. Final output step.