Build a Databricks Delta Live Tables Pipeline Using SQL and Medallion Architecture

From raw ingestion to gold KPI marts — a production-ready sales analytics pipeline

You have raw sales data arriving every day. Orders, customers, products. Different formats, different update frequencies, different levels of trust.

Your job is to turn this into something a business analyst can query confidently.

This is the exact problem Lakeflow Declarative Pipelines (formerly Delta Live Tables) was designed to solve. And in this guide, you will build a complete pipeline using SQL and the Medallion Architecture to create production-ready sales analytics.

No shortcuts. No black boxes. Every layer explained.

Why Lakeflow Declarative Pipelines

Before Lakeflow, building a medallion pipeline meant writing Spark jobs, managing checkpoints, handling retries, and orchestrating dependencies manually.

Lakeflow Declarative Pipelines changes this by making pipelines declarative. You define what each table should look like. The system figures out the execution order, manages state, and handles failures. Databricks Lakeflow is now GA, unifying ingestion, transformation, and orchestration.

For SQL-first data engineers, this is especially powerful. You write SQL. Lakeflow handles the rest.

The Architecture

This pipeline follows the standard Medallion pattern:

The source data consists of three tables:

  1. customers — Customer master data (customer_id, name, email, region, signup_date)
  2. products — Product catalog (product_id, product_name, category, unit_price)
  3. sales — Transaction records (sale_id, customer_id, product_id, quantity, sale_date, total_amount)

Bronze Layer: Raw Ingestion

The Bronze layer is your safety net. It captures everything from source, exactly as it arrives.

-- Bronze: Raw customer data
CREATE OR REFRESH STREAMING LIVE TABLE bronze_customers
COMMENT "Raw customer data ingested from source"
AS SELECT
  *,
  current_timestamp() AS _ingested_at,
  input_file_name() AS _source_file
FROM cloud_files(
  "/mnt/source/customers",
  "csv",
  map("header", "true", "inferSchema", "true")
);
-- Bronze: Raw product catalog
CREATE OR REFRESH STREAMING LIVE TABLE bronze_products
COMMENT "Raw product catalog from source system"
AS SELECT
  *,
  current_timestamp() AS _ingested_at,
  input_file_name() AS _source_file
FROM cloud_files(
  "/mnt/source/products",
  "csv",
  map("header", "true", "inferSchema", "true")
);
-- Bronze: Raw sales transactions
CREATE OR REFRESH STREAMING LIVE TABLE bronze_sales
COMMENT "Raw sales transactions ingested incrementally"
AS SELECT
  *,
  current_timestamp() AS _ingested_at,
  input_file_name() AS _source_file
FROM cloud_files(
  "/mnt/source/sales",
  "csv",
  map("header", "true", "inferSchema", "true")
);

Notice the pattern. Every Bronze table adds two metadata columns:

This is not optional in production. When something goes wrong, these columns tell you exactly when and where the data entered your system.

Data Quality Expectations

This is where Lakeflow separates itself from raw Spark jobs. Expectations let you define data quality rules declaratively.

-- Silver: Cleaned sales with quality gates
CREATE OR REFRESH STREAMING LIVE TABLE silver_sales (
  CONSTRAINT valid_sale_id
    EXPECT (sale_id IS NOT NULL) ON VIOLATION DROP ROW,
  CONSTRAINT valid_quantity
    EXPECT (quantity > 0) ON VIOLATION DROP ROW,
  CONSTRAINT valid_total
    EXPECT (total_amount >= 0) ON VIOLATION DROP ROW,
  CONSTRAINT valid_sale_date
    EXPECT (sale_date IS NOT NULL
      AND sale_date <= current_date())
    ON VIOLATION DROP ROW
)
COMMENT "Cleaned and validated sales transactions"
AS SELECT
  CAST(sale_id AS INT) AS sale_id,
  CAST(customer_id AS INT) AS customer_id,
  CAST(product_id AS INT) AS product_id,
  CAST(quantity AS INT) AS quantity,
  CAST(sale_date AS DATE) AS sale_date,
  CAST(total_amount AS DECIMAL(10,2)) AS total_amount,
  _ingested_at
