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.
Key capabilities
Section titled “Key capabilities”- Natural-language pipeline generation.
generate_pipelinedecomposes 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_pipelineruns 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_pipelinewrites 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_templatescovers 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.
Example prompts
Section titled “Example prompts”“Build a daily pipeline that loads new orders from Postgres into Snowflake, deduplicates on order_id, and merges into the
analytics.ordersIceberg 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.”
Connect it to your stack
Section titled “Connect it to your stack”- 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.
Works before you connect anything
Section titled “Works before you connect anything”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.
Limits, honestly
Section titled “Limits, honestly”- In 🟡 Evaluation,
deploy_pipelinerecords 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.