""" Seat - SWF execution analyzer tool. It loads executions for specific period and counts time spent on tasks. For each execution it sums time for same name tasks, so it display TOTAL time for each task name. Then it aggregates data by day (using current timezone). Finally it displays average time spent on each task for executions aggregated by days. The output is CSV-ready. First row is column names. """ import argparse from collections import defaultdict import csv from datetime import datetime, timedelta import logging import os import sys from feed_ingestion.util.aws import swf ONE_DAY = timedelta(days=1) logger = logging.getLogger(__name__) script_dir = os.path.dirname(__file__) def get_tasks_aggregated_durations(tasks): """Sum duration by task names.""" tasks_totals = defaultdict(timedelta) for task in tasks: tasks_totals[task['name']] += task['duration'] return {name: duration.total_seconds() for name, duration in tasks_totals.items()} def calculate_aggregated_totals( from_date: datetime, to_date: datetime, swf_domain: str, execution_type: str, aggregate_function: str = 'sum', aggregate_by: timedelta = ONE_DAY, ): """Calculate average durations aggregated by day.""" assert aggregate_function in ('sum', 'avg') result = [] while from_date <= to_date: current_to_date = from_date + aggregate_by executions = swf.list_closed_swf_executions( from_=from_date, to=current_to_date, swf_domain=swf_domain, execution_type=execution_type, ) n_executions = len(executions) logger.info( f'Found {n_executions} executions' f' from {from_date} to {current_to_date}') aggregated_tasks_totals = defaultdict(float) for execution in executions: events = swf.load_all_execution_events( swf_domain=swf_domain, workflow_id=execution['execution']['workflowId'], run_id=execution['execution']['runId'], ) tasks = swf.tasks_from_events(events) tasks_totals = get_tasks_aggregated_durations(tasks) for task_name, duration in tasks_totals.items(): value = duration if aggregate_function == 'avg': value /= n_executions aggregated_tasks_totals[task_name] += value aggregated_totals = { 'from': from_date, 'to': current_to_date, 'num_executions': n_executions, 'totals': aggregated_tasks_totals, } result.append(aggregated_totals) from_date = current_to_date return result def aggregated_totals_as_table(aggregated_totals: list): """ Build CSV-ready table. Columns are dates, first row is task names, other rows are dates """ tasks_names = {} for aggregated in aggregated_totals: tasks_names.update(aggregated['totals']) rows = [] header = ['task_name'] for aggregated in aggregated_totals: header.append(aggregated['from'].date()) rows.append(header) for task_name in tasks_names: row = [task_name] for aggregated in aggregated_totals: tasks_duration = aggregated['totals'].get(task_name, 0) value = int(tasks_duration) row.append(value) rows.append(row) num_executions_row = ['EXECUTIONS'] for aggregated in aggregated_totals: value = aggregated['num_executions'] num_executions_row.append(value) rows.append(num_executions_row) return rows def main(args): """Run main function.""" parser = argparse.ArgumentParser(description=__doc__) today = datetime.now().replace(hour=0, minute=0, second=0, microsecond=0) parser.add_argument( '--date-from', nargs='?', default=today - ONE_DAY, type=lambda s: datetime.strptime(s, '%Y-%m-%d'), help='start-date in form YYYY-MM-DD. Default yesterday', ) parser.add_argument( '--date-to', default=today, type=lambda s: datetime.strptime(s, '%Y-%m-%d'), help='end-date (inclusive) in form YYYY-MM-DD. Default today', ) parser.add_argument( '--swf-domain', default='prod_swf_feed_ingestion', type=str, ) parser.add_argument( '--workflow', type=str, required=True, help='Workflow type like "chartmetric_charts_feed_ingestion"', ) parser.add_argument( '--aggregate-function', type=str, default='sum', help='Aggregate by function. sum|avg. Default is sum', ) args = parser.parse_args(args) # with vcr.use_cassette('seat-cache.yml', record_mode='new_episodes'): aggregated_totals = calculate_aggregated_totals( from_date=args.date_from, to_date=args.date_to, swf_domain=args.swf_domain, execution_type=args.workflow, aggregate_function=args.aggregate_function, ) table = aggregated_totals_as_table(aggregated_totals) csv.writer(sys.stdout, delimiter='\t').writerows(table) if __name__ == '__main__': logging.basicConfig(level=logging.INFO) main(sys.argv[1:])