retail-etl-medallion-pipeline

Compare original and translation side by side

🇺🇸

Original

English
🇨🇳

Translation

Chinese

Retail ETL Medallion Pipeline Skill

零售ETL Medallion管道技能

Skill by ara.so — Data Skills collection.
技能来自ara.so — 数据技能合集。

Overview

概述

This project implements a production-grade Medallion Architecture ETL pipeline for retail/hypermarket data, handling complex business logic like inventory shrinkage, meat/poultry recipe conversions, supplier rebate tiers, and multi-branch sales consolidation. The architecture follows three data quality layers:
  • Bronze Layer: Raw data ingestion from CSV sources (sales, stock, products)
  • Silver Layer: Cleaned, standardized, and business-rule-applied data
  • Gold Layer: Aggregated, analytics-ready dimensional models
The pipeline processes:
  • Multi-branch sales transactions (Alex, Cairo, Giza)
  • Product catalogs with recipe/yield conversions
  • Stock/inventory tracking across locations
  • Supplier rebate calculations
本项目为零售/大卖场数据实现了一套生产级Medallion架构ETL管道,可处理库存损耗、肉禽配方转换、供应商返利层级、多门店销售合并等复杂业务逻辑。该架构遵循三层数据质量标准:
  • 青铜层:从CSV源(销售、库存、产品)导入原始数据
  • 白银层:经过清洗、标准化并应用业务规则的数据
  • 黄金层:聚合完成、可直接用于分析的维度模型
管道处理以下数据:
  • 多门店销售交易(Alex、Cairo、Giza)
  • 包含配方/出成率转换的产品目录
  • 跨地点的库存追踪
  • 供应商返利计算

Project Structure

项目结构

Retail-Data-Warehouse/
├── data_source/              # Raw CSV files (CRM/ERP exports)
│   ├── 000.Hypermarket Products.csv
│   ├── 001-003.*.Branch Sales.csv
│   └── 004-006.*.Stock.csv
├── sql_scripts/              # TSQL stored procedures for each layer
│   ├── 00_create_database_and_schemas.sql
│   ├── 01-04_bronze_*.sql
│   ├── 05-08_silver_*.sql
│   └── 09-12_gold_*.sql
├── BI_Team_Analysis/         # Power BI dashboards
└── docker-compose.yml        # SQL Server container setup
Retail-Data-Warehouse/
├── data_source/              # 原始CSV文件(CRM/ERP导出)
│   ├── 000.Hypermarket Products.csv
│   ├── 001-003.*.Branch Sales.csv
│   └── 004-006.*.Stock.csv
├── sql_scripts/              # 各层对应的TSQL存储过程
│   ├── 00_create_database_and_schemas.sql
│   ├── 01-04_bronze_*.sql
│   ├── 05-08_silver_*.sql
│   └── 09-12_gold_*.sql
├── BI_Team_Analysis/         # Power BI仪表盘
└── docker-compose.yml        # SQL Server容器配置

Installation & Setup

安装与配置

1. Infrastructure Setup (SQL Server)

1. 基础设施配置(SQL Server)

Using Docker:
bash
undefined
使用Docker:
bash
undefined

Start SQL Server container

启动SQL Server容器

docker-compose up -d
docker-compose up -d

Verify container is running

验证容器运行状态

docker ps | grep sqlserver

Or use an existing SQL Server instance (2017+).
docker ps | grep sqlserver

也可使用现有SQL Server实例(2017及以上版本)。

2. Database Initialization

2. 数据库初始化

bash
undefined
bash
undefined

Connect to SQL Server and create database structure

连接SQL Server并创建数据库结构

sqlcmd -S localhost -U sa -P $SQL_SA_PASSWORD -i sql_scripts/00_create_database_and_schemas.sql

This creates:
- Database: `RetailDataWarehouse`
- Schemas: `bronze`, `silver`, `gold`, `staging`
sqlcmd -S localhost -U sa -P $SQL_SA_PASSWORD -i sql_scripts/00_create_database_and_schemas.sql

该脚本会创建:
- 数据库:`RetailDataWarehouse`
- 架构:`bronze`、`silver`、`gold`、`staging`

3. Load Raw Data to Bronze Layer

3. 将原始数据加载到青铜层

Place CSV files in accessible location, then run:
sql
-- Execute Bronze layer ingestion procedures
EXEC bronze.usp_LoadProducts;
EXEC bronze.usp_LoadSales;
EXEC bronze.usp_LoadStock;
Or execute all Bronze scripts sequentially:
bash
for script in sql_scripts/01_bronze_*.sql sql_scripts/02_bronze_*.sql sql_scripts/03_bronze_*.sql sql_scripts/04_bronze_*.sql; do
    sqlcmd -S localhost -U sa -P $SQL_SA_PASSWORD -i "$script"
done
将CSV文件放置在可访问位置,然后执行:
sql
-- 执行青铜层导入存储过程
EXEC bronze.usp_LoadProducts;
EXEC bronze.usp_LoadSales;
EXEC bronze.usp_LoadStock;
或者按顺序执行所有青铜层脚本:
bash
for script in sql_scripts/01_bronze_*.sql sql_scripts/02_bronze_*.sql sql_scripts/03_bronze_*.sql sql_scripts/04_bronze_*.sql; do
    sqlcmd -S localhost -U sa -P $SQL_SA_PASSWORD -i "$script"
