import multiprocessing import traceback import sys import boto3 from boto3.s3.transfer import TransferConfig from botocore.exceptions import ClientError from dejavu import Dejavu from dejavu.logic import decoder import config PAGE_SIZE = 100 LIMIT = 1000 THREADS = 12 BUCKET = 'dev-orcd-raw-assets' PREFIX = '2020-07-20//' JOB_QUEUE = 'arn:aws:batch:us-east-1:103233932089:job-queue/dev-fingerprinting-batch-queue' JOB_DEFINITION = 'dev-fingerprinting-generator-job' paginator_config = { 'PageSize': PAGE_SIZE } client = boto3.client('s3') resource = boto3.resource('s3') batch = boto3.client('batch') transfer_config = TransferConfig(max_concurrency=50, use_threads=True) # example of a large file failing # files = ['2020-07-20//008e7fe2_7a0c_477a_99b6_0a82dca630e9.wav'] def submit_batch_job(arguments): k, f = arguments[0] try: response = batch.submit_job( jobName='generate-{}'.format(f), jobQueue=JOB_QUEUE, jobDefinition=JOB_DEFINITION, parameters={ 'S3bucket': BUCKET, 'S3key': k } ) print(response) return k, True except ClientError as e: print(e) return k, False if __name__ == '__main__': djv = Dejavu(config.DB_CONFIG) djv.load_fingerprinted_filenames() paginator = client.get_paginator('list_objects') page_iterator = paginator.paginate( Bucket=BUCKET, Prefix=PREFIX, PaginationConfig=paginator_config ) files = [] for page in page_iterator: items = page['Contents'] # One page of image files. for item in items: key = item['Key'] filename = decoder.get_audio_name_from_path(key)[0] if filename not in djv.filenames: files.append((key, filename)) if len(files) == LIMIT: break if len(files) == LIMIT: break pool = multiprocessing.Pool(THREADS) worker_input = list(zip(files)) iterator = pool.imap_unordered(submit_batch_job, worker_input) while True: try: key, done = next(iterator) except multiprocessing.TimeoutError: continue except StopIteration: break except Exception: print('Failed job submission for {}'.format(key)) # Print traceback because we can't reraise it here traceback.print_exc(file=sys.stdout) else: print('Submitted job {} {}'.format(key, done)) pool.close() pool.join()