"""Spotify Marketshare Garcon tasks. Tasks to ingest Spotify Marketshare data. """ from datetime import datetime import os import re from garcon import task from garcon_contrib.dynamo_feed_status import garcon_feed_status from feed_ingestion.flows.helpers import get_sf_config from feed_ingestion.flows.spotify_marketshare import config from feed_ingestion.flows.spotify_marketshare.snowflake_executor import\ SpotifyMarketshareSF from feed_ingestion.tasks import bootstrap as reload from feed_ingestion.util import task_status from feed_ingestion.util.aws import s3 as s3utils STOP_RESPONSE = {'stop': True} @task.decorate(timeout=1000) @reload.reset_dynamodb_status_on_reload(config.feed_name) def bootstrap(activity, date, dw_config=None): """Bootstrap workflow by injecting intial context from config. Args: activity (ActivityWorker): The Garcon activity worker. date (str): Reporting date (YYYY-MM-DD). Returns: dict: Initial context for the workflow. """ date_obj = datetime.strptime(date, '%Y-%m-%d') reports = {} file_patterns = [] for report in config.reports: file_pattern = config.file_pattern.format(report=report, date=date_obj) file_patterns.append(file_pattern) reports[report] = { 'file_pattern': file_pattern, 'temp_table_name': ( config.snowflake_table_names['temp_staging_raw'].format( report=report, date=date_obj)) } return { 'feed_name': config.feed_name, 'secrets_path': config.secrets_path, 'date': date, 's3_archive_path': config.s3['archive'].format(date=date_obj), 's3_download_path': config.s3['drop_path'], 'file_patterns': file_patterns, 'reports': reports, 'staging_raw_table': config.snowflake_table_names['staging_raw'] } @task.decorate(timeout=1000) def check_files_on_s3( activity, feed_name, date, s3_download_path, file_patterns): """Check if there are some new files on s3. Args: activity (ActivityWorker): The activity worker. feed_name (str): Feed name of workflow execution for status updates. date (str): Reporting date (YYYY-MM-DD). s3_download_path (str): Drop location on S3. file_patterns (list): List of file patterns for searching. Returns: source_files_dict (dict): Dict with files metadata. """ files_on_s3 = [] all_files_on_s3 = s3utils.get_list_of_files_and_directories( s3_download_path) for file_pattern in file_patterns: files = [] for file_path in all_files_on_s3: if re.search(file_pattern, file_path): files.append(os.path.basename(file_path)) if not files: garcon_feed_status.set_overall_status( feed_name, date, garcon_feed_status.STATUS_NOT_AVAILABLE) activity.logger.info( 'For {} {} there are no any new files for {} ' 'file pattern'.format(feed_name, date, file_pattern)) return STOP_RESPONSE files_on_s3.extend(files) ingested_files = task_status.get_values( feed_name, date, 'ingested_files_status') # if there are new files if set(files_on_s3) - set(ingested_files): activity.logger.info( 'New files were found for {} {}'.format(feed_name, date)) garcon_feed_status.delete_status(feed_name, date) return { 'source_files_dict': { 'files': [ {'file_name': file_name} for file_name in files_on_s3 ]}} activity.logger.info( 'For {} {} there are no any new files'.format(feed_name, date)) return STOP_RESPONSE @task.decorate(timeout=600) def drop_temp_table(activity, temp_table_name): """Drop temp staging raw table. Args: activity (ActivityWorker): The Garcon activity worker. temp_table_name (str): Table name to drop. """ sf_config = get_sf_config(config.feed_name) with SpotifyMarketshareSF(sf_config) as executor: executor.drop_table(temp_table_name) activity.logger.info('{} was dropped'.format(temp_table_name))