diff --git a/airflow/dags/config_interface_dag.py b/airflow/dags/config_interface_dag.py index 3ccf142..640cdc9 100644 --- a/airflow/dags/config_interface_dag.py +++ b/airflow/dags/config_interface_dag.py @@ -81,7 +81,7 @@ load_training_results = PythonOperator( task_id="load_request_submissions", python_callable=call_stored_proc, dag=dag, - op_kwargs={'proc_name':'', 'file_type':'request_submission'} + op_kwargs={'proc_name':'LOAD_REQUEST_SUBMISSION', 'file_type':'request_submission'} ) load_attempt_logs_sp = PythonOperator( diff --git a/airflow/dags/raw_training_data_dag.py b/airflow/dags/raw_training_data_dag.py index f2462f6..c337f65 100644 --- a/airflow/dags/raw_training_data_dag.py +++ b/airflow/dags/raw_training_data_dag.py @@ -93,7 +93,6 @@ load_attempt_logs_sp = PythonOperator( - end_job = EmptyOperator(task_id='End') 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 3af335a..9f5bd06 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,7 +30,7 @@ BEGIN call stg.log_audit(:procedure_name, 'Section 3', 99, 'START'); -- Copy command to load data - COPY INTO STG.TRAINING_DATA_RAW FROM FROM @STG.RAW_TRAINING_DATA_STAGE + COPY INTO STG.TRAINING_DATA_RAW FROM @STG.RAW_TRAINING_DATA_STAGE FILE_FORMAT = (FORMAT_NAME = 'STG.CSV_HEADER') ON_ERROR = ABORT_STATEMENT; @@ -39,7 +39,7 @@ BEGIN call stg.log_audit(:procedure_name, 'Section 4', 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.TRAINING_DATA_RAW',* + SELECT STG.AUDIT_SID.NEXTVAL, :procedure_name,* FROM (SELECT DISTINCT METADATA$FILENAME, CURRENT_TIMESTAMP(), max(METADATA$FILE_ROW_NUMBER) from @STG.RAW_TRAINING_DATA_STAGE group by 1,2);