done

Key Architecture Patterns

核心架构模式

Bronze Layer (Raw Ingestion)

青铜层(原始数据导入)

Purpose: Land raw data with minimal transformation. Add audit columns only.
sql
-- Example: Bronze Products Table Structure
CREATE TABLE bronze.Products (
    ProductID INT,
    ProductName NVARCHAR(255),
    Category NVARCHAR(100),
    SubCategory NVARCHAR(100),
    UnitPrice DECIMAL(10,2),
    SupplierID INT,
    RecipeYield DECIMAL(5,2),  -- For meat/poultry conversions
    LoadTimestamp DATETIME2 DEFAULT GETDATE(),
    SourceFile NVARCHAR(500)
);

-- Bronze Load Pattern
CREATE PROCEDURE bronze.usp_LoadProducts
AS
BEGIN
    TRUNCATE TABLE bronze.Products;
    
    BULK INSERT bronze.Products
    FROM '/data/000.Hypermarket Products.csv'
    WITH (
        FIELDTERMINATOR = ',',
        ROWTERMINATOR = '\n',
        FIRSTROW = 2,
        ERRORFILE = '/logs/products_errors.txt'
    );
    
    -- Add audit metadata
    UPDATE bronze.Products
    SET LoadTimestamp = GETDATE(),
        SourceFile = '000.Hypermarket Products.csv';
END;
目标:导入原始数据,仅添加审计字段,最小化转换操作。
sql
-- 示例:青铜层产品表结构
CREATE TABLE bronze.Products (
    ProductID INT,
    ProductName NVARCHAR(255),
    Category NVARCHAR(100),
    SubCategory NVARCHAR(100),
    UnitPrice DECIMAL(10,2),
    SupplierID INT,
    RecipeYield DECIMAL(5,2),  -- 用于肉禽转换
    LoadTimestamp DATETIME2 DEFAULT GETDATE(),
    SourceFile NVARCHAR(500)
);

-- 青铜层加载模式
CREATE PROCEDURE bronze.usp_LoadProducts
AS
BEGIN
    TRUNCATE TABLE bronze.Products;
    
    BULK INSERT bronze.Products
    FROM '/data/000.Hypermarket Products.csv'
    WITH (
        FIELDTERMINATOR = ',',
        ROWTERMINATOR = '\n',
        FIRSTROW = 2,
        ERRORFILE = '/logs/products_errors.txt'
    );
    
    -- 添加审计元数据
    UPDATE bronze.Products
    SET LoadTimestamp = GETDATE(),
        SourceFile = '000.Hypermarket Products.csv';
END;

Silver Layer (Cleaned & Standardized)

白银层(清洗与标准化)

Purpose: Apply data quality rules, deduplication, and business transformations.
sql
-- Example: Silver Sales with Business Rules
CREATE PROCEDURE silver.usp_TransformSales
AS
BEGIN
    TRUNCATE TABLE silver.Sales;
    
    INSERT INTO silver.Sales (
        SaleID,
        BranchID,
        ProductID,
        SaleDate,
        Quantity,
        UnitPrice,
        TotalAmount,
        AdjustedQuantity,  -- Recipe conversion applied
        DataQualityScore
    )
    SELECT 
        s.SaleID,
        s.BranchID,
        s.ProductID,
        CAST(s.SaleDate AS DATE) AS SaleDate,
        s.Quantity,
        s.UnitPrice,
        s.Quantity * s.UnitPrice AS TotalAmount,
        -- Apply recipe yield for meat/poultry
        CASE 
            WHEN p.Category = 'Meat & Poultry' 
            THEN s.Quantity * ISNULL(p.RecipeYield, 1.0)
            ELSE s.Quantity
        END AS AdjustedQuantity,
        -- Data quality scoring
        CASE 
            WHEN s.Quantity > 0 AND s.UnitPrice > 0 THEN 100
            WHEN s.Quantity IS NULL OR s.UnitPrice IS NULL THEN 0
            ELSE 50
        END AS DataQualityScore
    FROM bronze.Sales s
    INNER JOIN bronze.Products p ON s.ProductID = p.ProductID
    WHERE s.Quantity > 0  -- Filter invalid records
      AND s.SaleDate >= DATEADD(YEAR, -2, GETDATE());  -- Keep 2 years
END;
Key Silver Transformations:
  • Date standardization
  • Recipe yield conversions for perishables
  • Duplicate removal
  • Null handling and imputation
  • Data quality scoring
目标:应用数据质量规则、去重及业务转换。
sql
-- 示例:带业务规则的白银层销售表
CREATE PROCEDURE silver.usp_TransformSales
AS
BEGIN
    TRUNCATE TABLE silver.Sales;
    
    INSERT INTO silver.Sales (
        SaleID,
        BranchID,
        ProductID,
        SaleDate,
        Quantity,
        UnitPrice,
        TotalAmount,
        AdjustedQuantity,  -- 已应用配方转换
        DataQualityScore
    )
    SELECT 
        s.SaleID,
        s.BranchID,
        s.ProductID,
        CAST(s.SaleDate AS DATE) AS SaleDate,
        s.Quantity,
        s.UnitPrice,
        s.Quantity * s.UnitPrice AS TotalAmount,
        -- 对肉禽应用出成率转换
        CASE 
            WHEN p.Category = 'Meat & Poultry' 
            THEN s.Quantity * ISNULL(p.RecipeYield, 1.0)
            ELSE s.Quantity
        END AS AdjustedQuantity,
        -- 数据质量评分
        CASE 
            WHEN s.Quantity > 0 AND s.UnitPrice > 0 THEN 100
            WHEN s.Quantity IS NULL OR s.UnitPrice IS NULL THEN 0
            ELSE 50
        END AS DataQualityScore
    FROM bronze.Sales s
    INNER JOIN bronze.Products p ON s.ProductID = p.ProductID
    WHERE s.Quantity > 0  -- 过滤无效记录
      AND s.SaleDate >= DATEADD(YEAR, -2, GETDATE());  -- 保留最近2年数据
