From f8a255f7433ffb574d753cb55af176bb8e56656c Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Fri, 15 Mar 2024 16:31:43 -0500 Subject: [PATCH] Updated init objects in Snowflake --- airflow/dags/config_interface_dag.py | 4 +- airflow/dags/raw_training_data_dag.py | 100 ++++++++++++++++++ .../R__005_LOAD_REQUEST_SUBMISSION_SP.sql | 69 ++++++++++++ .../training_interface/R__001_INIT_OBJS.sql | 24 ++--- .../R__006_LOAD_TRAINING_DATA_RAW_SP.sql | 16 +-- 5 files changed, 191 insertions(+), 22 deletions(-) create mode 100644 airflow/dags/raw_training_data_dag.py create mode 100644 snowflake/scripts/config_interface/R__005_LOAD_REQUEST_SUBMISSION_SP.sql diff --git a/airflow/dags/config_interface_dag.py b/airflow/dags/config_interface_dag.py index dcbc12a..3ccf142 100644 --- a/airflow/dags/config_interface_dag.py +++ b/airflow/dags/config_interface_dag.py @@ -81,14 +81,14 @@ load_training_results = PythonOperator( task_id="load_request_submissions", python_callable=call_stored_proc, dag=dag, - op_kwargs={'proc_name':'LOAD_TRAINING_RESULTS', 'file_type':'request_submission'} + op_kwargs={'proc_name':'', 'file_type':'request_submission'} ) load_attempt_logs_sp = PythonOperator( task_id="load_contract_config", python_callable=call_stored_proc, dag=dag, - op_kwargs={'proc_name':'LOAD_ATTEMPT_LOGS', 'file_type':'contract_config'} + op_kwargs={'proc_name':'LOAD_CONTRACT_CONFIG', 'file_type':'contract_config'} ) diff --git a/airflow/dags/raw_training_data_dag.py b/airflow/dags/raw_training_data_dag.py new file mode 100644 index 0000000..f2462f6 --- /dev/null +++ b/airflow/dags/raw_training_data_dag.py @@ -0,0 +1,100 @@ + + +from __future__ import annotations + +import os +from datetime import datetime +import logging + +from airflow import DAG +from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator +from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook +from airflow.operators.python import PythonOperator +from airflow.utils.trigger_rule import TriggerRule +from airflow.operators.empty import EmptyOperator +from airflow.models import Variable +from airflow.exceptions import AirflowFailException + + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger(__name__) + + + +SNOWFLAKE_CONN_ID = "doczy_dev_snowflake" +DAG_ID = "load_raw_training_data" +DATABASE="DOCZY_DEV" +# bucket = "airflow-data-ingestion" + +TAGS=["dev","training_data","dataload"] + +# Trigger rules +ALL_SUCCESS = 'all_success' +ALL_FAILED = 'all_failed' +ALL_DONE = 'all_done' +ONE_SUCCESS = 'one_success' +ONE_FAILED = 'one_failed' + +# Passing empty params for now, this will be overridden by the payload from the trigger +# These params can also be set from the Airflow UI while manually triggering the DAG +default_params = {"column_config_file_name": "", "raw_training_data_file_name":""} +# This will be replaced with the payload from the event after API connection is setup +# training_results_file_name = "training_results_sample.csv" +# attempt_logs_file_name = "attempt_logs_sample.csv" + + +def call_stored_proc(proc_name,file_type, params): + dwh_hook = SnowflakeHook(snowflake_conn_id=SNOWFLAKE_CONN_ID) + with dwh_hook.get_conn() as conn: + # dwh_hook.set_autocommit(conn,autocommit=False) + cur = conn.cursor() + + # Added new parameter file_type to determine the file name to be passed to the stored procedure + # The bucket name is set by default to "doczy-dev-infra-raw-data-ingestion" and the files should be ALWAYS save under training_interface/ path for now + if file_type == 'column_config': + file_name = params['column_config_file_name'] + elif file_type == 'raw_training_data': + file_name = params['raw_training_data_file_name'] + + cur.execute(f"CALL {DATABASE}.STG.{proc_name}('{file_name}');") + result = cur.fetchone() + if result[0] == 'Setup, Load, and Audit Complete': + logger.info('PROCEDURE EXECUTED SUCCESSFULLY') + else: + raise AirflowFailException("Check the DAG logs for more information. ERROR FROM SNOWFLAKE: ", result) + logger.info(f"QUERY EXECUTION RESULT: {str(result)}") + +dag = DAG( + DAG_ID, + start_date=datetime(2024, 1, 1), + default_args={"snowflake_conn_id": SNOWFLAKE_CONN_ID, "retries":0}, + tags=TAGS, + catchup=False, + schedule=None, + params = default_params +) + +begin_job = EmptyOperator(task_id='Begin') + + +load_training_results = PythonOperator( + task_id="load_column_config", + python_callable=call_stored_proc, + dag=dag, + op_kwargs={'proc_name':'LOAD_COLUMN_CONFIG', 'file_type':'column_config'} +) + +load_attempt_logs_sp = PythonOperator( + task_id="load_raw_training_data", + python_callable=call_stored_proc, + dag=dag, + op_kwargs={'proc_name':'LOAD_TRAINING_DATA_RAW', 'file_type':'raw_training_data'} +) + + + + +end_job = EmptyOperator(task_id='End') + + +begin_job >> load_training_results >> load_attempt_logs_sp >> end_job \ No newline at end of file diff --git a/snowflake/scripts/config_interface/R__005_LOAD_REQUEST_SUBMISSION_SP.sql b/snowflake/scripts/config_interface/R__005_LOAD_REQUEST_SUBMISSION_SP.sql new file mode 100644 index 0000000..989eb2f --- /dev/null +++ b/snowflake/scripts/config_interface/R__005_LOAD_REQUEST_SUBMISSION_SP.sql @@ -0,0 +1,69 @@ +CREATE OR REPLACE PROCEDURE STG.LOAD_REQUEST_SUBMISSION(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_CONTRACT_CONFIG'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'START'); + + -- Create or replace stage with dynamic file name + EXECUTE IMMEDIATE 'CREATE OR REPLACE STAGE STG.CONTRACT_CONFIG_STAGE + STORAGE_INTEGRATION = dev_bucket_integration + URL = ''s3://doczy-dev-infra-raw-data-ingestion/config_interface/' || :file_name || ''' + FILE_FORMAT = (FORMAT_NAME = ''STG.CSV_HEADER'');'; + + call stg.log_audit(:procedure_name, 'Section 1', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'START'); + + COPY INTO STG.CONTRACT_CONFIG FROM ( + SELECT + NULLIF(TRIM($1),'') AS REQUEST_ID, + NULLIF(TRIM($2),'') AS CONTRACT_NAME, + NULLIF(TRIM($3),'') AS GROUP_1, + NULLIF(TRIM($4),'') AS GROUP_1_OVERRIDE, + NULLIF(TRIM($5),'') AS GROUP_2, + NULLIF(TRIM($6),'') AS GROUP_2_OVERRIDE, + NULLIF(TRIM($7),'') AS GROUP_3, + NULLIF(TRIM($8),'') AS GROUP_3_OVERRIDE, + NULLIF(TRIM($9),'') AS GROUP_4, + NULLIF(TRIM($10),'') AS GROUP_4_OVERRIDE, + NULLIF(TRIM($11),'') AS GROUP_5, + NULLIF(TRIM($12),'') AS GROUP_5_OVERRIDE, + NULLIF(TRIM($13),'') AS REQUEST_DATETIME, + NULLIF(TRIM($14),'') AS REQUEST_USER, + NULLIF(TRIM($15),'') AS OVERRIDE_DATETIME, + NULLIF(TRIM($16),'') AS LATEST_FLAG DEFAULT TRUE, + NULLIF(TRIM($17),'') AS PIPELINE_KICKOFF_DATETIME + FROM @STG.CONTRACT_CONFIG_STAGE + ) + FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + ON_ERROR = ABORT_STATEMENT; + + call stg.log_audit(:procedure_name, 'Section 2', 99, 'END'); + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); + + INSERT INTO DOCZY_DEV.STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, 'STG.CONTRACT_CONFIG',* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.CONTRACT_CONFIG_STAGE group by 1,2); + + + + RETURN 'Setup, Load, and Audit Complete'; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + + + +END; +$$; + diff --git a/snowflake/scripts/training_interface/R__001_INIT_OBJS.sql b/snowflake/scripts/training_interface/R__001_INIT_OBJS.sql index 64eda50..2413b5f 100644 --- a/snowflake/scripts/training_interface/R__001_INIT_OBJS.sql +++ b/snowflake/scripts/training_interface/R__001_INIT_OBJS.sql @@ -16,18 +16,18 @@ CREATE TABLE IF NOT EXISTS STG.TRAINING_RESULTS -- Create file format -CREATE FILE FORMAT IF NOT EXISTS STG.CSV_HEADER - TYPE = 'CSV' - FIELD_DELIMITER = ',' - FILE_EXTENSION = '.csv' - RECORD_DELIMITER = '\\n' - DATE_FORMAT = AUTO - TRIM_SPACE = TRUE - NULL_IF = ('NULL', '', 'N/A','?','~','\\N') - SKIP_HEADER = 1 - EMPTY_FIELD_AS_NULL = TRUE - FIELD_OPTIONALLY_ENCLOSED_BY = '"' - SKIP_BLANK_LINES = TRUE; +-- CREATE FILE FORMAT IF NOT EXISTS STG.CSV_HEADER +-- TYPE = 'CSV' +-- FIELD_DELIMITER = ',' +-- FILE_EXTENSION = '.csv' +-- RECORD_DELIMITER = '\\n' +-- DATE_FORMAT = AUTO +-- TRIM_SPACE = TRUE +-- NULL_IF = ('NULL', '', 'N/A','?','~','\\N') +-- SKIP_HEADER = 0 +-- EMPTY_FIELD_AS_NULL = TRUE +-- FIELD_OPTIONALLY_ENCLOSED_BY = '"' +-- SKIP_BLANK_LINES = TRUE; -- Create staging table for attempt logs diff --git a/snowflake/scripts/training_interface/R__006_LOAD_TRAINING_DATA_RAW_SP.sql b/snowflake/scripts/training_interface/R__006_LOAD_TRAINING_DATA_RAW_SP.sql index 4580a54..b893c32 100644 --- a/snowflake/scripts/training_interface/R__006_LOAD_TRAINING_DATA_RAW_SP.sql +++ b/snowflake/scripts/training_interface/R__006_LOAD_TRAINING_DATA_RAW_SP.sql @@ -30,14 +30,14 @@ BEGIN call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); -- Copy command to load data - COPY INTO STG.TRAINING_DATA_COLUMN_CONFIG FROM ( - SELECT - NULLIF(TRIM($1), '') AS COLUMN_NAME, - NULLIF(TRIM($2), '') AS COLUMN_DATATYPE - FROM @STG.RAW_TRAINING_DATA_STAGE - ) - FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') - ON_ERROR = ABORT_STATEMENT; + -- COPY INTO STG.TRAINING_DATA_COLUMN_CONFIG FROM ( + -- SELECT + -- NULLIF(TRIM($1), '') AS COLUMN_NAME, + -- NULLIF(TRIM($2), '') AS COLUMN_DATATYPE + -- FROM @STG.RAW_TRAINING_DATA_STAGE + -- ) + -- FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') + -- ON_ERROR = ABORT_STATEMENT; call stg.log_audit(:procedure_name, 'Section 3', 99, 'END');