import csv from collections import defaultdict import boto3 import datetime from datetime import timedelta start_date = '2021-05-26' end_date = '2021-05-31' start_date_obj = datetime.datetime.strptime(start_date, '%Y-%m-%d') end_date_obj = datetime.datetime.strptime(end_date, '%Y-%m-%d') SWF_DOMAIN = 'prod_swf_feed_ingestion' SWF_LOOKUP_DAYS_BACK = 50 def get_swf_executions( swf_client, domain, workflow_name, num_of_days_back=7): """Lookup SWF executions by workflow name. Args: swf_client (object) boto3 swf client domain (str) SWF domain workflow_name (str) name of workflow num_of_days_back (int) how many days back look for SWF executions """ oldest_date = datetime.date.today() - datetime.timedelta( days=num_of_days_back) response = swf_client.list_closed_workflow_executions( domain=domain, executionFilter={ 'workflowId': workflow_name, }, startTimeFilter={ 'oldestDate': datetime.datetime.combine( oldest_date, datetime.datetime.min.time()) }, ) return response['executionInfos'] def parse_swf_data(swf_data): """Explore swf data and create output dicts. Args: swf_data (dict) Result of get_workflow_execution_history request. Returns: activities (dict): Map eventId: eventType. activities_names (dict): Map eventType: eventId. timings (dict): Map eventId: {'start': timestamp, 'end': timestamp}. """ activities = {} activities_names = {} timings = defaultdict(dict) for i in swf_data['events']: if i.get('eventType') == 'ActivityTaskScheduled': activity_name = i[ 'activityTaskScheduledEventAttributes']['activityType']['name'] activities[i['eventId']] = activity_name activities_names[activity_name] = i['eventId'] if i.get('eventType') == 'ActivityTaskStarted': eventid = i['activityTaskStartedEventAttributes']['scheduledEventId'] timings[eventid]['start'] = i['eventTimestamp'] if i.get('eventType') == 'ActivityTaskCompleted': eventid = i['activityTaskCompletedEventAttributes']['scheduledEventId'] timings[eventid]['end'] = i['eventTimestamp'] return activities, activities_names, timings def process_swf_data(store, swf_data): """Process swf_data for given store and return major info about execution. Args: store (dict): Dist with {platform_name: platform_info}. swf_data (dict) Result of get_workflow_execution_history request. Returns: dict: Dist with main info about execution. """ def return_utc_datetime(timestamp): return timestamp.astimezone(datetime.timezone.utc).strftime( '%Y-%m-%d %H:%M:%S') activities, activities_names, timings = parse_swf_data(swf_data) if store['start_fact_analytics'] not in activities_names: return None first_activity = sorted(activities)[0] end_staging_raw_eventid = activities_names[store['end_staging_raw']] start_fact_analytics_eventid = activities_names[ store['start_fact_analytics']] end_fact_analytics_eventid = activities_names[store['end_fact_analytics']] last_activity = sorted(activities)[-1] ingested_in_fact_analytics_mins = ( timings[end_fact_analytics_eventid]['end'] - timings[start_fact_analytics_eventid]['start'] ).total_seconds() // 60 return { 'date_availability': return_utc_datetime( timings[first_activity]['start']), 'ingested_in_staring_raw': return_utc_datetime( timings[end_staging_raw_eventid]['start']), 'ingested_in_fact_analytics': return_utc_datetime( timings[end_fact_analytics_eventid]['end']), 'ingested_in_fact_analytics_mins': ingested_in_fact_analytics_mins, 'etl_finished': return_utc_datetime(timings[last_activity]['end']) } def write_to_csv(filename, data): """Write execution info in csv file. Args: filename (str): File name. data (dict) Info about execution. """ with open(filename, mode='w') as f: fieldnames = [ 'date', 'date_availability', 'ingested_in_staring_raw', 'ingested_in_fact_analytics', 'ingested_in_fact_analytics_mins', 'etl_finished'] writer = csv.DictWriter( f, fieldnames=fieldnames, delimiter=',', quoting=csv.QUOTE_MINIMAL) writer.writeheader() for row in sorted(data, key=lambda x: x['date']): writer.writerow(row) client = boto3.client('swf', region_name='us-east-1') workflows = { 'spotify': { 'workflow': 'spotify_feed_ingestion_{licensor}-{date}', 'end_staging_raw': 'spotify_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'spotify_feed_ingestion_unload_fact_data', 'end_fact_analytics': 'spotify_feed_ingestion_load_fact_tables' }, 'apple': { 'workflow': 'apple_music_feed_ingestion_{licensor}-{date}', 'end_staging_raw': 'apple_music_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'apple_music_feed_ingestion_load_fact_analytics', 'end_fact_analytics': 'apple_music_feed_ingestion_load_fact_analytics' }, 'itunes': { 'workflow': 'itunes_feed_ingestion_{licensor}-{date}', 'end_staging_raw': 'itunes_feed_ingestion_set_populated_status', 'start_fact_analytics': 'itunes_feed_ingestion_load_staging_fact', 'end_fact_analytics': 'itunes_feed_ingestion_load_fact_tables' }, 'pandora': { 'workflow': 'pandora_feed_ingestion_{licensor}-{date}', 'end_staging_raw': 'pandora_feed_ingestion_mark_staging_raw_table_tasks_complete', 'start_fact_analytics': 'pandora_feed_ingestion_load_staging_fact_table', 'end_fact_analytics': 'pandora_feed_ingestion_load_fact_tables' }, 'amazon_music_unlimited': { 'workflow': 'amazon_music_feed_ingestion_unlimited-{licensor}-{date}', 'end_staging_raw': 'amazon_music_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables', 'end_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables' }, 'amazon_music_prime': { 'workflow': 'amazon_music_feed_ingestion_prime-{licensor}-{date}', 'end_staging_raw': 'amazon_music_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables', 'end_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables' }, 'amazon_music_adsupported': { 'workflow': 'amazon_music_feed_ingestion_adsupported-{licensor}-{date}', 'end_staging_raw': 'amazon_music_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables', 'end_fact_analytics': 'amazon_music_feed_ingestion_load_fact_tables' }, 'rhapsody': { 'workflow': 'rhapsody_feed_ingestion_{licensor}-{date}', 'end_staging_raw': 'rhapsody_feed_ingestion_set_status_to_populated_raw_table', 'start_fact_analytics': 'rhapsody_feed_ingestion_load_staging_fact_table', 'end_fact_analytics': 'rhapsody_feed_ingestion_load_fact_tables' }, 'line': { 'workflow': 'line_feed_ingestion-{date}', 'end_staging_raw': 'line_feed_ingestion_mark_staging_raw_table_tasks_complete', 'start_fact_analytics': 'line_feed_ingestion_load_staging_fact', 'end_fact_analytics': 'line_feed_ingestion_load_fact_table' }, 'uma': { 'workflow': 'uma_feed_ingestion-{date}', 'end_staging_raw': 'uma_feed_ingestion_update_dim_tables', 'start_fact_analytics': 'uma_feed_ingestion_unload_fact_data', 'end_fact_analytics': 'uma_feed_ingestion_load_fact_tables' } } platform = 'spotify' licensor = 'theorchard' workflow_name_template = workflows[platform]['workflow'] platform_data = [] while start_date != end_date: workflow_id = workflow_name_template.format( date=end_date, licensor=licensor) executions = get_swf_executions( swf_client=client, domain=SWF_DOMAIN, workflow_name=workflow_id, num_of_days_back=SWF_LOOKUP_DAYS_BACK) if executions: for execution in executions: execution_history = client.get_workflow_execution_history( domain=SWF_DOMAIN, execution=execution['execution'], maximumPageSize=1000, reverseOrder=True ) result = process_swf_data(workflows[platform], execution_history) if not result: continue result['date'] = end_date platform_data.append(result) print(result) break end_date_obj = end_date_obj - timedelta(days=1) end_date = end_date_obj.strftime('%Y-%m-%d') else: break write_to_csv(f'{platform}_{licensor}.csv', platform_data)