FROM STREAM(LIVE.bronze_sales);

Three violation strategies exist in Lakeflow:

For most production pipelines, DROP ROW on critical fields and logging on soft constraints is the right balance.

Silver Layer: Cleaned and Enriched

Silver is where your data becomes trustworthy. You clean types, validate constraints, and enrich with dimension data.

-- Silver: Cleaned customers
CREATE OR REFRESH STREAMING LIVE TABLE silver_customers (
  CONSTRAINT valid_customer_id
    EXPECT (customer_id IS NOT NULL) ON VIOLATION DROP ROW,
  CONSTRAINT valid_email
    EXPECT (email IS NOT NULL
      AND email LIKE "%@%")
    ON VIOLATION DROP ROW
)
COMMENT "Cleaned customer dimension"
AS SELECT
  CAST(customer_id AS INT) AS customer_id,
  TRIM(name) AS customer_name,
  LOWER(TRIM(email)) AS email,
  TRIM(region) AS region,
  CAST(signup_date AS DATE) AS signup_date,
  _ingested_at
FROM STREAM(LIVE.bronze_customers);
-- Silver: Cleaned products
CREATE OR REFRESH STREAMING LIVE TABLE silver_products (
  CONSTRAINT valid_product_id
    EXPECT (product_id IS NOT NULL) ON VIOLATION DROP ROW,
  CONSTRAINT valid_price
    EXPECT (unit_price > 0) ON VIOLATION DROP ROW
)
COMMENT "Cleaned product dimension"
AS SELECT
  CAST(product_id AS INT) AS product_id,
  TRIM(product_name) AS product_name,
  TRIM(category) AS category,
  CAST(unit_price AS DECIMAL(10,2)) AS unit_price,
  _ingested_at
FROM STREAM(LIVE.bronze_products);

Now the enrichment. This is where Silver becomes genuinely useful. We join sales with customer and product dimensions to create a single, enriched fact table.

-- Silver: Enriched sales with customer and product details
CREATE OR REFRESH LIVE TABLE silver_sales_enriched
COMMENT "Sales enriched with customer and product dimensions"
AS SELECT
  s.sale_id,
  s.sale_date,
  s.quantity,
  s.total_amount,
  c.customer_name,
  c.email AS customer_email,
  c.region AS customer_region,
  p.product_name,
  p.category AS product_category,
  p.unit_price,
  (s.quantity * p.unit_price) AS calculated_revenue,
  DATEDIFF(s.sale_date, c.signup_date) AS days_since_signup
FROM LIVE.silver_sales s
LEFT JOIN LIVE.silver_customers c
  ON s.customer_id = c.customer_id
LEFT JOIN LIVE.silver_products p
  ON s.product_id = p.product_id;

A few things to notice:

We use LEFT JOIN, not INNER JOIN. In production, you want to see orphaned records, not silently lose them.

We compute calculated_revenue from quantity and unit price. This lets you compare against the source total_amount to catch pricing discrepancies.

The days_since_signup field is a simple feature that immediately unlocks cohort analysis downstream.

Gold Layer: KPI Marts

Gold tables are designed for specific business questions. They are pre-aggregated, fast to query, and easy for analysts to understand.

-- Gold: Daily sales summary
CREATE OR REFRESH LIVE TABLE gold_daily_sales
COMMENT "Daily sales KPI summary for BI dashboards"
AS SELECT
  sale_date,
  COUNT(DISTINCT sale_id) AS total_orders,
  SUM(quantity) AS total_units_sold,
  SUM(total_amount) AS total_revenue,
  ROUND(SUM(total_amount) / COUNT(DISTINCT sale_id), 2)
    AS avg_order_value,
  COUNT(DISTINCT customer_name) AS unique_customers