END;
白银层核心转换操作:
  • 日期标准化
  • 易腐品配方出成率转换
  • 去重
  • 空值处理与填充
  • 数据质量评分

Gold Layer (Analytics-Ready Aggregates)

黄金层(可分析聚合数据)

Purpose: Create dimensional models and pre-aggregated metrics for BI tools.
sql
-- Example: Gold Inventory Turnover Metrics
CREATE PROCEDURE gold.usp_BuildInventoryMetrics
AS
BEGIN
    TRUNCATE TABLE gold.InventoryTurnover;
    
    INSERT INTO gold.InventoryTurnover (
        ProductID,
        ProductName,
        Category,
        BranchID,
        Month,
        TotalSalesQty,
        AvgStockLevel,
        TurnoverRatio,
        ShrinkagePercent,
        ReorderAlert
    )
    SELECT 
        p.ProductID,
        p.ProductName,
        p.Category,
        s.BranchID,
        DATEPART(MONTH, s.SaleDate) AS Month,
        SUM(s.AdjustedQuantity) AS TotalSalesQty,
        AVG(st.StockQuantity) AS AvgStockLevel,
        -- Turnover = Sales / Avg Stock
        CASE 
            WHEN AVG(st.StockQuantity) > 0 
            THEN SUM(s.AdjustedQuantity) / AVG(st.StockQuantity)
            ELSE 0
        END AS TurnoverRatio,
        -- Shrinkage = (Expected - Actual) / Expected
        CASE 
            WHEN SUM(st.ExpectedStock) > 0
            THEN ((SUM(st.ExpectedStock) - SUM(st.StockQuantity)) * 100.0) / SUM(st.ExpectedStock)
            ELSE 0
        END AS ShrinkagePercent,
        -- Alert if turnover < 2 (slow-moving inventory)
        CASE 
            WHEN SUM(s.AdjustedQuantity) / NULLIF(AVG(st.StockQuantity), 0) < 2 
            THEN 'Reorder Needed'
            ELSE 'OK'
        END AS ReorderAlert
    FROM silver.Sales s
    INNER JOIN silver.Products p ON s.ProductID = p.ProductID
    LEFT JOIN silver.Stock st ON s.ProductID = st.ProductID AND s.BranchID = st.BranchID
    GROUP BY p.ProductID, p.ProductName, p.Category, s.BranchID, DATEPART(MONTH, s.SaleDate);
END;
Gold Layer Tables:
  • InventoryTurnover
    : Stock efficiency metrics
  • SalesPerformance
    : Revenue aggregates by branch/category
  • SupplierRebates
    : Tiered rebate calculations
  • ProductMargins
    : Profit analysis dimensions
目标:创建维度模型和预聚合指标,供BI工具使用。
sql
-- 示例:黄金层库存周转率指标
CREATE PROCEDURE gold.usp_BuildInventoryMetrics
AS
BEGIN
    TRUNCATE TABLE gold.InventoryTurnover;
    
    INSERT INTO gold.InventoryTurnover (
        ProductID,
        ProductName,
        Category,
        BranchID,
        Month,
        TotalSalesQty,
        AvgStockLevel,
        TurnoverRatio,
        ShrinkagePercent,
        ReorderAlert
    )
    SELECT 
        p.ProductID,
        p.ProductName,
        p.Category,
        s.BranchID,
        DATEPART(MONTH, s.SaleDate) AS Month,
        SUM(s.AdjustedQuantity) AS TotalSalesQty,
        AVG(st.StockQuantity) AS AvgStockLevel,
        -- 周转率 = 销量 / 平均库存
        CASE 
            WHEN AVG(st.StockQuantity) > 0 
            THEN SUM(s.AdjustedQuantity) / AVG(st.StockQuantity)
            ELSE 0
        END AS TurnoverRatio,
        -- 损耗率 = (预期库存 - 实际库存) / 预期库存
        CASE 
            WHEN SUM(st.ExpectedStock) > 0
            THEN ((SUM(st.ExpectedStock) - SUM(st.StockQuantity)) * 100.0) / SUM(st.ExpectedStock)
            ELSE 0
        END AS ShrinkagePercent,
        -- 若周转率 < 2则触发补货提醒(滞销库存)
        CASE 
            WHEN SUM(s.AdjustedQuantity) / NULLIF(AVG(st.StockQuantity), 0) < 2 
            THEN '需要补货'
            ELSE '正常'
        END AS ReorderAlert
    FROM silver.Sales s
    INNER JOIN silver.Products p ON s.ProductID = p.ProductID
    LEFT JOIN silver.Stock st ON s.ProductID = st.ProductID AND s.BranchID = st.BranchID
    GROUP BY p.ProductID, p.ProductName, p.Category, s.BranchID, DATEPART(MONTH, s.SaleDate);
