"""Docstring.""" import csv import json import sys import time import boto3 def main(): """Entrypoint.""" filename = sys.argv[1] env = sys.argv[2] after = sys.argv[3] if len(sys.argv) > 3 else None max_running = 15 sfn_arn = f'arn:aws:states:us-east-1:437795906767:stateMachine:{env}-sr-fingerprinter-sfn' # noqa:E501 sfn_client = boto3.client('stepfunctions', region_name='us-east-1') # read in CSV file and validate tracks = [] with open(filename, 'r', encoding='utf-8-sig') as f: reader = csv.DictReader(f) if set(reader.fieldnames) != set(['TRACK_ID']): # noqa:E501 print('Invalid input CSV column format!') exit(1) tracks = [x for x in reader] # sort to allow re-processing from middle process = False tracks = sorted(tracks, key=lambda x: x['TRACK_ID']) num_processed = 0 total_tracks = len(tracks) for track in tracks: track_id = track['TRACK_ID'] # skip if processing only after specific track if not process and after: if track_id == after: process = True num_processed += 1 continue # throttle based on num running state machines while True: sfn_list_response = sfn_client.list_executions( stateMachineArn=sfn_arn, statusFilter='RUNNING', maxResults=max_running ) num_running = len(sfn_list_response['executions']) if num_running >= max_running: time.sleep(5) else: break # kick-off new execution exc_name = '-'.join( [ 'Track', track_id, 'created', str(int(time.time())) ] ) payload = { 'id': track_id, 'label': 'Track', 'operation': 'created', 'payload': {} } print(exc_name) sfn_client.start_execution( stateMachineArn=sfn_arn, name=exc_name, input=json.dumps(payload) ) num_processed += 1 percent_processed = (num_processed / total_tracks) * 100 print(f'Processed {track_id} {num_processed} / {total_tracks} : {percent_processed}%') # noqa:E501 if __name__ == '__main__': main()