2 min read

    ๐Ÿ† Capstone Project: End-to-End Lakehouse

    #databricks#capstone#end-to-end

    Capstone Project: End-to-End Lakehouse

    This project connects the entire Databricks ecosystem. We will build a pipeline that ingests raw JSON sales data, cleans it, tracks historical changes, predicts future sales using Machine Learning, and deploys the entire workflow via CI/CD.

    Step 1: The Architecture

    mermaid

    Step 2: The DLT Pipeline (Bronze -> Silver -> Gold)

    We create a single file sales_pipeline.py deployed as a Delta Live Tables workflow.

    python
    import dlt
    from pyspark.sql.functions import col, sum, current_timestamp
    
    # 1. BRONZE (Ingest)
    @dlt.table(name="bronze_sales")
    def ingest_sales():
        return (
            spark.readStream.format("cloudFiles")
            .option("cloudFiles.format", "json")
            .option("cloudFiles.inferColumnTypes", "true")
            .load("s3://company-data/raw/sales/")
        )
    
    # 2. SILVER (Cleanse)
    @dlt.table(name="silver_sales")
    @dlt.expect_or_drop("valid_amount", "amount > 0")
    def clean_sales():
        return dlt.read_stream("bronze_sales").dropDuplicates(["transaction_id"])
    
    # 3. GOLD (Aggregate for BI)
    @dlt.table(name="gold_daily_sales")
    def aggregate_sales():
        return (
            dlt.read("silver_sales")
            .groupBy("store_id")
            .agg(sum("amount").alias("total_sales"))
        )
    

    Step 3: Model Training (MLflow)

    Once the gold_daily_sales table is populated, a Data Scientist runs the following notebook to train a model to predict tomorrow's sales.

    python
    import mlflow
    from mlflow.models.signature import infer_signature
    from sklearn.linear_model import LinearRegression
    
    # Load Gold Data
    df = spark.read.table("enterprise_catalog.sales_schema.gold_daily_sales").toPandas()
    X = df'store_id'
    y = df['total_sales']
    
    with mlflow.start_run():
        model = LinearRegression().fit(X, y)
        
        # Log model to Unity Catalog Registry
        signature = infer_signature(X, y)
        mlflow.sklearn.log_model(
            sk_model=model,
            artifact_path="sales_model",
            signature=signature,
            registered_model_name="enterprise_catalog.sales_schema.store_predictor"
        )
    

    Step 4: CI/CD Deployment (Databricks Asset Bundles)

    We deploy both the DLT Pipeline and the Model Training Notebook automatically via Git and a databricks.yml file.

    yaml
    bundle:
      name: end_to_end_sales
    
    resources:
      jobs:
        sales_master_workflow:
          name: "Sales E2E Pipeline"
          tasks:
            # Task 1: Run the DLT Pipeline
            - task_key: "run_dlt"
              pipeline_task:
                pipeline_id: "your_dlt_pipeline_id"
                
            # Task 2: Train the Model (Runs ONLY if DLT succeeds)
            - task_key: "train_model"
              depends_on:
                - task_key: "run_dlt"
              notebook_task:
                notebook_path: "./src/train_model.py"
              job_cluster_key: "standard_cluster"
    

    Deployment:

    1. git push to the main branch.
    2. A GitHub Action automatically triggers databricks bundle deploy -t prod.
    3. The enterprise pipeline is successfully deployed to production!

    โ† ๐Ÿ’ป Hands-On Lab & Implementation Guide | Syllabus โ†’ Syllabus