END;
黄金层表:
  • InventoryTurnover
    :库存效率指标
  • SalesPerformance
    :按门店/品类聚合的收入数据
  • SupplierRebates
    :层级返利计算
  • ProductMargins
    :利润分析维度

Complete Pipeline Execution

完整管道执行

Manual Execution (Sequential)

手动执行(顺序执行)

sql
-- 1. Bronze: Load raw data
EXEC bronze.usp_LoadProducts;
EXEC bronze.usp_LoadSales;
EXEC bronze.usp_LoadStock;

-- 2. Silver: Apply transformations
EXEC silver.usp_TransformProducts;
EXEC silver.usp_TransformSales;
EXEC silver.usp_TransformStock;

-- 3. Gold: Build analytics aggregates
EXEC gold.usp_BuildInventoryMetrics;
EXEC gold.usp_BuildSalesPerformance;
EXEC gold.usp_BuildSupplierRebates;

-- 4. Verify row counts
SELECT 'Bronze Products' AS Layer, COUNT(*) AS RowCount FROM bronze.Products
UNION ALL
SELECT 'Silver Products', COUNT(*) FROM silver.Products
UNION ALL
SELECT 'Gold Inventory', COUNT(*) FROM gold.InventoryTurnover;
sql
-- 1. 青铜层:加载原始数据
EXEC bronze.usp_LoadProducts;
EXEC bronze.usp_LoadSales;
EXEC bronze.usp_LoadStock;

-- 2. 白银层:应用转换规则
EXEC silver.usp_TransformProducts;
EXEC silver.usp_TransformSales;
EXEC silver.usp_TransformStock;

-- 3. 黄金层:构建分析聚合数据
EXEC gold.usp_BuildInventoryMetrics;
EXEC gold.usp_BuildSalesPerformance;
EXEC gold.usp_BuildSupplierRebates;

-- 4. 验证行数
SELECT '青铜层产品' AS 层级, COUNT(*) AS 行数 FROM bronze.Products
UNION ALL
SELECT '白银层产品', COUNT(*) FROM silver.Products
UNION ALL
SELECT '黄金层库存', COUNT(*) FROM gold.InventoryTurnover;

Automated Pipeline Script

自动化管道脚本

bash
#!/bin/bash
bash
#!/bin/bash

run_etl_pipeline.sh

run_etl_pipeline.sh

set -e
SQL_SERVER="${SQL_SERVER:-localhost}" SQL_USER="${SQL_USER:-sa}" SQL_PASSWORD="${SQL_PASSWORD}"
echo "Starting Retail ETL Pipeline..."
set -e
SQL_SERVER="${SQL_SERVER:-localhost}" SQL_USER="${SQL_USER:-sa}" SQL_PASSWORD="${SQL_PASSWORD}"
echo "启动零售ETL管道..."

Bronze Layer

青铜层

echo "[1/3] Loading Bronze Layer..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadProducts;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadSales;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadStock;"
echo "[1/3] 加载青铜层..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadProducts;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadSales;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC bronze.usp_LoadStock;"

Silver Layer

白银层

echo "[2/3] Transforming Silver Layer..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformProducts;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformSales;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformStock;"
echo "[2/3] 转换白银层..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformProducts;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformSales;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC silver.usp_TransformStock;"

Gold Layer

黄金层

echo "[3/3] Building Gold Layer..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC gold.usp_BuildInventoryMetrics;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC gold.usp_BuildSalesPerformance;"
echo "Pipeline completed successfully!"
undefined
echo "[3/3] 构建黄金层..." sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC gold.usp_BuildInventoryMetrics;" sqlcmd -S "$SQL_SERVER" -U "$SQL_USER" -P "$SQL_PASSWORD" -d RetailDataWarehouse -Q "EXEC gold.usp_BuildSalesPerformance;"
echo "管道执行成功!"
undefined

Configuration

配置

Environment Variables

环境变量

bash
undefined
bash
undefined

.env file for pipeline configuration

管道配置.env文件

SQL_SERVER=localhost SQL_USER=sa SQL_PASSWORD=${SQL_SA_PASSWORD} SQL_DATABASE=RetailDataWarehouse
SQL_SERVER=localhost SQL_USER=sa SQL_PASSWORD=${SQL_SA_PASSWORD} SQL_DATABASE=RetailDataWarehouse

Data source paths

数据源路径

DATA_SOURCE_PATH=/path/to/data_source LOGS_PATH=/var/log/retail-etl
DATA_SOURCE_PATH=/path/to/data_source LOGS_PATH=/var/log/retail-etl

Airflow (if using orchestration)

Airflow(若使用编排工具)

AIRFLOW_HOME=/opt/airflow AIRFLOW__CORE__DAGS_FOLDER=${AIRFLOW_HOME}/dags
undefined
AIRFLOW_HOME=/opt/airflow AIRFLOW__CORE__DAGS_FOLDER=${AIRFLOW_HOME}/dags
undefined

Docker Compose Configuration

Docker Compose配置

