retail-etl-medallion-pipeline
Compare original and translation side by side
🇺🇸
Original
English🇨🇳
Translation
ChineseRetail ETL Medallion Pipeline Skill
零售ETL Medallion管道技能
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 setupRetail-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
undefinedStart 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
undefinedbash
undefinedConnect 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"
doneKey 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:
- : Stock efficiency metrics
InventoryTurnover - : Revenue aggregates by branch/category
SalesPerformance - : Tiered rebate calculations
SupplierRebates - : Profit analysis dimensions
ProductMargins
目标:创建维度模型和预聚合指标,供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/bashbash
#!/bin/bashrun_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!"
undefinedecho "[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 "管道执行成功!"
undefinedConfiguration
配置
Environment Variables
环境变量
bash
undefinedbash
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
undefinedAIRFLOW_HOME=/opt/airflow
AIRFLOW__CORE__DAGS_FOLDER=${AIRFLOW_HOME}/dags
undefinedDocker 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管道,包含真实业务逻辑和生产级模式。