"""Utility functions for triggering and monitoring AWS Glue jobs.""" import time from typing import Dict, Any, Optional, List, Union from enum import Enum import boto3 import dagster as dg from dagster import OpExecutionContext from src.aws_clients import glue_client, get_glue_client class GlueJobState(str, Enum): """AWS Glue job states.""" STARTING = "STARTING" RUNNING = "RUNNING" STOPPING = "STOPPING" STOPPED = "STOPPED" SUCCEEDED = "SUCCEEDED" FAILED = "FAILED" TIMEOUT = "TIMEOUT" def trigger_glue_job( job_name: str, context: OpExecutionContext, arguments: Optional[Dict[str, str]] = None, wait_for_completion: bool = True, poll_interval: int = 10, timeout: int = 3600, client = None, ) -> Dict[str, Any]: """ Trigger an AWS Glue job and optionally wait for completion. Args: job_name: Name of the Glue job to run context: Dagster execution context for logging arguments: Optional dictionary of arguments to pass to the Glue job wait_for_completion: If True, poll until job completes (default: True) poll_interval: Seconds between status checks (default: 10) timeout: Maximum seconds to wait for job completion (default: 3600) client: Optional boto3 Glue client (uses default if not provided) Returns: Dictionary with job run details including: - job_run_id: The Glue job run ID - state: Final state of the job - execution_time: Total execution time in seconds (if completed) Raises: Exception: If the Glue job fails or times out Example: >>> # Using default client >>> result = trigger_glue_job( ... job_name="my_glue_job", ... context=context, ... arguments={"--date": "2023-01-01"}, ... ) >>> >>> # Using custom client >>> from src.aws_clients import get_glue_client >>> custom_client = get_glue_client(region_name="us-west-2") >>> result = trigger_glue_job( ... job_name="my_glue_job", ... context=context, ... client=custom_client, ... ) """ logger = context.log # Use provided client or default _client = client if client is not None else glue_client # Log client configuration logger.info(f"AWS Glue Client Configuration:") logger.info(f" Region: {_client.meta.region_name}") # Get AWS account ID using STS try: sts_client = boto3.client('sts', region_name=_client.meta.region_name) identity = sts_client.get_caller_identity() logger.info(f" Account ID: {identity['Account']}") logger.info(f" User ARN: {identity['Arn']}") except Exception as e: logger.warning(f" Could not retrieve account info: {e}") logger.info(f"Listing jobs...") response = _client.list_jobs(MaxResults=1) logger.info(f"Jobs found: {response.get('JobNames', [])}") # Prepare arguments job_arguments = arguments or {} # Start the Glue job logger.info(f"Starting Glue job: {job_name}") logger.info(f"Arguments: {job_arguments}") response = _client.start_job_run( JobName=job_name, Arguments=job_arguments, ) job_run_id = response["JobRunId"] logger.info(f"Started Glue job {job_name} with JobRunId: {job_run_id}") # If not waiting, return immediately if not wait_for_completion: return { "job_run_id": job_run_id, "state": "RUNNING", "job_name": job_name, } # Poll for job completion start_time = time.time() elapsed_time = 0 while elapsed_time < timeout: # Get current job status job_status = _client.get_job_run(JobName=job_name, RunId=job_run_id) job_run = job_status["JobRun"] state = job_run["JobRunState"] logger.info(f"Job {job_name} status: {state} (elapsed: {int(elapsed_time)}s)") # Check if job is in a terminal state if state in [GlueJobState.SUCCEEDED, GlueJobState.FAILED, GlueJobState.STOPPED, GlueJobState.TIMEOUT]: execution_time = int(elapsed_time) if state == GlueJobState.SUCCEEDED: logger.info(f"Glue job {job_name} completed successfully in {execution_time}s") return { "job_run_id": job_run_id, "state": state, "execution_time": execution_time, "job_name": job_name, } else: error_message = job_run.get("ErrorMessage", "No error message provided") logger.error(f"Glue job {job_name} failed with state: {state}") logger.error(f"Error message: {error_message}") raise Exception(f"Glue job {job_name} failed with state: {state}. Error: {error_message}") # Wait before next poll time.sleep(poll_interval) elapsed_time = time.time() - start_time # Timeout reached logger.error(f"Glue job {job_name} timed out after {timeout}s") raise Exception(f"Glue job {job_name} timed out after {timeout}s") def create_glue_job_asset( job_name: str, asset_name: Optional[str] = None, description: Optional[str] = None, group_name: str = "glue_jobs", tags: Optional[Dict[str, str]] = None, arguments: Optional[Dict[str, str]] = None, deps: Optional[List[dg.AssetKey]] = None, wait_for_completion: bool = True, ) -> callable: """ Factory function to create a Dagster asset that triggers a Glue job. Args: job_name: Name of the Glue job asset_name: Name of the Dagster asset (defaults to job_name) description: Asset description (defaults to generic description) group_name: Asset group name (default: "glue_jobs") tags: Asset tags arguments: Arguments to pass to Glue job deps: Asset dependencies wait_for_completion: Wait for Glue job to complete (default: True) Returns: A Dagster asset that triggers the Glue job Example: >>> my_glue_asset = create_glue_job_asset( ... job_name="rds_to_s3_export", ... asset_name="rds_export", ... description="Export RDS data to S3", ... arguments={"--table": "users"}, ... ) """ asset_name = asset_name or job_name description = description or f"Trigger AWS Glue job: {job_name}" tags = tags or {} arguments = arguments or {} deps = deps or [] @dg.asset( name=asset_name, description=description, group_name=group_name, tags=tags, deps=deps, ) def glue_job_asset(context: OpExecutionContext) -> Dict[str, Any]: """Asset that triggers a Glue job.""" result = trigger_glue_job( job_name=job_name, context=context, arguments=arguments, wait_for_completion=wait_for_completion, ) # Return metadata for Dagster UI context.add_output_metadata({ "job_run_id": result["job_run_id"], "state": result["state"], "execution_time": result.get("execution_time", "N/A"), }) return result return glue_job_asset def create_glue_job_op( job_name: str, op_name: Optional[str] = None, description: Optional[str] = None, tags: Optional[Dict[str, str]] = None, arguments: Optional[Dict[str, str]] = None, wait_for_completion: bool = True, ) -> callable: """ Factory function to create a Dagster op that triggers a Glue job. Useful for more complex workflows where you need ops instead of assets. Args: job_name: Name of the Glue job op_name: Name of the Dagster op (defaults to job_name) description: Op description tags: Op tags arguments: Arguments to pass to Glue job wait_for_completion: Wait for Glue job to complete (default: True) Returns: A Dagster op that triggers the Glue job """ op_name = op_name or f"{job_name}_op" description = description or f"Trigger AWS Glue job: {job_name}" tags = tags or {} arguments = arguments or {} @dg.op( name=op_name, description=description, tags=tags, ) def glue_job_op(context: OpExecutionContext) -> Dict[str, Any]: """Op that triggers a Glue job.""" result = trigger_glue_job( job_name=job_name, context=context, arguments=arguments, wait_for_completion=wait_for_completion, ) return result return glue_job_op def trigger_glue_job_from_sensor( job_name: str, context: dg.SensorEvaluationContext, arguments: Optional[Dict[str, str]] = None, ) -> bool: """ Trigger a Glue job from a sensor (non-blocking). Note: This triggers the job but doesn't wait for completion. Use this in sensors to kick off async jobs. Args: job_name: Name of the Glue job context: Sensor evaluation context arguments: Arguments to pass to the Glue job Returns: True if job was started successfully, False otherwise Example: >>> @dg.sensor(...) ... def my_sensor(context): ... if should_trigger_glue(): ... trigger_glue_job_from_sensor("my_job", context) ... return dg.RunRequest(...) """ logger = context.log job_arguments = arguments or {} try: logger.info(f"Triggering Glue job: {job_name}") logger.info(f"Arguments: {job_arguments}") response = glue_client.start_job_run( JobName=job_name, Arguments=job_arguments, ) job_run_id = response["JobRunId"] logger.info(f"Started Glue job {job_name} with JobRunId: {job_run_id}") return True except Exception as e: logger.error(f"Failed to trigger Glue job {job_name}: {str(e)}") return False