yaml
version: '3.8'
services:
  sqlserver:
    image: mcr.microsoft.com/mssql/server:2019-latest
    environment:
      ACCEPT_EULA: Y
      SA_PASSWORD: ${SQL_SA_PASSWORD}
      MSSQL_PID: Developer
    ports:
      - "1433:1433"
    volumes:
      - ./data_source:/data
      - ./sql_scripts:/scripts
      - sqlserver_data:/var/opt/mssql
    restart: unless-stopped

volumes:
  sqlserver_data:
yaml
version: '3.8'
services:
  sqlserver:
    image: mcr.microsoft.com/mssql/server:2019-latest
    environment:
      ACCEPT_EULA: Y
      SA_PASSWORD: ${SQL_SA_PASSWORD}
      MSSQL_PID: Developer
    ports:
      - "1433:1433"
    volumes:
      - ./data_source:/data
      - ./sql_scripts:/scripts
      - sqlserver_data:/var/opt/mssql
    restart: unless-stopped

volumes:
  sqlserver_data:

Business Logic Examples

业务逻辑示例

Recipe Conversion for Meat Products

肉类产品配方转换

sql
-- Handle meat/poultry yield conversions
-- Example: 1kg raw chicken → 0.65kg cooked meat
CREATE FUNCTION dbo.fn_ApplyRecipeYield(
    @Quantity DECIMAL(10,2),
    @RecipeYield DECIMAL(5,2),
    @Category NVARCHAR(100)
)
RETURNS DECIMAL(10,2)
AS
BEGIN
    DECLARE @AdjustedQty DECIMAL(10,2);
    
    IF @Category IN ('Meat & Poultry', 'Seafood')
        SET @AdjustedQty = @Quantity * ISNULL(@RecipeYield, 1.0);
    ELSE
        SET @AdjustedQty = @Quantity;
    
    RETURN @AdjustedQty;
END;
sql
-- 处理肉禽出成率转换
-- 示例:1kg生鸡肉 → 0.65kg熟肉
CREATE FUNCTION dbo.fn_ApplyRecipeYield(
    @Quantity DECIMAL(10,2),
    @RecipeYield DECIMAL(5,2),
    @Category NVARCHAR(100)
)
RETURNS DECIMAL(10,2)
AS
BEGIN
    DECLARE @AdjustedQty DECIMAL(10,2);
    
    IF @Category IN ('Meat & Poultry', 'Seafood')
        SET @AdjustedQty = @Quantity * ISNULL(@RecipeYield, 1.0);
    ELSE
        SET @AdjustedQty = @Quantity;
    
    RETURN @AdjustedQty;
END;

Supplier Rebate Tiers

供应商返利层级

sql
-- Calculate dynamic rebate percentages based on purchase volume
CREATE PROCEDURE gold.usp_CalculateSupplierRebates
AS
BEGIN
    INSERT INTO gold.SupplierRebates (
        SupplierID,
        TotalPurchaseAmount,
        RebateTier,
        RebatePercent,
        RebateAmount
    )
    SELECT 
        SupplierID,
        SUM(TotalAmount) AS TotalPurchaseAmount,
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 'Platinum'
            WHEN SUM(TotalAmount) >= 50000 THEN 'Gold'
            WHEN SUM(TotalAmount) >= 25000 THEN 'Silver'
            ELSE 'Bronze'
        END AS RebateTier,
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 5.0
            WHEN SUM(TotalAmount) >= 50000 THEN 3.0
            WHEN SUM(TotalAmount) >= 25000 THEN 1.5
            ELSE 0.0
        END AS RebatePercent,
        SUM(TotalAmount) * 
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 0.05
            WHEN SUM(TotalAmount) >= 50000 THEN 0.03
            WHEN SUM(TotalAmount) >= 25000 THEN 0.015
            ELSE 0.0
        END AS RebateAmount
    FROM silver.Sales s
    INNER JOIN silver.Products p ON s.ProductID = p.ProductID
    GROUP BY SupplierID;
END;
sql
-- 根据采购量计算动态返利比例
CREATE PROCEDURE gold.usp_CalculateSupplierRebates
AS
BEGIN
    INSERT INTO gold.SupplierRebates (
        SupplierID,
        TotalPurchaseAmount,
        RebateTier,
        RebatePercent,
        RebateAmount
    )
    SELECT 
        SupplierID,
        SUM(TotalAmount) AS TotalPurchaseAmount,
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 'Platinum'
            WHEN SUM(TotalAmount) >= 50000 THEN 'Gold'
            WHEN SUM(TotalAmount) >= 25000 THEN 'Silver'
            ELSE 'Bronze'
        END AS RebateTier,
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 5.0
            WHEN SUM(TotalAmount) >= 50000 THEN 3.0
            WHEN SUM(TotalAmount) >= 25000 THEN 1.5
            ELSE 0.0
        END AS RebatePercent,
        SUM(TotalAmount) * 
        CASE 
            WHEN SUM(TotalAmount) >= 100000 THEN 0.05
            WHEN SUM(TotalAmount) >= 50000 THEN 0.03
            WHEN SUM(TotalAmount) >= 25000 THEN 0.015
            ELSE 0.0
        END AS RebateAmount
    FROM silver.Sales s
    INNER JOIN silver.Products p ON s.ProductID = p.ProductID
    GROUP BY SupplierID;
END;

Inventory Shrinkage Detection

库存损耗检测

