# Databricks DLT Sales Analytics Pipeline

**A production-ready Delta Live Tables pipeline using SQL and Medallion Architecture**

Built for data engineers learning lakehouse ETL design with Databricks.

---

## Architecture Overview

```
┌─────────────────────────────────────────────────────────────────┐
│                    DLT PIPELINE: Sales Analytics                │
├─────────────────────────────────────────────────────────────────┤
│                                                                 │
│  SOURCE FILES          BRONZE              SILVER         GOLD  │
│  (/mnt/source/)        (Raw)             (Cleaned)      (KPIs) │
│                                                                 │
│  customers.csv ──→ bronze_customers ──→ silver_customers ──┐    │
│                                                            │    │
│  products.csv  ──→ bronze_products  ──→ silver_products  ──┤    │
│                                                            │    │
│  sales.csv     ──→ bronze_sales     ──→ silver_sales     ──┤    │
│                                                            │    │
│                                        silver_sales_enriched    │
│                                              │                  │
│                                   ┌──────────┼──────────┐       │
│                                   │          │          │       │
│                              gold_daily  gold_product  gold_    │
│                              _sales      _performance  regional │
│                                                        _insights│
└─────────────────────────────────────────────────────────────────┘
```

## DLT Lineage DAG

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

## Project Structure

```
dlt-sales-analytics/
│
├── notebooks/
│   ├── 01_bronze_layer.sql        # Raw ingestion tables
│   ├── 02_silver_layer.sql        # Cleaned + enriched tables
│   ├── 03_gold_layer.sql          # KPI mart tables
│   └── 00_full_pipeline.sql       # All-in-one pipeline definition
│
├── data/
│   ├── customers.csv              # Customer master data
│   ├── products.csv               # Product catalog
│   └── sales.csv                  # Transaction records
│
├── tests/
│   └── quality_checks.sql         # Standalone quality validation
│
└── README.md                      # This file
```

## Source Data

### customers.csv
| Column | Type | Description |
|--------|------|-------------|
| customer_id | INT | Unique customer identifier |
| name | STRING | Customer full name |
| email | STRING | Customer email address |
| region | STRING | Geographic region |
| signup_date | DATE | Account creation date |

### products.csv
| Column | Type | Description |
|--------|------|-------------|
| product_id | INT | Unique product identifier |
| product_name | STRING | Product display name |
| category | STRING | Product category |
| unit_price | DECIMAL | Price per unit |

### sales.csv
| Column | Type | Description |
|--------|------|-------------|
| sale_id | INT | Unique sale identifier |
| customer_id | INT | FK to customers |
| product_id | INT | FK to products |
| quantity | INT | Units purchased |
| sale_date | DATE | Transaction date |
| total_amount | DECIMAL | Total sale value |

---

## Complete Pipeline Code

### Bronze Layer — Raw Ingestion

```sql
-- ============================================
-- BRONZE LAYER: Raw Data Ingestion
-- ============================================

-- 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")
);
```

### Silver Layer — Cleaned, Validated, Enriched

```sql
-- ============================================
-- SILVER LAYER: Cleaned + Enriched
-- ============================================

-- Silver: Cleaned customers with quality gates
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);

-- 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);

-- Silver: Enriched sales with customer + 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;
```

### Gold Layer — KPI Marts

```sql
-- ============================================
-- GOLD LAYER: Business KPI Marts
-- ============================================

-- 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;
```

---

## Deployment Steps

1. **Upload source data** to `/mnt/source/` (or your preferred cloud storage path)
2. **Create a new DLT pipeline** from Databricks Workflows
3. **Point to the SQL notebook** containing all table definitions
4. **Set the target schema** (e.g., `sales_analytics`)
5. **Choose pipeline mode:**
   - **Development** — Faster iteration, better error messages
   - **Production** — Optimized, reliable, auto-retry on failures
6. **Click Start** and monitor the DAG in the DLT UI

## Data Quality Monitoring

DLT automatically tracks expectation metrics. After each run, check:

- **Event log** — See which records were dropped and why
- **Expectation results** — Pass/fail rates per constraint
- **Pipeline health** — Overall success rate and processing times

Query the event log:

```sql
SELECT
  details:flow_definition.output_dataset AS table_name,
  details:flow_progress.data_quality.expectations AS expectations
FROM event_log(TABLE(your_pipeline_name))
WHERE event_type = 'flow_progress';
```

## Key Design Decisions

| Decision | Rationale |
|----------|-----------|
| `STREAMING` for Bronze/Silver facts | Incremental processing, no reprocessing of old data |
| `LIVE TABLE` for enriched + Gold | Simpler, full recompute ensures consistency across joins |
| `LEFT JOIN` for enrichment | Preserves orphaned records for debugging |
| `DROP ROW` on critical fields | Production safety — bad data does not propagate |
| Metadata columns in Bronze | Audit trail for data lineage and troubleshooting |
| Separate Gold marts per domain | Focused, fast queries for specific business questions |

## Related BricksNotes Chapters

- [Medallion Architecture](/lessons/medallion-architecture) — Full pattern explanation
- [Delta Lake](/lessons/delta-lake) — Table format fundamentals
- [Data Quality](/lessons/data-quality) — Validation strategies
- [Workflows](/lessons/workflows) — Pipeline orchestration

---

**Built with BricksNotes** — [bricksnotes.com](https://bricksnotes.com)
