from dagster import AutomationCondition from dagster import AssetSelection # This condition runs asset materialization on a cron schedule for the latest partition # after all upstream dependencies have been updated. See on_cron documentation for more details. hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="0 * * * *", cron_timezone="UTC" ) web_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="15 * * * *", cron_timezone="UTC" ) ddb_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="0 * * * *", cron_timezone="UTC" ) backend_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="15 * * * *", cron_timezone="UTC" ) frontend_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="0 * * * *", cron_timezone="UTC" ) video_hook_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="0 * * * *", cron_timezone="UTC" ) # This condition runs asset materialization on a cron schedule for the latest partition of RDS data rds_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="1 * * * *", cron_timezone="UTC" ) # This condition runs asset materialization on a cron schedule for the latest partition of RDS data long_running_rds_hourly_cron_condition = AutomationCondition.on_cron( cron_schedule="2 * * * *", cron_timezone="UTC" ) daily_cron_condition = AutomationCondition.on_cron( cron_schedule="0 0 * * *", cron_timezone="UTC" ) daily_5am_cron_condition = AutomationCondition.on_cron( cron_schedule="0 5 * * *", cron_timezone="UTC" ) every_6_hours_cron_condition = AutomationCondition.on_cron( cron_schedule="0 */6 * * *", cron_timezone="UTC" ) agg_user_daily_cron_condition = AutomationCondition.on_cron( cron_schedule="0 0 * * *", cron_timezone="UTC" ) # This condition eagerly materializes whenever upstream dependencies have been updated, # excluding the latest time window. eager_without_latest_time_window_condition = AutomationCondition.eager().without( AutomationCondition.in_latest_time_window() ).with_label("eager_without_latest_time_window") # This condition: # (1) For the latest partition, runs the asset on a cron schedule after all upstream dependencies have been updated. # (2) For all older partitions, eagerly materializes whenever upstream dependencies have been updated. hourly_cron_with_eager_historical_backfill_condition = ( (hourly_cron_condition | eager_without_latest_time_window_condition) & AutomationCondition.all_deps_blocking_checks_passed() ) # This condition: # (1) For the latest partition, runs the asset on a cron schedule after all upstream dependencies have been updated. # (2) For all older partitions, eagerly materializes whenever upstream dependencies have been updated. daily_cron_with_eager_historical_backfill_condition = ( (daily_cron_condition | eager_without_latest_time_window_condition) & AutomationCondition.all_deps_blocking_checks_passed() ) # This condition: # (1) For the latest partition, runs the asset on a 6-hour cron schedule after all upstream dependencies have been updated. # (2) For all older partitions, eagerly materializes whenever upstream dependencies have been updated. every_6_hours_cron_with_eager_historical_backfill_condition = ( (every_6_hours_cron_condition | eager_without_latest_time_window_condition) & AutomationCondition.all_deps_blocking_checks_passed() ) # This condition is similar to hourly_cron_with_eager_historical_backfill_condition but excludes rds_usage_plan # from the blocking checks. This allows dim_user to materialize even if rds_usage_plan has issues. hourly_cron_with_eager_historical_backfill_except_usage_plan_condition = ( ( hourly_cron_condition.ignore(AssetSelection.assets("rds_usage_plan")) | eager_without_latest_time_window_condition.ignore(AssetSelection.assets("rds_usage_plan")) ) & AutomationCondition.all_deps_blocking_checks_passed() ) # This condition runs daily at midnight UTC for the latest partition, ignoring rds_usage_plan dependency. # Does NOT backfill historical partitions automatically. daily_cron_except_usage_plan_condition = ( daily_cron_condition.ignore(AssetSelection.assets("rds_usage_plan")) & AutomationCondition.all_deps_blocking_checks_passed() ) # This condition runs daily at 5am UTC for the latest partition, ignoring rds_usage_plan dependency. # Does NOT backfill historical partitions automatically. # Used for assets that need rds_discord_info data (available hourly starting at 00:02 UTC). daily_5am_cron_except_usage_plan_condition = ( daily_5am_cron_condition.ignore(AssetSelection.assets("rds_usage_plan")) & AutomationCondition.all_deps_blocking_checks_passed() ) # This maps dbt tags to corresponding automation conditions DBT_AUTOMATION_CONDITION_MAP = { "hourly_cron_with_eager_historical_backfill_condition": hourly_cron_with_eager_historical_backfill_condition, "daily_cron_with_eager_historical_backfill_condition": daily_cron_with_eager_historical_backfill_condition, "hourly_cron_with_eager_historical_backfill_except_usage_plan_condition": hourly_cron_with_eager_historical_backfill_except_usage_plan_condition, "hourly_cron": hourly_cron_condition, "daily_cron": daily_cron_condition, "daily_cron_except_usage_plan": daily_cron_except_usage_plan_condition, "daily_5am_cron_except_usage_plan": daily_5am_cron_except_usage_plan_condition, "on_cron_15min": AutomationCondition.on_cron(cron_schedule="*/15 * * * *", cron_timezone="UTC"), "on_cron_hourly": AutomationCondition.on_cron(cron_schedule="0 * * * *", cron_timezone="UTC"), "on_cron_daily": AutomationCondition.on_cron(cron_schedule="0 0 * * *", cron_timezone="UTC"), "on_cron_weekly": AutomationCondition.on_cron(cron_schedule="0 0 * * 0", cron_timezone="UTC"), "agg_user_daily_cron": agg_user_daily_cron_condition, }