sql
-- Identify products with abnormal shrinkage
SELECT 
    p.ProductName,
    p.Category,
    st.BranchID,
    st.ExpectedStock,
    st.StockQuantity AS ActualStock,
    ((st.ExpectedStock - st.StockQuantity) * 100.0) / st.ExpectedStock AS ShrinkagePercent
FROM silver.Stock st
INNER JOIN silver.Products p ON st.ProductID = p.ProductID
WHERE st.ExpectedStock > 0
  AND ((st.ExpectedStock - st.StockQuantity) * 100.0) / st.ExpectedStock > 5.0  -- >5% shrinkage threshold
ORDER BY ShrinkagePercent DESC;
sql
-- 识别异常损耗的产品
SELECT 
    p.ProductName,
    p.Category,
    st.BranchID,
    st.ExpectedStock,
    st.StockQuantity AS ActualStock,
    ((st.ExpectedStock - st.StockQuantity) * 100.0) / st.ExpectedStock AS ShrinkagePercent
FROM silver.Stock st
INNER JOIN silver.Products p ON st.ProductID = p.ProductID
WHERE st.ExpectedStock > 0
  AND ((st.ExpectedStock - st.StockQuantity) * 100.0) / st.ExpectedStock > 5.0  -- 损耗阈值>5%
ORDER BY ShrinkagePercent DESC;

Data Quality Checks

数据质量检查

Validation Queries

验证查询

sql
-- Check for duplicate sales records
SELECT SaleID, COUNT(*) AS Duplicates
FROM bronze.Sales
GROUP BY SaleID
HAVING COUNT(*) > 1;

-- Validate price consistency
SELECT 
    p.ProductID,
    p.ProductName,
    COUNT(DISTINCT s.UnitPrice) AS PriceVariations
FROM silver.Products p
INNER JOIN silver.Sales s ON p.ProductID = s.ProductID
GROUP BY p.ProductID, p.ProductName
HAVING COUNT(DISTINCT s.UnitPrice) > 3;  -- More than 3 price points

-- Check for negative stock
SELECT ProductID, BranchID, StockQuantity
FROM silver.Stock
WHERE StockQuantity < 0;

-- Data completeness metrics
SELECT 
    'Products' AS TableName,
    COUNT(*) AS TotalRows,
    SUM(CASE WHEN ProductName IS NULL THEN 1 ELSE 0 END) AS NullProductNames,
    SUM(CASE WHEN UnitPrice IS NULL THEN 1 ELSE 0 END) AS NullPrices
FROM silver.Products;
sql
-- 检查重复销售记录
SELECT SaleID, COUNT(*) AS Duplicates
FROM bronze.Sales
GROUP BY SaleID
HAVING COUNT(*) > 1;

-- 验证价格一致性
SELECT 
    p.ProductID,
    p.ProductName,
    COUNT(DISTINCT s.UnitPrice) AS PriceVariations
FROM silver.Products p
INNER JOIN silver.Sales s ON p.ProductID = s.ProductID
GROUP BY p.ProductID, p.ProductName
HAVING COUNT(DISTINCT s.UnitPrice) > 3;  -- 价格变动超过3种

-- 检查负库存
SELECT ProductID, BranchID, StockQuantity
FROM silver.Stock
WHERE StockQuantity < 0;

-- 数据完整性指标
SELECT 
    '产品表' AS 表名,
    COUNT(*) AS 总行数,
    SUM(CASE WHEN ProductName IS NULL THEN 1 ELSE 0 END) AS 产品名称空值数,
    SUM(CASE WHEN UnitPrice IS NULL THEN 1 ELSE 0 END) AS 价格空值数
FROM silver.Products;

Troubleshooting

故障排查

Common Issues

常见问题

Issue: BULK INSERT fails with permission error
sql
-- Solution: Grant read permissions to SQL Server service account
-- Or use OPENROWSET with explicit credentials
INSERT INTO bronze.Products
SELECT * FROM OPENROWSET(
    BULK '/data/000.Hypermarket Products.csv',
    FORMATFILE = '/data/products_format.xml',
    ERRORFILE = '/logs/errors.txt'
) AS DataFile;
Issue: Recipe yield conversions producing NULL values
sql
-- Check for missing RecipeYield in Products table
SELECT ProductID, ProductName, Category, RecipeYield
FROM bronze.Products
WHERE Category IN ('Meat & Poultry', 'Seafood')
  AND RecipeYield IS NULL;

-- Fix: Set default yield to 1.0
UPDATE bronze.Products
SET RecipeYield = 1.0
WHERE RecipeYield IS NULL;
Issue: Silver layer procedure times out on large datasets
sql
-- Solution: Add batch processing with cursor or temp tables
CREATE PROCEDURE silver.usp_TransformSalesBatch
    @BatchSize INT = 10000
AS
BEGIN
    DECLARE @MinID INT, @MaxID INT;
    
    SELECT @MinID = MIN(SaleID), @MaxID = MAX(SaleID) FROM bronze.Sales;
    
    WHILE @MinID <= @MaxID
    BEGIN
        INSERT INTO silver.Sales (...)
        SELECT ...
        FROM bronze.Sales
        WHERE SaleID BETWEEN @MinID AND (@MinID + @BatchSize - 1);
        
        SET @MinID = @MinID + @BatchSize;
    END;
