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.
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.
This pipeline follows the standard Medallion pattern:
The source data consists of three tables:
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:
_ingested_at — When the record was processed_source_file — Which file it came fromThis is not optional in production. When something goes wrong, these columns tell you exactly when and where the data entered your system.
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 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 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.
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_insightsThis 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.
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:
To deploy this pipeline in Databricks:
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.
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.
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.