0% found this document useful (0 votes)
3 views14 pages

Snowflake Database Setup for July Project

The document outlines the creation and management of a Snowflake database named JULY_PROJECT, including schemas for raw, stage, and standard layers. It details the processes for integrating data from AWS S3, creating tables and views, and implementing Slowly Changing Dimensions (SCD) for effective data management. Additionally, it includes procedures and tasks for loading and updating data in dimension tables, ensuring data integrity and historical tracking.

Uploaded by

Harsha Vardhanan
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd
0% found this document useful (0 votes)
3 views14 pages

Snowflake Database Setup for July Project

The document outlines the creation and management of a Snowflake database named JULY_PROJECT, including schemas for raw, stage, and standard layers. It details the processes for integrating data from AWS S3, creating tables and views, and implementing Slowly Changing Dimensions (SCD) for effective data management. Additionally, it includes procedures and tasks for loading and updating data in dimension tables, ensuring data integrity and historical tracking.

Uploaded by

Harsha Vardhanan
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd

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

You might also like