"""Job and schedule definitions for frontend assets.""" import dagster as dg from dagster import AssetSelection from datetime import timedelta # Define a job that materializes the frontend assets frontend_hourly_job = dg.define_asset_job( name="frontend_hourly_job", selection=AssetSelection.keys("web_user_event", "web_hook_event", "web_audio_player_actions", "app_event", "app_audio_actions"), description="Hourly job to materialize frontend data (web user events, hook events, web audio player actions, app events, and app audio actions)", ) # Define an hourly schedule for the job @dg.schedule( job=frontend_hourly_job, cron_schedule="15 * * * *", # Runs at 15 minutes past every hour execution_timezone="UTC", ) def frontend_hourly_schedule(context: dg.ScheduleEvaluationContext): """ Schedule that runs the frontend job hourly. Runs at 15 minutes past every hour to allow source data to be available. For example, for the 14:00 hour, this runs at 15:15 UTC. """ scheduled_time = context.scheduled_execution_time previous_hour = scheduled_time - timedelta(hours=1) partition_key = previous_hour.strftime("%Y-%m-%d-%H:00") return dg.RunRequest( partition_key=partition_key, tags={ "schedule_name": "frontend_hourly_schedule", "partition_hour": previous_hour.strftime("%Y-%m-%d %H:00"), "duration": "10 minutes", # duration for slack alerts "dagster/max_runtime": str(15 * 60), # 15 minutes timeout (in seconds) } )