From f1cc222821ee2b971093ac7da99a2a2e87fcdf7e Mon Sep 17 00:00:00 2001 From: Umang Mistry Date: Mon, 10 Jun 2024 19:34:27 -0500 Subject: [PATCH] Added pipeline outputs data resources --- .../dags/pipeline_processed_outputs_dag.py | 83 +++++++++++++++++++ .../R__008_PIPELINE_PROCESSED_OUTPUT.sql | 12 +++ ...009_LOAD_PIPELINE_PROCESSED_OUTPUTS_SP.sql | 73 ++++++++++++++++ 3 files changed, 168 insertions(+) create mode 100644 airflow/dags/pipeline_processed_outputs_dag.py create mode 100644 snowflake/DEV/config_interface/R__008_PIPELINE_PROCESSED_OUTPUT.sql create mode 100644 snowflake/DEV/config_interface/R__009_LOAD_PIPELINE_PROCESSED_OUTPUTS_SP.sql diff --git a/airflow/dags/pipeline_processed_outputs_dag.py b/airflow/dags/pipeline_processed_outputs_dag.py new file mode 100644 index 0000000..981230f --- /dev/null +++ b/airflow/dags/pipeline_processed_outputs_dag.py @@ -0,0 +1,83 @@ + + +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_doczy_processed_outputs" +DATABASE="DOCZY_DEV" + +TAGS=["dev","config_interface","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 +# This will be replaced with the payload from the event after API connection is setup +default_params = {"pipeline_processed_output": ""} + +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 + if file_type == 'pipeline_processed_output': + file_name = params['pipeline_processed_output'] + + 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_pipeline_processed_output = PythonOperator( + task_id="load_pipeline_processed_output", + python_callable=call_stored_proc, + dag=dag, + op_kwargs={'proc_name':'LOAD_DOCZY_PIPELINE_PROCESSED_OUTPUT', 'file_type':'pipeline_processed_output'} +) + + +end_job = EmptyOperator(task_id='End') + +begin_job >> load_pipeline_processed_output >> end_job \ No newline at end of file diff --git a/snowflake/DEV/config_interface/R__008_PIPELINE_PROCESSED_OUTPUT.sql b/snowflake/DEV/config_interface/R__008_PIPELINE_PROCESSED_OUTPUT.sql new file mode 100644 index 0000000..daa0ecf --- /dev/null +++ b/snowflake/DEV/config_interface/R__008_PIPELINE_PROCESSED_OUTPUT.sql @@ -0,0 +1,12 @@ +-- CREATING THE CONFIG TABLE +CREATE TABLE IF NOT EXISTS STG.DOCZY_PIPELINE_PROCESSED_OUTPUT ( + CONTRACT_NAME VARCHAR(16777216) NOT NULL, + FIELD_NAME VARCHAR(255), + RAW_VALUE VARCHAR(16777216), + NEW_EXTRACTED_VALUE VARCHAR(16777216), + CONFIDENCE_LEVEL VARCHAR(255), + SNIPPET VARCHAR(16777216), + NEW_PAGE_NUMBER VARCHAR(255), + BATCH_ID VARCHAR(255), + CREATED_TIME TIMESTAMP_NTZ(9) +); \ No newline at end of file diff --git a/snowflake/DEV/config_interface/R__009_LOAD_PIPELINE_PROCESSED_OUTPUTS_SP.sql b/snowflake/DEV/config_interface/R__009_LOAD_PIPELINE_PROCESSED_OUTPUTS_SP.sql new file mode 100644 index 0000000..374c3b4 --- /dev/null +++ b/snowflake/DEV/config_interface/R__009_LOAD_PIPELINE_PROCESSED_OUTPUTS_SP.sql @@ -0,0 +1,73 @@ +-- columns to be ingested +-- CONTRACT_NAME VARCHAR(16777216) NOT NULL, +-- FIELD_NAME VARCHAR(255), +-- RAW_VALUE VARCHAR(16777216), +-- NEW_EXTRACTED_VALUE VARCHAR(16777216), +-- CONFIDENCE_LEVEL VARCHAR(255), +-- SNIPPET VARCHAR(16777216), +-- NEW_PAGE_NUMBER VARCHAR(255), +-- BATCH_ID VARCHAR(255), +-- CREATED_TIME TIMESTAMP_NTZ(9) + +CREATE OR REPLACE PROCEDURE STG.LOAD_DOCZY_PIPELINE_PROCESSED_OUTPUT(file_name VARCHAR) +RETURNS STRING +LANGUAGE SQL +EXECUTE AS CALLER +AS +$$ +DECLARE + procedure_name varchar; +BEGIN + + procedure_name := 'LOAD_DOCZY_PIPELINE_PROCESSED_OUTPUT'; + + 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.DOCZY_PIPELINE_PROCESSED_OUTPUT_STAGE + STORAGE_INTEGRATION = dev_bucket_integration + URL = ''s3://doczy-dev-infra-raw-data-ingestion/doczy_pipeline_output/' || :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.DOCZY_PIPELINE_RAW_OUTPUT FROM ( + SELECT + NULLIF(TRIM($1),'') AS CONTRACT_NAME, + NULLIF(TRIM($2),'') AS FIELD_NAME, + NULLIF(TRIM($3),'') AS RAW_VALUE, + NULLIF(TRIM($4),'') AS NEW_EXTRACTED_VALUE, + NULLIF(TRIM($5),'') AS CONFIDENCE_LEVEL, + NULLIF(TRIM($6),'') AS SNIPPET, + NULLIF(TRIM($7),'') AS NEW_PAGE_NUMBER, + NULLIF(TRIM($8),'') AS BATCH_ID, + CURRENT_TIMESTAMP() AS CREATED_TIME + + FROM @STG.DOCZY_PIPELINE_PROCESSED_OUTPUT_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 STG.DIM_AUDIT (AUDIT_SID, TABLE_NAME, SOURCE_FILE_NAME, LOAD_DATE, SOURCE_COUNT) + SELECT STG.AUDIT_SID.NEXTVAL, 'STG.DOCZY_PIPELINE_PROCESSED_OUTPUT',* + FROM + (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.DOCZY_PIPELINE_PROCESSED_OUTPUT_STAGE group by 1,2); + + + + RETURN 'Setup, Load, and Audit Complete'; + + call stg.log_audit(:procedure_name, 'Section 3', 99, 'END'); + + + + +END; +$$; +