-- Database
CREATE OR REPLACE DATABASE JULY_PROJECT;
-- Schemas
CREATE OR REPLACE SCHEMA RAW_LAYER;
CREATE OR REPLACE SCHEMA STAGE_LAYER;
CREATE OR REPLACE SCHEMA STANDARD_LAYER;
--------------------------------------------------------- RAW LAYER
-------------------------------------------------------------------
-- Stoarge Integration
create or replace storage integration S3_INTEGRATION_JULY
TYPE = EXTERNAL_STAGE
STORAGE_PROVIDER = S3
ENABLED = TRUE
STORAGE_AWS_ROLE_ARN = 'arn:aws:iam::394094767033:role/snowflake_training_july'
STORAGE_ALLOWED_LOCATIONS = ('s3://snowflake-training-july/json/')
COMMENT = 'Integration with aws s3 buckets';
-- To get the metadata of the storage integration
DESC integration S3_INTEGRATION_JULY_PROJECT;
-- File Format
CREATE OR REPLACE FILE FORMAT JULY_PROJECT.RAW_LAYER.json_format
TYPE = 'JSON';
-- External Stage
CREATE OR REPLACE STAGE JULY_PROJECT.RAW_LAYER.aws_s3_json
URL = 's3://snowflake-training-july/json/'
STORAGE_INTEGRATION = S3_INTEGRATION_JULY_PROJECT
FILE_FORMAT = JULY_PROJECT.RAW_LAYER.json_format;
-- List of the files from the external stage
LIST @JULY_PROJECT.RAW_LAYER.aws_s3_json;
SELECT * FROM @JULY_PROJECT.RAW_LAYER.aws_s3_json;
-- Raw orders data
CREATE OR REPLACE TABLE JULY_PROJECT.RAW_LAYER.orders_raw (
raw_json VARIANT,
-- row_inserted_date default current_timestamp
);
TRUNCATE TABLE JULY_PROJECT.RAW_LAYER.orders_raw;
SELECT * FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- COPY COMMAND --
COPY INTO JULY_PROJECT.RAW_LAYER.orders_raw
FROM @JULY_PROJECT.RAW_LAYER.aws_s3_json;
-- SNOWPIPE --
CREATE OR REPLACE pipe JULY_PROJECT.RAW_LAYER.customer_pipe
AUTO_INGEST = TRUE
AS
COPY INTO JULY_PROJECT.RAW_LAYER.orders_raw
FROM @JULY_PROJECT.RAW_LAYER.aws_s3_json;
-- To get the metadata of the pipe
desc pipe JULY_PROJECT.RAW_LAYER.customer_pipe;
-- File Name --
Orders_Data_July_Project.ndjson;
-- History --
SELECT * FROM SNOWFLAKE.ACCOUNT_USAGE.COPY_HISTORY;
--------------------------------------------------------- STAGE LAYER
----------------------------------------------------------------------------
-- Orders
CREATE OR REPLACE VIEW JULY_PROJECT.STAGE_LAYER.orders
AS
SELECT
raw_json:"order_id"::STRING AS order_id,
raw_json:"order_date"::TIMESTAMP_NTZ AS order_date,
raw_json:"customer":"customer_id"::STRING AS customer_id,
raw_json:"payment":"transaction_id"::STRING AS payment_id
FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- Customers
CREATE OR REPLACE VIEW JULY_PROJECT.STAGE_LAYER.customers AS
SELECT
raw_json:"customer":"customer_id"::STRING AS customer_id,
raw_json:"customer":"name"::STRING AS name,
raw_json:"customer":"email"::STRING AS email,
raw_json:"customer":"phone"::STRING AS phone,
raw_json:"customer":"address":"street"::STRING AS street,
raw_json:"customer":"address":"city"::STRING AS city,
raw_json:"customer":"address":"state"::STRING AS state,
raw_json:"customer":"address":"zip"::STRING AS zip,
raw_json:"customer":"address":"country"::STRING AS country
FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- Order Items (flatten array)
CREATE OR REPLACE VIEW JULY_PROJECT.STAGE_LAYER.order_items AS
SELECT
raw_json:"order_id"::STRING AS order_id,
[Link]:"product_id"::STRING AS product_id,
[Link]:"product_name"::STRING AS product_name,
[Link]:"category"::STRING AS category,
[Link]:"price"::FLOAT AS price,
[Link]:"quantity"::INT AS quantity
FROM JULY_PROJECT.RAW_LAYER.orders_raw,
LATERAL FLATTEN(input => raw_json:"items") AS item;
-- Payments
CREATE OR REPLACE VIEW JULY_PROJECT.STAGE_LAYER.payments AS
SELECT
raw_json:"order_id"::STRING AS order_id,
raw_json:"payment":"method"::STRING AS method,
raw_json:"payment":"transaction_id"::STRING AS transaction_id,
raw_json:"payment":"amount"::FLOAT AS amount,
raw_json:"payment":"currency"::STRING AS currency
FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- Delivery
CREATE OR REPLACE VIEW JULY_PROJECT.STAGE_LAYER.delivery AS
SELECT
raw_json:"order_id"::STRING AS order_id,
raw_json:"delivery":"status"::STRING AS status,
raw_json:"delivery":"expected_date"::DATE AS expected_date,
raw_json:"delivery":"tracking_number"::STRING AS tracking_number
FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- RAW TABLE --
SELECT * FROM JULY_PROJECT.RAW_LAYER.orders_raw;
-- truncate table JULY_PROJECT.RAW_LAYER.orders_raw;
-- STAGE VIEWS --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.orders;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.customers;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.order_items;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.payments;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.delivery;
------------------------------------------------------------ STANDARD LAYER
-------------------------------------------------------------------------
-- Facts --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.orders;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.order_items;
-- Dimensions --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.customers;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.delivery;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.payments;
-- Append Only --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.orders;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.order_items;
-- SCD Type 1 --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.payments;
-- SCD Type 2 --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.delivery;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.customers;
-- FACT TABLES --
CREATE OR REPLACE TABLE JULY_PROJECT.STANDARD_LAYER.fact_orders
LIKE JULY_PROJECT.STAGE_LAYER.orders;
CREATE OR REPLACE TABLE JULY_PROJECT.STANDARD_LAYER.fact_order_items
LIKE JULY_PROJECT.STAGE_LAYER.order_items;
-- DIMENSION TABLES --
CREATE OR REPLACE TABLE JULY_PROJECT.STANDARD_LAYER.dim_payments
LIKE JULY_PROJECT.STAGE_LAYER.payments;
CREATE OR REPLACE TABLE JULY_PROJECT.STANDARD_LAYER.dim_customers (
customer_sk INT AUTOINCREMENT,
customer_id STRING,
name STRING,
email STRING,
phone STRING,
street STRING,
city STRING,
state STRING,
zip STRING,
country STRING,
effective_start_date DATE,
effective_end_date DATE,
current_flag BOOLEAN
);
CREATE OR REPLACE TABLE JULY_PROJECT.STANDARD_LAYER.dim_delivery (
delivery_sk INT AUTOINCREMENT,
order_id STRING,
status STRING,
expected_date DATE,
tracking_number STRING,
effective_start_date DATE,
effective_end_date DATE,
current_flag BOOLEAN
);
-- STANDARD TABLES --
SELECT * FROM JULY_PROJECT.STANDARD_LAYER.fact_orders;
SELECT * FROM JULY_PROJECT.STANDARD_LAYER.fact_order_items;
SELECT * FROM JULY_PROJECT.STANDARD_LAYER.dim_payments;
SELECT * FROM JULY_PROJECT.STANDARD_LAYER.dim_customers;
SELECT * FROM JULY_PROJECT.STANDARD_LAYER.dim_delivery;
-- STAGE VIEWS --
SELECT * FROM JULY_PROJECT.STAGE_LAYER.orders;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.order_items;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.payments;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.customers;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.delivery;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.dim_payments;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.fact_orders;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.fact_order_items;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.dim_payments;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.dim_customers;
TRUNCATE TABLE JULY_PROJECT.STANDARD_LAYER.dim_delivery;
-- COLUMN NAMES AND DATA TYPES --
DESC TABLE JULY_PROJECT.STANDARD_LAYER.fact_orders;
DESC TABLE JULY_PROJECT.STANDARD_LAYER.dim_customers;
--------------------------------------------------------- FACT TABLE IMPLEMENTATION
----------------------------------------------------------------------
-- fact_orders --
INSERT INTO JULY_PROJECT.STANDARD_LAYER.fact_orders
SELECT SRC.ORDER_ID, SRC.ORDER_DATE, SRC.CUSTOMER_ID, SRC.PAYMENT_ID
FROM JULY_PROJECT.STAGE_LAYER.orders AS SRC
LEFT JOIN JULY_PROJECT.STANDARD_LAYER.fact_orders AS TGT
ON SRC.ORDER_ID = TGT.ORDER_ID
WHERE TGT.ORDER_ID IS NULL;
SELECT * FROM JULY_PROJECT.STAGE_LAYER.orders
WHERE ORDER_ID NOT IN (
SELECT ORDER_ID FROM JULY_PROJECT.STANDARD_LAYER.fact_orders
);
-- fact_order_items --
INSERT INTO JULY_PROJECT.STANDARD_LAYER.fact_order_items
SELECT SRC.ORDER_ID, SRC.PRODUCT_ID, SRC.PRODUCT_NAME, [Link],
[Link], [Link]
FROM JULY_PROJECT.STAGE_LAYER.order_items AS SRC
LEFT JOIN JULY_PROJECT.STANDARD_LAYER.fact_order_items AS TGT
ON SRC.ORDER_ID = TGT.ORDER_ID
WHERE TGT.ORDER_ID IS NULL;
---------------------------- DIMENSION SCD TYPE 1 TABLE IMPLEMENTATION
----------------------------------------------------
-- dim_payments --
MERGE INTO JULY_PROJECT.STANDARD_LAYER.dim_payments TGT
USING JULY_PROJECT.STAGE_LAYER.payments SRC
ON TGT.TRANSACTION_ID = SRC.TRANSACTION_ID
AND TGT.ORDER_ID = SRC.ORDER_ID
WHEN MATCHED
AND ([Link] <> [Link] OR [Link] <> [Link] OR
[Link] <> [Link])
THEN
UPDATE SET
[Link] = [Link],
[Link] = [Link],
[Link] = [Link]
WHEN NOT MATCHED THEN
INSERT (ORDER_ID, METHOD, TRANSACTION_ID, AMOUNT, CURRENCY)
VALUES (SRC.ORDER_ID, [Link], SRC.TRANSACTION_ID, [Link],
[Link]);
-- STORED PROCEDURE --
CREATE OR REPLACE PROCEDURE
JULY_PROJECT.STANDARD_LAYER.SP_LOAD_DIM_PAYMENTS_SCD1()
RETURNS VARCHAR
LANGUAGE SQL
AS
$$
BEGIN
MERGE INTO JULY_PROJECT.STANDARD_LAYER.dim_payments TGT
USING JULY_PROJECT.STAGE_LAYER.payments SRC
ON TGT.TRANSACTION_ID = SRC.TRANSACTION_ID
AND TGT.ORDER_ID = SRC.ORDER_ID
WHEN MATCHED
AND ([Link] <> [Link] OR [Link] <> [Link] OR
[Link] <> [Link])
THEN
UPDATE SET
[Link] = [Link],
[Link] = [Link],
[Link] = [Link]
WHEN NOT MATCHED THEN
INSERT (ORDER_ID, METHOD, TRANSACTION_ID, AMOUNT, CURRENCY)
VALUES (SRC.ORDER_ID, [Link], SRC.TRANSACTION_ID, [Link],
[Link]);
END;
$$;
CALL JULY_PROJECT.STANDARD_LAYER.SP_LOAD_DIM_PAYMENTS_SCD1();
-- TASK --
CREATE OR REPLACE TASK
JULY_PROJECT.STANDARD_LAYER.task_load_dim_payments_scd1
WAREHOUSE = COMPUTE_WH
SCHEDULE = '1 MINUTE'
AS
CALL JULY_PROJECT.STANDARD_LAYER.SP_LOAD_DIM_PAYMENTS_SCD1();
ALTER TASK JULY_PROJECT.STANDARD_LAYER.task_load_dim_payments_scd1 RESUME;
ALTER TASK JULY_PROJECT.STANDARD_LAYER.task_load_dim_payments_scd1
SUSPEND;
SHOW TASKS;
---------------------------------------------- SCD TYPE 2 ----------------------------------------------
CREATE OR REPLACE DATABASE SCD;
TRUNCATE TABLE raw_products_scd2;
TRUNCATE TABLE dim_product_scd2;
SELECT * FROM raw_products_scd2;
SELECT * FROM dim_product_scd2;
SELECT * FROM raw_products_stream;
-- Create the source table
CREATE OR REPLACE TABLE raw_products_scd2 (
product_id INT,
product_name VARCHAR(100),
product_category VARCHAR(50),
product_price DECIMAL(10, 2),
-- This timestamp helps in ordering changes and can be used as start_date for new versions
last_updated_at TIMESTAMP_NTZ DEFAULT CURRENT_TIMESTAMP()
);
INSERT INTO raw_products_scd2 (product_id, product_name, product_category,
product_price) VALUES
(1, 'Laptop Pro X', 'Electronics', 1300.00),
(2, 'Gaming Mouse RGB', 'Electronics', 80.00),
(3, 'Mechanical Keyboard', 'Peripherals', 110.00);
SELECT * FROM raw_products_scd2;
-- Create the target SCD Type 2 dimension table
CREATE OR REPLACE TABLE dim_product_scd2 (
-- Surrogate Key: Unique identifier for each version of a product
sk_product_id INT IDENTITY(1,1),
-- Natural Key: The business key from the source system
product_id INT,
-- Dimension Attributes: The columns you are tracking changes for
product_name VARCHAR(100),
product_category VARCHAR(50),
product_price DECIMAL(10, 2),
-- SCD Type 2 Attributes:
start_date TIMESTAMP_NTZ, -- The date/time this version became effective
end_date TIMESTAMP_NTZ, -- The date/time this version ceased to be effective (NULL if
current)
is_current BOOLEAN, -- Flag: TRUE if this is the currently active version, FALSE
otherwise
load_timestamp TIMESTAMP_NTZ -- When this record was loaded into the dimension table
);
SELECT * FROM dim_product_scd2;
TRUNCATE TABLE dim_product_scd2;
-- Create a stream on the raw_products table
CREATE OR REPLACE STREAM raw_products_stream ON TABLE raw_products_scd2;
-- Check if the stream is empty (it should be, as no changes have occurred yet)
SELECT * FROM raw_products_stream;
-- Initial Load using MERGE (all rows will be 'NOT MATCHED')
MERGE INTO dim_product_scd2 AS T
USING raw_products_scd2 AS S
ON T.product_id = S.product_id AND T.is_current = TRUE -- Match on natural key and current
version
WHEN NOT MATCHED THEN
INSERT (
product_id,
product_name,
product_category,
product_price,
start_date,
end_date,
is_current,
load_timestamp
VALUES (
S.product_id,
S.product_name,
S.product_category,
S.product_price,
S.last_updated_at, -- Use source's last_updated_at as start_date
NULL, -- No end_date for current records
TRUE, -- This is the current version
CURRENT_TIMESTAMP()
);
-- Check the dim_product table after initial load
SELECT * FROM dim_product_scd2 ORDER BY product_id, start_date;
-- Check the stream again - it should be empty now as it's been consumed by the MERGE
SELECT * FROM raw_products_stream;
-- Simulate an update to an existing product
UPDATE raw_products_scd2
SET product_price = 1250.00,
last_updated_at = CURRENT_TIMESTAMP()
WHERE product_id = 1;
select * from raw_products_scd2;
-- STORED PROCEDURE --
CREATE OR REPLACE PROCEDURE SP_LOAD_DIM_PRODUCT_SCD2()
RETURNS VARCHAR
LANGUAGE SQL
AS
$$
BEGIN
-- Step-by-step MERGE for SCD Type 2:
-- 1. Close out old records for updates (WHEN MATCHED)
-- 2. Insert new versions (WHEN NOT MATCHED, or WHEN MATCHED but different relevant
attributes)
-- Ensure atomicity by wrapping in a transaction
BEGIN TRANSACTION;
-- Step 1: Close out old versions (for updates) and soft-delete records
-- This update statement uses the stream to identify records to expire in the dimension.
UPDATE dim_product_scd2 AS T
SET
end_date = S.last_updated_at, -- Use the timestamp from the stream's changed record
is_current = FALSE,
load_timestamp = CURRENT_TIMESTAMP()
FROM raw_products_stream AS S
WHERE T.product_id = S.product_id
AND T.is_current = TRUE
-- This condition captures both:
-- a) 'DELETE' actions (which mean the record was removed from source)
-- b) 'INSERT' actions with ISUPDATE=TRUE (which means the record was updated, this is
the old state)
AND ([Link]$ACTION = 'DELETE' OR [Link]$ISUPDATE = TRUE)
-- Optional: Add conditions to only expire if relevant attributes have changed for updates.
-- For example, if [Link]$ACTION = 'INSERT' AND [Link]$ISUPDATE = TRUE,
-- you might add: AND (T.product_name != S.product_name OR T.product_price !=
S.product_price)
-- However, for simplicity and typical stream usage, just METADATA$ISUPDATE=TRUE is
often enough here.
-- Step 2: Insert new versions (for new records and the new state of updated records)
-- This INSERT statement uses the stream to find new or updated records to add as current.
INSERT INTO dim_product_scd2 (
product_id,
product_name,
product_category,
product_price,
start_date,
end_date,
is_current,
load_timestamp
SELECT
S.product_id,
S.product_name,
S.product_category,
S.product_price,
S.last_updated_at, -- Start date for the new version
NULL, -- Current version, so no end_date
TRUE, -- This is the current version
CURRENT_TIMESTAMP()
FROM raw_products_stream AS S
WHERE [Link]$ACTION = 'INSERT'
-- Optional: Add a check here if you want to skip inserting a new version
-- if the attributes haven't meaningfully changed. This requires joining back to T.
-- This check is usually done in the first UPDATE statement or at the source level.
COMMIT; -- Commit the transaction to make both changes permanent and atomic
-- Return a message indicating success (optional)
RETURN 'SCD Type 2 load for DIM_PRODUCT_SCD2 completed successfully.';
EXCEPTION
WHEN OTHER THEN
ROLLBACK;
RETURN 'SCD Type 2 load for DIM_PRODUCT_SCD2 failed: ' || SQLERRM;
END;
$$;
CALL SP_LOAD_DIM_PRODUCT_SCD2();