FROM LIVE.silver_sales_enriched
GROUP BY sale_date;
-- Gold: Product performance
CREATE OR REFRESH LIVE TABLE gold_product_performance
COMMENT "Product-level KPIs for category analysis"
AS SELECT
  product_category,
  product_name,
  COUNT(DISTINCT sale_id) AS total_orders,
  SUM(quantity) AS total_units_sold,
  SUM(total_amount) AS total_revenue,
  ROUND(AVG(total_amount), 2) AS avg_sale_amount,
  MIN(sale_date) AS first_sale_date,
  MAX(sale_date) AS last_sale_date
FROM LIVE.silver_sales_enriched
GROUP BY product_category, product_name;
-- Gold: Regional customer insights
CREATE OR REFRESH LIVE TABLE gold_regional_insights
COMMENT "Regional performance KPIs"
AS SELECT
  customer_region,
  COUNT(DISTINCT customer_name) AS total_customers,
  COUNT(DISTINCT sale_id) AS total_orders,
  SUM(total_amount) AS total_revenue,
  ROUND(SUM(total_amount) / 
    COUNT(DISTINCT customer_name), 2)
    AS revenue_per_customer,
  ROUND(AVG(days_since_signup), 0)
    AS avg_days_since_signup
FROM LIVE.silver_sales_enriched
GROUP BY customer_region;

Each Gold table answers a specific question:

This is the pattern that works in production. Not one massive table that tries to answer everything. Focused, purpose-built marts.

The Pipeline Lineage DAG

When you deploy this pipeline in Databricks, the UI automatically generates a lineage DAG showing how data flows through each layer.

Your DAG will look like this:

bronze_customers → silver_customers ──┐
bronze_products  → silver_products  ──┤
bronze_sales     → silver_sales     ──┼→ silver_sales_enriched
                                      │
                                      ├→ gold_daily_sales
                                      ├→ gold_product_performance
                                      └→ gold_regional_insights

This is not just documentation. The DAG is functional. Lakeflow uses it to determine execution order, parallelism, and dependency resolution. If bronze_customers fails, the system knows not to update silver_sales_enriched.

Incremental Processing

Notice that Bronze and Silver streaming tables use CREATE OR REFRESH STREAMING LIVE TABLE with STREAM(). This means Lakeflow processes only new data on each run.

No manual checkpoint management. No "find the last processed timestamp" logic. Lakeflow tracks state internally using its own checkpointing mechanism.

For the enriched and Gold tables, we use standard LIVE TABLE (not streaming). These are recomputed fully on each pipeline run. This is intentional. Aggregations and joins across multiple streams are simpler and more reliable when recomputed.

The pattern is:

Running the Pipeline

To deploy this pipeline in Databricks:

  1. Create a new pipeline from the Jobs & Pipelines page
  2. Point it to your SQL notebook containing all the table definitions
  3. Set the target schema where output tables will be created
  4. Choose between Triggered (batch) or Continuous mode
  5. Click Start

For development, use Development mode. It provides faster iteration and better error messages. Switch to Production mode when you are ready for reliable, cost-optimized runs.

What You Learned

This pipeline covers the core patterns you will use in every production project:

The beauty of Lakeflow is that the pipeline definition is the documentation. Anyone can read your SQL and understand exactly what each table does, where data comes from, and what quality rules apply.

The best data pipelines are the ones where the code reads like a specification. Lakeflow makes that possible.

If you are building your first production pipeline or refactoring an existing one, this pattern scales. Start with three layers. Keep each table focused. Let Lakeflow handle the orchestration.

That is how production lakehouse ETL is built.


This pipeline connects directly to concepts covered in the Medallion Architecture, Delta Lake, Data Quality, and Jobs & Pipelines chapters of BricksNotes. The complete project README with architecture diagrams, DAG flow, and all SQL code is available as a downloadable reference.


Download the Complete Project

Want the full project README with architecture diagrams, DAG flow, folder structure, and every SQL file ready to use?

📥 Download the DLT Sales Analytics Pipeline README

This document includes everything you need to recreate this pipeline in your own Databricks workspace — from source table schemas to Gold KPI mart definitions.