import json import os from abc import ABC from typing import Any, Dict, List, Optional from src.utils.database import invoke_postgres_lambda import pandas as pd from dagster import ( AssetSelection, OpExecutionContext, ScheduleDefinition, asset, define_asset_job, ) class PostgresJob(ABC): """Base class for simple Postgres assets and jobs.""" def __init__( self, name: str, description: str, query_file: str, is_reader: bool = True, schedule: Optional[str] = None, group_name: str = "postgres", monitored: bool = True, owners: Optional[List[str]] = None, metadata: Optional[Dict[str, Any]] = None, tags: Optional[Dict[str, Any]] = None, deps: Optional[List[str]] = None, ): self.name = name self.description = description self.schedule = schedule self.query_file = query_file self.is_reader = is_reader self.group_name = group_name self.monitored = monitored self.owners = owners self.metadata = metadata self.tags = tags or {} # Add tech-alerts tag if not present, default to alerting on all failures if "tech-alerts" not in self.tags: self.tags["tech-alerts"] = "true" self.deps = deps # Create asset and job self.asset = self.create_asset() self.job = self.create_job() if schedule else None self.schedule_def = self.create_schedule() if schedule else None # Merge monitored tag with existing tags if monitored is True if self.monitored: self.tags["monitored"] = "true" def get_query(self) -> str: """Load and return the SQL query from file with table name substituted.""" query_path = os.path.join(os.path.dirname(__file__), self.query_file) with open(query_path, "r") as f: return f.read() def post_execute(self, result: Optional[pd.DataFrame], context: OpExecutionContext) -> None: """Optional hook for child classes to process query results. Args: result: DataFrame for SELECT queries, None for DDL queries context: Dagster execution context for logging and metadata """ pass # Default implementation does nothing def create_asset(self): """Create the Dagster asset.""" instance = self @asset( name=self.name, description=self.description, group_name=self.group_name, owners=self.owners, metadata=self.metadata, tags=self.tags, deps=self.deps, ) def postgres_asset(context: OpExecutionContext) -> Optional[pd.DataFrame]: dagster_run_id = context.run_id context.log.info(f"DAGSTER_RUN_ID: {dagster_run_id}") try: # Get query and params query = instance.get_query() response = invoke_postgres_lambda(query, is_reader=self.is_reader) if response.get("statusCode") != 200: raise Exception(f"Error saving bot records to Postgres: {response.get('body')}") context.log.info(f"Query results:\n{json.dumps(response, indent=2)}") # Parse JSON and pass to post_execute if response.get("body"): parsed_data = json.loads(response["body"]) result = instance.post_execute(parsed_data, context) return result else: instance.post_execute(None, context) return None except Exception as e: context.log.error(f"Query failed: {str(e)}") raise return postgres_asset def create_job(self): """Create a job.""" return define_asset_job( name=f"{self.name}_job", description=self.description, selection=AssetSelection.groups(self.group_name).downstream(), ) def create_schedule(self): """Create a schedule with parameters.""" return ScheduleDefinition( job=self.job, cron_schedule=self.schedule, name=f"{self.name}_schedule", )