"""Dim Record Refresh Tasks. ========================== Garcon task for recording results of Dim (Refresh or Sync) Flow. """ from garcon import task from snowflake_connector.etl_connector import SnowflakeSQLExecutor @task.decorate(timeout=7200) def record_refresh( activity, sf_config, dim_type, timestamp, insert_count_sql, update_count_sql): """Garcon task for generating refresh stats for current dim refresh. Args: activity (ActivityWorker): the activity worker. sf_config (dict): Snowflake credentials. dim_type (str): dimension being refreshed. timestamp (str): YYYY-MM-DD H:i:s time refresh is recorded as insert_count_sql (str): sql to count how many new dimensions have during this dimension refresh. update_count_sql (str): sql to count how many new dimensions have been updated during this dimension refresh. Returns: dict: context variables for sending success message to sns {'record_refresh.sns_message': '' 'record_refresh.sns_subject': ''} """ with SnowflakeSQLExecutor(sf_config=sf_config) as executor: insert_count_row = executor.fetchone(insert_count_sql) update_count_row = executor.fetchone(update_count_sql) sns_message = ( '{dim_type} dimension refresh complete for timestamp: {timestamp}\n' 'Existing dimensions updated: {update_count}\n' 'New dimensions added: {insert_count}\n').format( dim_type=dim_type, timestamp=timestamp, insert_count=insert_count_row[0], update_count=update_count_row[0]) sns_subject = 'Dimension Refresh Complete: {dim_type}'.format( dim_type=dim_type) # TODO(jpenner): save refresh info in dynamodb return { 'sns_message': sns_message, 'sns_subject': sns_subject }