END;
Issue: Gold aggregates not updating incrementally
sql
-- Solution: Implement incremental load with watermark
CREATE TABLE gold.ETL_Watermark (
    TableName NVARCHAR(100),
    LastProcessedDate DATETIME2
);

CREATE PROCEDURE gold.usp_IncrementalInventoryMetrics
AS
BEGIN
    DECLARE @LastRun DATETIME2;
    SELECT @LastRun = LastProcessedDate FROM gold.ETL_Watermark WHERE TableName = 'InventoryMetrics';
    
    -- Delete and recalculate only changed data
    DELETE FROM gold.InventoryTurnover
    WHERE Month >= DATEPART(MONTH, @LastRun);
    
    INSERT INTO gold.InventoryTurnover (...)
    SELECT ...
    FROM silver.Sales
    WHERE SaleDate >= @LastRun;
    
    -- Update watermark
    UPDATE gold.ETL_Watermark 
    SET LastProcessedDate = GETDATE()
    WHERE TableName = 'InventoryMetrics';
END;
问题:BULK INSERT因权限错误失败
sql
-- 解决方案:为SQL Server服务账号授予读取权限
-- 或使用带明确凭据的OPENROWSET
INSERT INTO bronze.Products
SELECT * FROM OPENROWSET(
    BULK '/data/000.Hypermarket Products.csv',
    FORMATFILE = '/data/products_format.xml',
    ERRORFILE = '/logs/errors.txt'
) AS DataFile;
问题:配方出成率转换产生NULL值
sql
-- 检查产品表中缺失的RecipeYield字段
SELECT ProductID, ProductName, Category, RecipeYield
FROM bronze.Products
WHERE Category IN ('Meat & Poultry', 'Seafood')
  AND RecipeYield IS NULL;

-- 修复:将默认出成率设置为1.0
UPDATE bronze.Products
SET RecipeYield = 1.0
WHERE RecipeYield IS NULL;
问题:白银层存储过程处理大数据集时超时
sql
-- 解决方案:使用游标或临时表实现批量处理
CREATE PROCEDURE silver.usp_TransformSalesBatch
    @BatchSize INT = 10000
AS
BEGIN
    DECLARE @MinID INT, @MaxID INT;
    
    SELECT @MinID = MIN(SaleID), @MaxID = MAX(SaleID) FROM bronze.Sales;
    
    WHILE @MinID <= @MaxID
    BEGIN
        INSERT INTO silver.Sales (...)
        SELECT ...
        FROM bronze.Sales
        WHERE SaleID BETWEEN @MinID AND (@MinID + @BatchSize - 1);
        
        SET @MinID = @MinID + @BatchSize;
    END;
END;
问题:黄金层聚合数据未增量更新
sql
-- 解决方案:使用水印实现增量加载
CREATE TABLE gold.ETL_Watermark (
    TableName NVARCHAR(100),
    LastProcessedDate DATETIME2
);

CREATE PROCEDURE gold.usp_IncrementalInventoryMetrics
AS
BEGIN
    DECLARE @LastRun DATETIME2;
    SELECT @LastRun = LastProcessedDate FROM gold.ETL_Watermark WHERE TableName = 'InventoryMetrics';
    
    -- 删除并重新计算变更数据
    DELETE FROM gold.InventoryTurnover
    WHERE Month >= DATEPART(MONTH, @LastRun);
    
    INSERT INTO gold.InventoryTurnover (...)
    SELECT ...
    FROM silver.Sales
    WHERE SaleDate >= @LastRun;
    
    -- 更新水印
    UPDATE gold.ETL_Watermark 
    SET LastProcessedDate = GETDATE()
    WHERE TableName = 'InventoryMetrics';
END;

Performance Optimization

性能优化

sql
-- Add indexes for Bronze layer queries
CREATE CLUSTERED INDEX IX_Sales_SaleID ON bronze.Sales(SaleID);
CREATE NONCLUSTERED INDEX IX_Sales_ProductID ON bronze.Sales(ProductID);
CREATE NONCLUSTERED INDEX IX_Sales_SaleDate ON bronze.Sales(SaleDate);

-- Partition Gold tables by month for faster queries
CREATE PARTITION FUNCTION pf_MonthPartition (INT)
AS RANGE RIGHT FOR VALUES (1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12);

CREATE PARTITION SCHEME ps_MonthPartition
AS PARTITION pf_MonthPartition ALL TO ([PRIMARY]);

CREATE TABLE gold.InventoryTurnover (
    ...
    Month INT
) ON ps_MonthPartition(Month);

-- Enable query store for performance monitoring
ALTER DATABASE RetailDataWarehouse SET QUERY_STORE = ON;
sql
-- 为青铜层查询添加索引
CREATE CLUSTERED INDEX IX_Sales_SaleID ON bronze.Sales(SaleID);
CREATE NONCLUSTERED INDEX IX_Sales_ProductID ON bronze.Sales(ProductID);
CREATE NONCLUSTERED INDEX IX_Sales_SaleDate ON bronze.Sales(SaleDate);

-- 按月份分区黄金层表以提升查询速度
CREATE PARTITION FUNCTION pf_MonthPartition (INT)
AS RANGE RIGHT FOR VALUES (1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12);

CREATE PARTITION SCHEME ps_MonthPartition
AS PARTITION pf_MonthPartition ALL TO ([PRIMARY]);

CREATE TABLE gold.InventoryTurnover (
    ...
    Month INT
) ON ps_MonthPartition(Month);

