Skip to content

Pipelines & Ingestion (dw-pipelines)

The Pipelines & Ingestion agent turns a plain-English description of what you want moved and transformed into a working pipeline. It decomposes your description into extraction, transformation, loading, testing, and deployment tasks, generates the code (SQL, Python, or dbt), and targets the orchestrator you already run — Airflow, Dagster, or Prefect. It also covers the ingestion side: EL, CDC, and replication patterns, including Iceberg MERGE INTO.

It doesn’t work alone. When it generates a pipeline, it registers the new asset in the catalog, asks the Quality agent to create quality tests for it, and checks schema compatibility — so a pipeline born here arrives with context, tests, and lineage instead of as an orphan.

  • Natural-language pipeline generation. generate_pipeline decomposes a description into extract/transform/load/test/deploy tasks and generates the code, with template fallback when generation isn’t confident.
  • Real validation before anything ships. validate_pipeline runs sandbox execution: Python AST parsing, SQL syntax checking, and YAML schema validation, plus semantic-layer validation when the catalog agent is reachable.
  • Airflow deployment with verification. deploy_pipeline writes DAG files via filesystem, S3, or git-sync and verifies the deployment through the Airflow REST API — it doesn’t just drop a file and hope.
  • Versioned specs in Git. Deployment can commit the pipeline specification as YAML to your repo, so every deployed pipeline has a reviewable history.
  • A template library for common patterns. list_pipeline_templates covers ETL, ELT, CDC, streaming, reverse-ETL, and data-quality patterns, filterable by orchestrator — use one as a starting point instead of generating from scratch.
  • Cross-agent registration. Generated pipelines are registered in the catalog, get quality tests created, and are checked for schema compatibility automatically.

“Build a daily pipeline that loads new orders from Postgres into Snowflake, deduplicates on order_id, and merges into the analytics.orders Iceberg table.”

“Validate this pipeline spec before I deploy it — check the SQL and the Python.”

“What CDC templates do you have for Airflow?”

“Deploy the validated orders pipeline to staging and commit the spec to the main branch.”

  • Orchestration — Airflow (deployment and verification), plus Dagster and Prefect as generation targets.
  • Warehouses and lakehouses — Snowflake, BigQuery, Databricks as pipeline sources and targets.
  • dbt — as a code language for generated transformations.

See the connector catalog for setup.

The agent starts in 🟡 Evaluation on built-in sample data — the generation, templating, and validation logic is the real thing, run against a realistic sample estate. It earns 🟢 Connected per system through a passing live test. See Verify your setup.

  • In 🟡 Evaluation, deploy_pipeline records the deployment locally — nothing reaches a real orchestrator until Airflow is configured and verified.
  • The verified deployment path is Airflow today. Dagster and Prefect are supported as generation targets; deployment to them is not yet implemented.
  • Semantic-layer validation only runs when the Catalog & Context agent is reachable — the syntax and sandbox checks still run without it.
  • Deploying is a governed write: it follows the propose–approve–receipt path like every other write in the swarm.