"""Tasks of the FeatureFMFacebook Ingestion Workflow.""" from datetime import datetime, timedelta from os import path import re from garcon import task from garcon_contrib.dynamo_feed_status import garcon_feed_status import smart_open from feed_ingestion.flows import FeatureFmFacebook from feed_ingestion.flows.feature_fm_facebook import config from feed_ingestion.flows.helpers import get_sf_config from feed_ingestion.util.aws import s3 @task.decorate(timeout=1000) def bootstrap(activity, date, reload, report_type, dw_config=None): """Bootstrap workflow by getting the correct configurations. Args: activity (ActivityWorker): The activity worker. date (str): Reporting date (YYYY-MM-DD). report_type (str): Report type 'daily', 'monthly' or both separated by comma. Returns: dict: Context. """ date_obj = datetime.strptime( date, '%Y-%m-%d').date() if date else datetime.today().date() previous_month_date = date_obj - timedelta(days=30) if reload == 'True': activity.logger.info( 'Delete status for feed: {} {} '.format(config.feed_name, date)) garcon_feed_status.delete_status(config.feed_name, date) else: overall_status = garcon_feed_status.\ get_overall_status(config.feed_name, date) if overall_status == garcon_feed_status.STATUS_INGESTED: return { 'stop': True, 'message': '{feed_name} is already ingested for {date}' .format( feed_name=config.feed_name, date=date_obj) } if not report_type: error_message = ('report_type is required ' + '(daily, monthly or both separated by comma)') else: report_type = [i.strip() for i in report_type.split(',')] if not all(report in ['daily', 'monthly'] for report in report_type): error_message = ('Invalid report_type. ' + 'Use: daily, monthly or both separated by comma') else: error_message = None if error_message: return { 'stop': True, 'message': error_message } activity.logger.info('Bootstrap flow: {}'.format(date_obj)) activity.logger.info('Report type: {}'.format(report_type)) # Short circuit flow if overall status is already INGESTED if garcon_feed_status.get_overall_status( config.feed_name, date_obj.strftime( '%Y-%m-%d')) == garcon_feed_status.STATUS_INGESTED: activity.logger.info('Feed already ingested for {}'.format(date)) return { 'stop': True, 'message': '{feed_name} is already ingested for {date:%Y-%m-%d}'.format( feed_name=config.feed_name, date=date_obj)} reports_to_run = [] if 'daily' in report_type: file_pattern = config.source_daily_filename_template.format( year=date_obj.year, month_number=date_obj.strftime('%m'), day_number=date_obj.strftime('%d')) reports_to_run.append(file_pattern) if 'monthly' in report_type: file_pattern = config.source_filename_template.format( year=previous_month_date.year, month_number=previous_month_date.strftime('%m') ) reports_to_run.append(file_pattern) activity.logger.info('Report patterns to run: {}'.format(reports_to_run)) try: file_list = [] for report_pattern in reports_to_run: try: file_list += s3.get_list_of_files_and_directories( config.s3_dir_path_template + report_pattern ) except: # noqa activity.logger.info( '{report_pattern} CSV file does not exist for {date_obj}' .format(report_pattern=report_pattern, date_obj=date_obj)) except: # noqa activity.logger.info('CSV file does not exist {}'.format(date_obj)) garcon_feed_status.set_overall_status( config.feed_name, date_obj.strftime( '%Y-%m-%d'), garcon_feed_status.STATUS_NOT_AVAILABLE) return { 'stop': True, 'message': 'Did not find csv file to ' 'ingest for {feed_name} {report_type} on {date:%Y-%m-%d}' .format( feed_name=config.feed_name, report_type=report_type, date=date_obj)} file_dict_list = [_create_file_dict(s3_key) for s3_key in file_list] filtered_list = [f for f in file_dict_list if f['file_size'] > 300] # empty file is 189 bytes if (len(filtered_list) == 0): garcon_feed_status.set_overall_status( config.feed_name, date_obj.strftime( '%Y-%m-%d'), garcon_feed_status.STATUS_NOT_AVAILABLE) return { 'stop': True, 'message': 'CSV files empty for' ' {feed_name} {report_type} on {date:%Y-%m-%d}' .format( feed_name=config.feed_name, report_type=report_type, date=date_obj) } source_files_dict = { 'files': filtered_list} preprocess_file_name_map = { f['file_name']: f['file_name'].replace( '.csv', '_clean.csv' ) for f in filtered_list } s3_dir_path = config.s3_dir_path_template return dict( feed_name=config.feed_name, secrets_path=config.feed_name, staging_raw_table=config.snowflake_table_names['staging_raw'], s3_dir_path=s3_dir_path, source_files_dict=source_files_dict, preprocess_file_name_map=preprocess_file_name_map, date=date_obj.strftime('%Y-%m-%d') ) def _create_file_dict(s3_key): filename = s3_key.split('/')[-1] file_size = s3.get_key_size( config.s3_dir_path_template + filename, config.expected_bucket_owner ) return {'file_name': filename, 'found': True, 'file_size': file_size} @task.decorate(timeout=1000) def preprocess_source_files( activity, date, file_name_map, fetch_files_dict, s3_dir_path): """Strip non-utf8 chars. Args: activity (Activity): Activity instance. date (str): date being processed. file_name_map (dict): dict containing file_name -> file_name_clean map. fetch_files_dict (dict): dict containing metadata of fetched files. s3_archive_path (str): S3 directory to write files to. """ def preprocess_file(src_file_name): dst_file_name = file_name_map[src_file_name] src_file_url = path.join(s3_dir_path, src_file_name) dst_file_url = path.join(s3_dir_path, dst_file_name) with smart_open.smart_open(src_file_url) as sfd: with smart_open.smart_open(dst_file_url, 'w') as dfd: for line in sfd: cleaned = re.sub( '(?