import warnings from datetime import datetime import dagster as dg from dagster import EnvVar from dagster_snowflake import SnowflakeResource from src.assets.snowflake.dim.user.assets import dim_user from src.assets.snowflake.dim.hook.assets import dim_hook from src.assets.snowflake.raw_table.frontend.app_event.assets import app_event, APP_EVENT_TABLE_NAME from src.assets.snowflake.raw_table.frontend.web_hook_event.assets import web_hook_event, WEB_HOOK_EVENT_START_DATE, WEB_HOOK_EVENT_TABLE_NAME from src.utils.automation_conditions import hourly_cron_with_eager_historical_backfill_condition from src.utils.snowflake.constants import TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, PartitionExpr, Team, Warehouse from src.utils.snowflake.query import JinjaSQLFormatter, PythonStringSQLFormatter warnings.filterwarnings("ignore", category=dg.BetaWarning) IOS_HOOK_EVENT_START_DATE = datetime.strptime('2025-08-20', '%Y-%m-%d') ANDROID_HOOK_EVENT_START_DATE = datetime.strptime('2025-09-02', '%Y-%m-%d') FACT_HOOK_PLAY_START_DATE = min(WEB_HOOK_EVENT_START_DATE, IOS_HOOK_EVENT_START_DATE, ANDROID_HOOK_EVENT_START_DATE) FACT_HOOK_PLAY_TABLE_NAME = "FACT_HOOK_PLAY" class StgWebHookPlayConfig(dg.Config): dedup_table_name: str = "TEMP_WEB_HOOK_PLAYER_ACTIONS_DEDUP" unique_segments_table_name: str = "TEMP_WEB_UNIQUE_SEGMENTS" sessions_table_name: str = "TEMP_WEB_HOOK_SESSIONS" stg_web_hook_play_table_name: str = "STG_WEB_HOOK_PLAY" warehouse: str = Warehouse.FACT_HOOK_PLAY_SMALL.value @dg.asset( name="stg_web_hook_play", description="Staging table for web hook play events. This table deduplicates and sessionizes raw web hook events.", group_name="hooks", partitions_def=dg.HourlyPartitionsDefinition(start_date=WEB_HOOK_EVENT_START_DATE, end_offset=-1), deps=[ dg.AssetDep(web_hook_event, partition_mapping=dg.TimeWindowPartitionMapping(start_offset=-1, end_offset=1)), dg.AssetDep(dim_hook), dg.AssetDep(["PROD", "dim_clip"]), dg.AssetDep(dim_user), ], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24*7*2), owners=[Team.DATA_POD.value], metadata={ "database": EnvVar("SNOWFLAKE_DB").get_value(), "schema": EnvVar("SNOWFLAKE_SCHEMA").get_value(), "table_name": "stg_web_hook_play", "data_start_date": WEB_HOOK_EVENT_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "transient": True, "sla_minutes": 120, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def stg_web_hook_play(context: dg.AssetExecutionContext, snowflake: SnowflakeResource, config: StgWebHookPlayConfig) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() jinja_formatter = JinjaSQLFormatter() python_formatter = PythonStringSQLFormatter() # Because hook sessions can span multiple hour partitions, we use a buffered window to fetch data. # We then use the actual partition window to filter and insert data that belongs to the target partition. is_multi_partition_range = context.has_partition_key_range fetch_window = context.asset_partitions_time_window_for_input(web_hook_event.key.to_user_string()) fetch_start_ts = fetch_window.start fetch_end_ts = fetch_window.end fetch_params = { "buffered_partition_start_date": fetch_start_ts.strftime("%Y-%m-%d"), "buffered_partition_end_date": fetch_end_ts.strftime("%Y-%m-%d"), "buffered_partition_start_hour": fetch_start_ts.hour, "buffered_partition_end_hour": fetch_end_ts.hour, "partition_start_date": context.partition_time_window.start.strftime("%Y-%m-%d"), "partition_end_date": context.partition_time_window.end.strftime("%Y-%m-%d"), "partition_start_hour": context.partition_time_window.start.hour, "partition_end_hour": context.partition_time_window.end.hour, "web_hook_event_table_name": WEB_HOOK_EVENT_TABLE_NAME, "dedup_table_name": config.dedup_table_name, "unique_segments_table_name": config.unique_segments_table_name, "sessions_table_name": config.sessions_table_name, "fact_hook_play_table_name": config.stg_web_hook_play_table_name, } logger.info(f"Processing target partition(s): {context.partition_time_window.start} to {context.partition_time_window.end}") logger.info(f"Fetching raw data from buffered window: {fetch_start_ts} to {fetch_end_ts}") logger.info(f"Fetch params: {fetch_params}") with snowflake.get_connection() as conn: cursor = conn.cursor() logger.info(f"Using warehouse {config.warehouse}") warehouse_query = python_formatter.load("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": config.warehouse}, logger=logger) cursor.execute(warehouse_query) # 1. Create a temporary table with deduplicated data from the buffered window. logger.info(f"Creating temporary table {config.dedup_table_name} with deduplicated data...") dedup_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_web_hook_events_deduplicated.sql", params=fetch_params, logger=logger) cursor.execute(dedup_query) # 2. Create a temporary table with unique play segments logger.info(f"Creating temporary table {config.unique_segments_table_name} with unique play segments.") unique_segments_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_web_unique_hook_play_segments.sql", params=fetch_params, logger=logger) cursor.execute(unique_segments_query) # 3. Create a temporary table with aggregated sessions. logger.info(f"Creating temporary table {config.sessions_table_name} with aggregated hook sessions.") sessions_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_web_hook_sessions.sql", params=fetch_params, logger=logger) cursor.execute(sessions_query) # 4. Delete existing data from the target table for the partition window logger.info(f"Deleting existing data from {config.stg_web_hook_play_table_name} for partition window {context.partition_time_window.start} to {context.partition_time_window.end}.") delete_query = python_formatter.load("src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params, "delete_partition_table_name": config.stg_web_hook_play_table_name}, logger=logger) cursor.execute(delete_query) # 5. Insert new web hook play data into the target table for the partition window logger.info(f"Inserting new web hook play data into {config.stg_web_hook_play_table_name}.") insert_fact_hook_play_web_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_web_hook_play.sql", params=fetch_params, logger=logger) cursor.execute(insert_fact_hook_play_web_query) rows_inserted = cursor.rowcount logger.info(f"Successfully processed partition(s). Inserted {rows_inserted} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": config.stg_web_hook_play_table_name, "partition_time_window_start": dg.MetadataValue.text(context.partition_time_window.start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(context.partition_time_window.end.isoformat()), "dagster/row_count": rows_inserted if not is_multi_partition_range else 0, }, ) class StgIosHookPlayConfig(dg.Config): dedup_table_name: str = "TEMP_IOS_HOOK_PLAYER_ACTIONS_DEDUP" unique_segments_table_name: str = "TEMP_IOS_UNIQUE_SEGMENTS" sessions_table_name: str = "TEMP_IOS_HOOK_SESSIONS" stg_ios_hook_play_table_name: str = "STG_IOS_HOOK_PLAY" warehouse: str = Warehouse.FACT_HOOK_PLAY_SMALL.value @dg.asset( name="stg_ios_hook_play", description="Staging table for iOS hook play events. This table deduplicates and sessionizes raw iOS hook events.", group_name="hooks", partitions_def=dg.HourlyPartitionsDefinition(start_date=IOS_HOOK_EVENT_START_DATE, end_offset=-1), deps=[ dg.AssetDep(app_event, partition_mapping=dg.TimeWindowPartitionMapping(start_offset=-1, end_offset=1)), dg.AssetDep(dim_hook), dg.AssetDep(["PROD", "dim_clip"]), dg.AssetDep(dim_user), ], owners=[Team.DATA_POD.value], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24*7*2), metadata={ "database": EnvVar("SNOWFLAKE_DB").get_value(), "schema": EnvVar("SNOWFLAKE_SCHEMA").get_value(), "table_name": "stg_ios_hook_play", "data_start_date": IOS_HOOK_EVENT_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "transient": True, "sla_minutes": 120, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def stg_ios_hook_play(context: dg.AssetExecutionContext, snowflake: SnowflakeResource, config: StgIosHookPlayConfig) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() jinja_formatter = JinjaSQLFormatter() python_formatter = PythonStringSQLFormatter() # Because hook sessions can span multiple hour partitions, we use a buffered window to fetch data. # We then use the actual partition window to filter and insert data that belongs to the target partition. is_multi_partition_range = context.has_partition_key_range fetch_window = context.asset_partitions_time_window_for_input(app_event.key.to_user_string()) fetch_start_ts = fetch_window.start fetch_end_ts = fetch_window.end fetch_params = { "buffered_partition_start_date": fetch_start_ts.strftime("%Y-%m-%d"), "buffered_partition_end_date": fetch_end_ts.strftime("%Y-%m-%d"), "buffered_partition_start_hour": fetch_start_ts.hour, "buffered_partition_end_hour": fetch_end_ts.hour, "partition_start_date": context.partition_time_window.start.strftime("%Y-%m-%d"), "partition_end_date": context.partition_time_window.end.strftime("%Y-%m-%d"), "partition_start_hour": context.partition_time_window.start.hour, "partition_end_hour": context.partition_time_window.end.hour, "ios_hook_event_table_name": APP_EVENT_TABLE_NAME, "dedup_table_name": config.dedup_table_name, "unique_segments_table_name": config.unique_segments_table_name, "sessions_table_name": config.sessions_table_name, "fact_hook_play_table_name": config.stg_ios_hook_play_table_name, } logger.info(f"Processing target partition(s): {context.partition_time_window.start} to {context.partition_time_window.end}") logger.info(f"Fetching raw data from buffered window: {fetch_start_ts} to {fetch_end_ts}") logger.info(f"Fetch params: {fetch_params}") with snowflake.get_connection() as conn: cursor = conn.cursor() logger.info(f"Using warehouse {config.warehouse}") warehouse_query = python_formatter.load("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": config.warehouse}, logger=logger) cursor.execute(warehouse_query) # 1. Create a temporary table with deduplicated data from the buffered window. logger.info(f"Creating temporary table {config.dedup_table_name} with deduplicated data...") dedup_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_ios_hook_events_deduplicated.sql", params=fetch_params, logger=logger) cursor.execute(dedup_query) # 2. Create a temporary table with unique play segments logger.info(f"Creating temporary table {config.unique_segments_table_name} with unique play segments.") unique_segments_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_ios_unique_hook_play_segments.sql", params=fetch_params, logger=logger) cursor.execute(unique_segments_query) # 3. Create a temporary table with aggregated sessions. logger.info(f"Creating temporary table {config.sessions_table_name} with aggregated hook sessions.") sessions_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_ios_hook_sessions.sql", params=fetch_params, logger=logger) cursor.execute(sessions_query) # 4. Delete existing data from the target table for the partition window logger.info(f"Deleting existing data from {config.stg_ios_hook_play_table_name} for partition window {context.partition_time_window.start} to {context.partition_time_window.end}.") delete_query = python_formatter.load("src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params, "delete_partition_table_name": config.stg_ios_hook_play_table_name}, logger=logger) cursor.execute(delete_query) # 5. Insert new iOS hook play data into the target table for the partition window logger.info(f"Inserting new iOS hook play data into {config.stg_ios_hook_play_table_name}.") insert_fact_hook_play_ios_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_ios_hook_play.sql", params=fetch_params, logger=logger) cursor.execute(insert_fact_hook_play_ios_query) rows_inserted = cursor.rowcount logger.info(f"Successfully processed partition(s). Inserted {rows_inserted} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": config.stg_ios_hook_play_table_name, "partition_time_window_start": dg.MetadataValue.text(context.partition_time_window.start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(context.partition_time_window.end.isoformat()), "dagster/row_count": rows_inserted if not is_multi_partition_range else 0, }, ) class StgAndroidHookPlayConfig(dg.Config): dedup_table_name: str = "TEMP_ANDROID_HOOK_PLAYER_ACTIONS_DEDUP" unique_segments_table_name: str = "TEMP_ANDROID_UNIQUE_SEGMENTS" sessions_table_name: str = "TEMP_ANDROID_HOOK_SESSIONS" stg_android_hook_play_table_name: str = "STG_ANDROID_HOOK_PLAY" warehouse: str = Warehouse.FACT_HOOK_PLAY_SMALL.value @dg.asset( name="stg_android_hook_play", description="Staging table for Android hook play events. This table deduplicates and sessionizes raw Android hook events.", group_name="hooks", partitions_def=dg.HourlyPartitionsDefinition(start_date=ANDROID_HOOK_EVENT_START_DATE, end_offset=-1), deps=[ dg.AssetDep(app_event, partition_mapping=dg.TimeWindowPartitionMapping(start_offset=-1, end_offset=1)), dg.AssetDep(dim_hook), dg.AssetDep(["PROD", "dim_clip"]), dg.AssetDep(dim_user), ], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24*7*2), owners=[Team.DATA_POD.value], metadata={ "database": EnvVar("SNOWFLAKE_DB").get_value(), "schema": EnvVar("SNOWFLAKE_SCHEMA").get_value(), "table_name": "stg_android_hook_play", "data_start_date": ANDROID_HOOK_EVENT_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, "transient": True, "sla_minutes": 120, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def stg_android_hook_play(context: dg.AssetExecutionContext, snowflake: SnowflakeResource, config: StgAndroidHookPlayConfig) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() jinja_formatter = JinjaSQLFormatter() python_formatter = PythonStringSQLFormatter() # Because hook sessions can span multiple hour partitions, we use a buffered window to fetch data. # We then use the actual partition window to filter and insert data that belongs to the target partition. is_multi_partition_range = context.has_partition_key_range fetch_window = context.asset_partitions_time_window_for_input(app_event.key.to_user_string()) fetch_start_ts = fetch_window.start fetch_end_ts = fetch_window.end fetch_params = { "buffered_partition_start_date": fetch_start_ts.strftime("%Y-%m-%d"), "buffered_partition_end_date": fetch_end_ts.strftime("%Y-%m-%d"), "buffered_partition_start_hour": fetch_start_ts.hour, "buffered_partition_end_hour": fetch_end_ts.hour, "partition_start_date": context.partition_time_window.start.strftime("%Y-%m-%d"), "partition_end_date": context.partition_time_window.end.strftime("%Y-%m-%d"), "partition_start_hour": context.partition_time_window.start.hour, "partition_end_hour": context.partition_time_window.end.hour, "android_hook_event_table_name": APP_EVENT_TABLE_NAME, "dedup_table_name": config.dedup_table_name, "unique_segments_table_name": config.unique_segments_table_name, "sessions_table_name": config.sessions_table_name, "fact_hook_play_table_name": config.stg_android_hook_play_table_name, } logger.info(f"Processing target partition(s): {context.partition_time_window.start} to {context.partition_time_window.end}") logger.info(f"Fetching raw data from buffered window: {fetch_start_ts} to {fetch_end_ts}") logger.info(f"Fetch params: {fetch_params}") with snowflake.get_connection() as conn: cursor = conn.cursor() logger.info(f"Using warehouse {config.warehouse}") warehouse_query = python_formatter.load("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": config.warehouse}, logger=logger) cursor.execute(warehouse_query) # 1. Create a temporary table with deduplicated data from the buffered window. logger.info(f"Creating temporary table {config.dedup_table_name} with deduplicated data...") dedup_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_android_hook_events_deduplicated.sql", params=fetch_params, logger=logger) cursor.execute(dedup_query) # 2. Create a temporary table with unique play segments logger.info(f"Creating temporary table {config.unique_segments_table_name} with unique play segments.") unique_segments_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_android_unique_hook_play_segments.sql", params=fetch_params, logger=logger) cursor.execute(unique_segments_query) # 3. Create a temporary table with aggregated sessions. logger.info(f"Creating temporary table {config.sessions_table_name} with aggregated hook sessions.") sessions_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_android_hook_sessions.sql", params=fetch_params, logger=logger) cursor.execute(sessions_query) # 4. Delete existing data from the target table for the partition window logger.info(f"Deleting existing data from {config.stg_android_hook_play_table_name} for partition window {context.partition_time_window.start} to {context.partition_time_window.end}.") delete_query = python_formatter.load("src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params, "delete_partition_table_name": config.stg_android_hook_play_table_name}, logger=logger) cursor.execute(delete_query) # 5. Insert new Android hook play data into the target table for the partition window logger.info(f"Inserting new Android hook play data into {config.stg_android_hook_play_table_name}.") insert_fact_hook_play_android_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/stg_android_hook_play.sql", params=fetch_params, logger=logger) cursor.execute(insert_fact_hook_play_android_query) rows_inserted = cursor.rowcount logger.info(f"Successfully processed partition(s). Inserted {rows_inserted} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": config.stg_android_hook_play_table_name, "partition_time_window_start": dg.MetadataValue.text(context.partition_time_window.start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(context.partition_time_window.end.isoformat()), "dagster/row_count": rows_inserted if not is_multi_partition_range else 0, }, ) class FactHookPlayConfig(dg.Config): stg_web_hook_play_table_name: str = "STG_WEB_HOOK_PLAY" stg_ios_hook_play_table_name: str = "STG_IOS_HOOK_PLAY" stg_android_hook_play_table_name: str = "STG_ANDROID_HOOK_PLAY" fact_hook_play_table_name: str = FACT_HOOK_PLAY_TABLE_NAME warehouse: str = Warehouse.FACT_HOOK_PLAY_SMALL.value @dg.asset( name="fact_hook_play", description="Fact table for hook plays.", group_name="hooks", partitions_def=dg.HourlyPartitionsDefinition(start_date=FACT_HOOK_PLAY_START_DATE, end_offset=-1), deps=[ dg.AssetDep(stg_web_hook_play, partition_mapping=dg.TimeWindowPartitionMapping(allow_nonexistent_upstream_partitions=True)), dg.AssetDep(stg_ios_hook_play, partition_mapping=dg.TimeWindowPartitionMapping(allow_nonexistent_upstream_partitions=True)), dg.AssetDep(stg_android_hook_play, partition_mapping=dg.TimeWindowPartitionMapping(allow_nonexistent_upstream_partitions=True)), ], backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=24*7*2), owners=[Team.DATA_POD.value], metadata={ "database": EnvVar("SNOWFLAKE_DB").get_value(), "schema": EnvVar("SNOWFLAKE_SCHEMA").get_value(), "table_name": FACT_HOOK_PLAY_TABLE_NAME, "data_start_date": FACT_HOOK_PLAY_START_DATE.strftime("%Y-%m-%d"), "cluster_by": "[p_date, p_hour]", "partition_expr": PartitionExpr.HOURLY.value, }, automation_condition=hourly_cron_with_eager_historical_backfill_condition, freshness_policy=TIME_WINDOW_FRESHNESS_POLICY_WARN_1H_FAIL_2H, ) def fact_hook_play(context: dg.AssetExecutionContext, snowflake: SnowflakeResource, config: FactHookPlayConfig) -> dg.MaterializeResult: run_id = context.run.run_id logger = dg.get_dagster_logger() jinja_formatter = JinjaSQLFormatter() python_formatter = PythonStringSQLFormatter() logger.info(f"Initial partition time window for job: {context.partition_time_window.start.strftime('%Y-%m-%d %H:%M:%S')} to {context.partition_time_window.end.strftime('%Y-%m-%d %H:%M:%S')}") is_multi_partition_range = context.has_partition_key_range fetch_params = { 'final_table': { "partition_start_date": context.partition_time_window.start.strftime("%Y-%m-%d"), "partition_end_date": context.partition_time_window.end.strftime("%Y-%m-%d"), "partition_start_hour": context.partition_time_window.start.hour, "partition_end_hour": context.partition_time_window.end.hour, "fact_hook_play_table_name": config.fact_hook_play_table_name, }, } insert_web_hook_plays = context.partition_time_window.end.timestamp() >= WEB_HOOK_EVENT_START_DATE.timestamp() insert_ios_hook_plays = context.partition_time_window.end.timestamp() >= IOS_HOOK_EVENT_START_DATE.timestamp() insert_android_hook_plays = context.partition_time_window.end.timestamp() >= ANDROID_HOOK_EVENT_START_DATE.timestamp() if insert_web_hook_plays: web_hook_play_window = context.asset_partitions_time_window_for_input(stg_web_hook_play.key.to_user_string()) web_start_ts = min(max(web_hook_play_window.start.replace(tzinfo=None), WEB_HOOK_EVENT_START_DATE.replace(tzinfo=None)), web_hook_play_window.end.replace(tzinfo=None)) logger.info(f"Web hook play window: {web_hook_play_window.start.strftime('%Y-%m-%d %H:%M:%S')} to {web_hook_play_window.end.strftime('%Y-%m-%d %H:%M:%S')}") fetch_params['web'] = { "partition_start_date": web_start_ts.strftime("%Y-%m-%d"), "partition_end_date": web_hook_play_window.end.strftime("%Y-%m-%d"), "partition_start_hour": web_start_ts.hour, "partition_end_hour": web_hook_play_window.end.hour, "stg_fact_hook_play_table_name": config.stg_web_hook_play_table_name, "fact_hook_play_table_name": config.fact_hook_play_table_name, } else: logger.info(f"Skipping web hook plays - requested partition time window ({context.partition_time_window}) ends before the start date of web hook events ({WEB_HOOK_EVENT_START_DATE}).") if insert_ios_hook_plays: ios_hook_play_window = context.asset_partitions_time_window_for_input(stg_ios_hook_play.key.to_user_string()) ios_start_ts = min(max(ios_hook_play_window.start.replace(tzinfo=None), IOS_HOOK_EVENT_START_DATE.replace(tzinfo=None)), ios_hook_play_window.end.replace(tzinfo=None)) logger.info(f"iOS hook play window: {ios_hook_play_window.start.strftime('%Y-%m-%d %H:%M:%S')} to {ios_hook_play_window.end.strftime('%Y-%m-%d %H:%M:%S')}") fetch_params['ios'] = { "partition_start_date": ios_start_ts.strftime("%Y-%m-%d"), "partition_end_date": ios_hook_play_window.end.strftime("%Y-%m-%d"), "partition_start_hour": ios_start_ts.hour, "partition_end_hour": ios_hook_play_window.end.hour, "stg_fact_hook_play_table_name": config.stg_ios_hook_play_table_name, "fact_hook_play_table_name": config.fact_hook_play_table_name, } else: logger.info(f"Skipping iOS hook plays - requested partition time window ({context.partition_time_window}) ends before the start date of iOS hook events ({IOS_HOOK_EVENT_START_DATE}).") if insert_android_hook_plays: android_hook_play_window = context.asset_partitions_time_window_for_input(stg_android_hook_play.key.to_user_string()) android_start_ts = min(max(android_hook_play_window.start.replace(tzinfo=None), ANDROID_HOOK_EVENT_START_DATE.replace(tzinfo=None)), android_hook_play_window.end.replace(tzinfo=None)) logger.info(f"Android hook play window: {android_hook_play_window.start.strftime('%Y-%m-%d %H:%M:%S')} to {android_hook_play_window.end.strftime('%Y-%m-%d %H:%M:%S')}") fetch_params['android'] = { "partition_start_date": android_start_ts.strftime("%Y-%m-%d"), "partition_end_date": android_hook_play_window.end.strftime("%Y-%m-%d"), "partition_start_hour": android_start_ts.hour, "partition_end_hour": android_hook_play_window.end.hour, "stg_fact_hook_play_table_name": config.stg_android_hook_play_table_name, "fact_hook_play_table_name": config.fact_hook_play_table_name, } else: logger.info(f"Skipping Android hook plays - requested partition time window ({context.partition_time_window}) ends before the start date of Android hook events ({ANDROID_HOOK_EVENT_START_DATE}).") with snowflake.get_connection() as conn: cursor = conn.cursor() logger.info(f"Using warehouse {config.warehouse}") warehouse_query = python_formatter.load("src/utils/snowflake/queries/use_warehouse.sql", params={"warehouse": config.warehouse}, logger=logger) cursor.execute(warehouse_query) # 1. Delete partition from existing target table logger.info(f"Deleting existing data from {config.fact_hook_play_table_name} for partition(s) {context.partition_time_window.start} to {context.partition_time_window.end}.") delete_query = python_formatter.load("src/utils/snowflake/queries/delete_hourly_partitions.sql", params={**fetch_params['final_table'], "delete_partition_table_name": config.fact_hook_play_table_name}, logger=logger) cursor.execute(delete_query) # 2. Insert web hook play data into target table if insert_web_hook_plays: logger.info(f"Inserting new web hook play data into {config.fact_hook_play_table_name} for partition(s) {context.partition_time_window.start} to {context.partition_time_window.end}.") insert_stg_to_final_table_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/insert_stg_to_final_table.sql", params={**fetch_params['web']}, logger=logger) cursor.execute(insert_stg_to_final_table_query) # 3. Insert iOS hook play data into target table if insert_ios_hook_plays: logger.info(f"Inserting new iOS hook play data into {config.fact_hook_play_table_name} for partition(s) {context.partition_time_window.start} to {context.partition_time_window.end}.") insert_fact_hook_play_ios_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/insert_stg_to_final_table.sql", params={**fetch_params['ios']}, logger=logger) cursor.execute(insert_fact_hook_play_ios_query) # # 4. Insert Android hook play data into target table if insert_android_hook_plays: logger.info(f"Inserting new Android hook play data into {config.fact_hook_play_table_name} for partition(s) {context.partition_time_window.start} to {context.partition_time_window.end}.") insert_fact_hook_play_android_query = jinja_formatter.load("src/assets/snowflake/fact/fact_hook_play/queries/insert_stg_to_final_table.sql", params={**fetch_params['android']}, logger=logger) cursor.execute(insert_fact_hook_play_android_query) rows_inserted = cursor.rowcount logger.info(f"Successfully processed partition(s). Inserted {rows_inserted} rows.") return dg.MaterializeResult( metadata={ "run_id": dg.MetadataValue.text(run_id), "table_name": config.fact_hook_play_table_name, "partition_time_window_start": dg.MetadataValue.text(context.partition_time_window.start.isoformat()), "partition_time_window_end": dg.MetadataValue.text(context.partition_time_window.end.isoformat()), "dagster/row_count": rows_inserted if not is_multi_partition_range else 0, }, )