from streamlit_app.common.local_connection import get_connection_parameters from snowflake.core import Root, CreateMode # type: ignore[attr-defined] from snowflake.core.task.dagv1 import DAG, DAGTask, DAGOperation from datetime import timedelta from snowflake.snowpark import Session connection_params = get_connection_parameters("sme_merch") session = Session.builder.configs(connection_params).create() root = Root(session) tasks = root.databases[session.get_current_database()].schemas[session.get_current_schema()].tasks # https://docs.snowflake.com/en/developer-guide/snowflake-python-api/snowflake-python-managing-functions-procedures#managing-stored-procedures # list existing procedures procedure_iter = ( root.databases[session.get_current_database()] .schemas[session.get_current_schema()] .procedures.iter(like="%") ) for procedure_obj in procedure_iter: print(procedure_obj.name) store_ids = ['S00030','S00088','S00113','S00115','S00121','S00122','S00123'] artists = ['LISA','SZA','Rex Orange County'] basket_creation_in_table = 'BASKET_ANALYSIS_INPUT' basket_creation_store_ref = 'CRM_ECOMMERCE_DATA.CONSOLIDATION_DATA.ECOMMERCE_STORES' basket_creation_order_ref = 'CRM_ECOMMERCE_DATA.CONSOLIDATION_DATA.ECOMMERCE_ORDERS' basket_creation_products_ref = 'CRM_ECOMMERCE_DATA.CONSOLIDATION_DATA.ECOMMERCE_PRODUCTS' basket_creation_fans_ref = 'DELPHI_CRM_DATA.RAW_SALESFORCE_SALES_CLOUD.FAN_C' basket_creation_in_tables = ['BASKET_ANALYSIS_INPUT', 'merch_basket_metrics', 'merch_name_basket_metrics', 'merch_basket_metrics_gender', 'merch_basket_metrics_age'] basket_analysis_out_tables = ['BASKET_ANALYSIS_OUTPUT', 'BASKET_ANALYSIS_OUTPUT_ASSOC','BASKET_ANALYSIS_OUTPUT_ASSOC_ITEM'] with DAG(name="merch_pipeline_dag", warehouse="QA_ETL_WH", schedule=timedelta(hours=4), user_task_timeout_ms=600000) as dag: # Create a task that runs our first step in pipeline main_metrics_task = DAGTask( name="main_metrics_pipeline_task", definition=f"CALL CREATE_SALES_METRICS({store_ids})", ) # Create a task that runs our second step in pipeline dag_basket_creation_task = DAGTask( name="basket_output_task", definition=f"CALL BASKET_INPUT_CREATION({store_ids}, {basket_creation_in_tables}, '{basket_creation_store_ref}', '{basket_creation_order_ref}', '{basket_creation_products_ref}', '{basket_creation_fans_ref}')", ) # Create a task that runs our third step in pipeline dag_basket_analysis_task = DAGTask( name="basket_input_task", definition=f"CALL BASKET_ANALYSIS('{basket_creation_in_table}', {basket_analysis_out_tables}, {artists})", ) # Create a task that runs our fourth step in pipeline dag_fan_insights_task = DAGTask( name="fan_insights_task", definition=f"CALL FAN_INSIGHTS_CREATION()", ) # Shift right and left operators can specify task relationships. main_metrics_task >> dag_basket_creation_task >> dag_basket_analysis_task >> dag_fan_insights_task # type: ignore[operator] schema = root.databases[session.get_current_database()].schemas[session.get_current_schema()] dag_op = DAGOperation(schema) dag_op.deploy(dag, mode=CreateMode.or_replace) # A DAG is not suspended by default so we will suspend the root task that will suspend the full DAG root_task = tasks["merch_pipeline_dag"] root_task.resume() # root_task.suspend()