""" PYTHONPATH=. env/bin/python3.6 dev/backfill_copy_filename.py -env qa """ import argparse import multiprocessing import boto3 from podcast.models.ad_action import AdAction from podcast.connectors import mysql SUCCESS = 'success' FAILURE = 'failure' s3_client = boto3.client('s3') def get_copy_file_urls(): """Return all the ad read copy_path_url that doesn't have a copy_filename.""" with mysql.pod_db_session(read_only=True) as session: query = session.query(AdAction.copy_url_path).filter(AdAction.copy_filename == None) rows = query.all() return [row.copy_url_path for row in rows] def ingest_copy_filename(arguments): copy_url, bucket = arguments try: s3_input_file = s3_client.head_object(Bucket=bucket, Key=copy_url) with mysql.pod_db_session() as session: query = session.query(AdAction).filter(AdAction.copy_url_path == copy_url) query.update({'copy_filename':s3_input_file['Metadata']['original_filename']}) return copy_url, SUCCESS except Exception as e: error_message = 'Backfill failed for {key} in bucket {bucket} as {error}'.format(key=copy_url, bucket=bucket, error=e) print(error_message) return copy_url, FAILURE def start_ingestion(bucket): copy_file_urls = get_copy_file_urls() pool = multiprocessing.Pool(12) worker_input = list( zip( copy_file_urls, [bucket] * len(copy_file_urls) ) ) ingestion_iterator = pool.imap_unordered(ingest_copy_filename, worker_input) for ingestor in (ingestion_iterator): print('Ingestion for {} was a {}'.format(ingestor[0], ingestor[1])) pool.close() pool.join() def main(): """Extract shell arguments, fetch copy original_filenames from s3 & start backfilling.""" argparser = argparse.ArgumentParser(prog='Backfill copy_filename') argparser.add_argument('-env', help='env', default='qa') args = argparser.parse_args() bucket = '{}-orcd-podcast-output-assets'.format(args.env) start_ingestion(bucket) if __name__ == '__main__': main()