-- 启用查询存储以监控性能
ALTER DATABASE RetailDataWarehouse SET QUERY_STORE = ON;

Integration with BI Tools

与BI工具集成

Power BI Connection

Power BI连接

sql
-- Create view optimized for Power BI
CREATE VIEW gold.vw_SalesDashboard AS
SELECT 
    s.SaleDate,
    p.ProductName,
    p.Category,
    b.BranchName,
    s.Quantity,
    s.UnitPrice,
    s.TotalAmount,
    i.TurnoverRatio,
    i.ShrinkagePercent
FROM gold.InventoryTurnover i
INNER JOIN silver.Sales s ON i.ProductID = s.ProductID AND i.BranchID = s.BranchID
INNER JOIN silver.Products p ON s.ProductID = p.ProductID
INNER JOIN silver.Branches b ON s.BranchID = b.BranchID;

-- Grant read-only access to BI service account
CREATE USER [bi_service] WITH PASSWORD = '${BI_SERVICE_PASSWORD}';
GRANT SELECT ON SCHEMA::gold TO [bi_service];
sql
-- 创建为Power BI优化的视图
CREATE VIEW gold.vw_SalesDashboard AS
SELECT 
    s.SaleDate,
    p.ProductName,
    p.Category,
    b.BranchName,
    s.Quantity,
    s.UnitPrice,
    s.TotalAmount,
    i.TurnoverRatio,
    i.ShrinkagePercent
FROM gold.InventoryTurnover i
INNER JOIN silver.Sales s ON i.ProductID = s.ProductID AND i.BranchID = s.BranchID
INNER JOIN silver.Products p ON s.ProductID = p.ProductID
INNER JOIN silver.Branches b ON s.BranchID = b.BranchID;

-- 为BI服务账号授予只读权限
CREATE USER [bi_service] WITH PASSWORD = '${BI_SERVICE_PASSWORD}';
GRANT SELECT ON SCHEMA::gold TO [bi_service];

Monitoring & Logging

监控与日志

sql
-- Create audit log table
CREATE TABLE dbo.ETL_AuditLog (
    LogID INT IDENTITY(1,1) PRIMARY KEY,
    ProcedureName NVARCHAR(255),
    LayerName NVARCHAR(50),
    StartTime DATETIME2,
    EndTime DATETIME2,
    RowsProcessed INT,
    Status NVARCHAR(50),
    ErrorMessage NVARCHAR(MAX)
);

-- Example audit logging in procedures
CREATE PROCEDURE silver.usp_TransformSalesWithLogging
AS
BEGIN
    DECLARE @StartTime DATETIME2 = GETDATE();
    DECLARE @RowCount INT;
    
    BEGIN TRY
        -- Transform logic
        INSERT INTO silver.Sales (...) SELECT ...;
        SET @RowCount = @@ROWCOUNT;
        
        -- Log success
        INSERT INTO dbo.ETL_AuditLog (ProcedureName, LayerName, StartTime, EndTime, RowsProcessed, Status)
        VALUES ('usp_TransformSales', 'Silver', @StartTime, GETDATE(), @RowCount, 'Success');
    END TRY
    BEGIN CATCH
        -- Log failure
        INSERT INTO dbo.ETL_AuditLog (ProcedureName, LayerName, StartTime, EndTime, Status, ErrorMessage)
        VALUES ('usp_TransformSales', 'Silver', @StartTime, GETDATE(), 'Failed', ERROR_MESSAGE());
        
        THROW;
    END CATCH;
END;
This skill provides comprehensive guidance for implementing and extending the Retail ETL Medallion Pipeline with real-world business logic and production-ready patterns.
sql
-- 创建审计日志表
CREATE TABLE dbo.ETL_AuditLog (
    LogID INT IDENTITY(1,1) PRIMARY KEY,
    ProcedureName NVARCHAR(255),
    LayerName NVARCHAR(50),
    StartTime DATETIME2,
    EndTime DATETIME2,
    RowsProcessed INT,
    Status NVARCHAR(50),
    ErrorMessage NVARCHAR(MAX)
);

-- 存储过程中添加审计日志示例
CREATE PROCEDURE silver.usp_TransformSalesWithLogging
AS
BEGIN
    DECLARE @StartTime DATETIME2 = GETDATE();
    DECLARE @RowCount INT;
    
    BEGIN TRY
        -- 转换逻辑
        INSERT INTO silver.Sales (...) SELECT ...;
        SET @RowCount = @@ROWCOUNT;
        
        -- 记录成功日志
        INSERT INTO dbo.ETL_AuditLog (ProcedureName, LayerName, StartTime, EndTime, RowsProcessed, Status)
        VALUES ('usp_TransformSales', 'Silver', @StartTime, GETDATE(), @RowCount, 'Success');
    END TRY
    BEGIN CATCH
        -- 记录失败日志
        INSERT INTO dbo.ETL_AuditLog (ProcedureName, LayerName, StartTime, EndTime, Status, ErrorMessage)
        VALUES ('usp_TransformSales', 'Silver', @StartTime, GETDATE(), 'Failed', ERROR_MESSAGE());
        
        THROW;
    END CATCH;
END;
本技能提供了全面的指导,帮助实现和扩展零售ETL Medallion管道,包含真实业务逻辑和生产级模式。