"""Store Specific Youtube Adjustment Workflow.""" import os from garcon import task from snowflake_connector.etl_connector import SnowflakeSQLExecutor from snowflake_connector.etl_connector import SQLLoader from feed_ingestion.common.base_executor import SnowflakeAWSExecutor from feed_ingestion.flows.helpers import get_sf_config from feed_ingestion.flows.youtube_adjustment import config from feed_ingestion.util import datalytics_util from feed_ingestion.util.aws import s3 as s3utils from feed_ingestion.util.aws.s3 import copy_s3_key, delete_file sql_loader = SQLLoader(__file__) @task.decorate(timeout=1000) def bootstrap(activity): """Bootstrap workflow. Args: activity (ActivityWorker): The Garcon activity worker. Returns: dict: Initial context for the workflow. """ return { 'cms_dict': config.CMS_DICT, 'adjustment_dir': config.ADJUSTMENT_DIR, 's3_root': config.S3_ROOT, 'ya_format': config.STORE_SPECIFIC_YOUTUBE_ADJUSTMENT_FORMAT, 'dev_table': config.ADJUSTMENT_DEV_TABLE, 'prod_table': config.ADJUSTMENT_PROD_TABLE, 'report_table': config.ADJUSTMENT_REPORT_TABLE, 'asset_view_table': config.ASSET_VIEW_TABLE, 'email_list': config.EMAIL_LIST} @task.decorate(timeout=10000) def process_files( activity, cms_dict, adjustment_dir, s3_root, ya_format, dev_table, prod_table): """Process files found on S3. Args: activity (ActivityWorker): The Garcon activity worker. cms_dict (list): List of cmd licensors. adjustment_dir (str): S3 bucket path. s3_root (str): S3 bucket path. ya_format (str): Youtube adjustment format. dev_table (str): Adjustment table name. prod_table (str): Adjustment v2 table name. Returns: dict: Number of rows added """ rows_added = 0 sf_config = get_sf_config(config.secrets_path) executor = SnowflakeAWSExecutor(sf_config) aws_params = executor.get_aws_params() for cms_info in cms_dict: files_on_s3 = [] for file_path in s3utils.get_list_of_files_and_directories( os.path.join(adjustment_dir, cms_info['adjustment_root'])): files_on_s3.append(os.path.basename(file_path)) if not files_on_s3: continue for file_in in files_on_s3: file_date_string = file_in.split('_')[-2] year = file_date_string[:4] month = file_date_string[4:6] day = file_date_string[6:] report_date = year + '-' \ + month + '-' \ + day activity.logger.debug('File: {}, report date: {}.'.format( file_in, report_date)) s3_adjust = os.path.join( s3_root, 'adjustments/archive/', year, month, file_in.replace(adjustment_dir, '') ) copy_s3_key(file_in, s3_adjust) executor.execute_query( sql_loader, 'truncate_table', { 'table_name': dev_table } ) executor.execute_query( sql_loader, 'load_from_s3', { 'table_name': dev_table, 's3_path': s3_adjust, 'format_name': ya_format, **aws_params } ) rows_added += executor.execute_query( sql_loader, 'insert', { 'dev': dev_table, 'prod': prod_table, 'report_date': report_date, 'content_owner': cms_info['CMS'].upper() } ) executor.execute_query( sql_loader, 'truncate_table', { 'table_name': dev_table } ) delete_file(file_in, config.expected_bucket_owner) executor.execute_query( sql_loader, 'update_1', { 'table_name': prod_table, 'content_owner': cms_info['CMS'].upper() } ) executor.execute_query( sql_loader, 'update_2', { 'table_name': prod_table, 'content_owner': cms_info['CMS'].upper() } ) return { 'rows_added': rows_added} @task.decorate(timeout=10000) def send_emails( activity, rows_added, prod_table, report_table, asset_view_table, email_list): """Process files found on S3. Args: activity (ActivityWorker): The Garcon activity worker. rows_added (int): Number of rows processed. prod_table (str): Adjustment v2 table name. report_table (str): Adjusment report table name. asset_view_table (str): Asset view table name. email_list (list): List of emails to send notification to. """ sf_config = get_sf_config(config.secrets_path) SnowflakeSQLExecutor(sf_config).execute_query( sql_loader, 'make_view', { 'table_name': prod_table, 'adjustment_report_table': report_table, 'asset_view_table': asset_view_table, } ) BODY = ( '{} new adjustments.\nLooker:' 'https://theorchard.looker.com/x/7Bh4knZ'.format(rows_added)) datalytics_util.send_emails( to=email_list, sub='Adjustments Updated', body=BODY)