import re from snowflake.snowpark.session import Session def matches_pattern(s): return re.match(r'^USERS_(MESSAGES|CANVAS|CAMPAIGNS|BEHAVIORS)_.*$', s) is not None def get_braze_tables(session: Session): sql = """ SELECT table_name FROM braze_raw_events.INFORMATION_SCHEMA.VIEWS ORDER BY table_name; """ return session.sql(sql).collect() def create_braze_share_tables(session: Session): tables = get_braze_tables(session) tables = get_braze_tables(session) for row in tables: table_name = row["TABLE_NAME"] if matches_pattern(table_name): try: print (table_name) sql = f""" create or replace transient table bz_{table_name} as SELECT * FROM braze_raw_events.DATALAKE_SHARING.{table_name}; """ session.sql(sql).collect() sql = f""" create or replace task bz_{table_name}_hourly_insert schedule = 'USING CRON 0 * * * * UTC' WAREHOUSE = BRAZE_DYNAMIC_TABLE_SYNC_X_SMALL AS DECLARE start_time TIMESTAMP; sql STRING; BEGIN start_time := CURRENT_TIMESTAMP(); sql := ' insert into bz_{table_name} select * from braze_raw_events.DATALAKE_SHARING.{table_name} where SF_CREATED_AT >= TO_CHAR(current_timestamp() - interval ''1 hour'', ''YYYY-MM-DD HH24:00:00'') and SF_CREATED_AT < TO_CHAR(current_timestamp(), ''YYYY-MM-DD HH24:00:00'')'; EXECUTE IMMEDIATE :sql; -- Log the task execution CALL BZ_TASK_MONITOR_INSERT_PROC( 'bz_{table_name}_hourly_insert', :start_time, -1, 'SUCCESS', 'hourly' ); end; """ session.sql(sql).collect() print (f"Task bz_{table_name}_hourly_insert created") sql = f""" alter task bz_{table_name}_hourly_insert resume; """ session.sql(sql).collect() # sql = f""" # create or replace secure view nb_braze_{table_name} as # select * from bz_{table_name}; # """ # session.sql(sql).collect() print (f"Finish table {table_name}") except Exception as e: print (table_name, e)