"""Methods for accessing S3 with backoff protection.""" import os import sys import math from botocore.exceptions import ClientError from integration_scripts.utils.awsretry import AWSRetry from integration_scripts.connectors.s3 import s3_client, s3_resource, S3_CONFIG class ProgressPercentage(object): def __init__(self, o_s3bucket, key_name): self._key_name = key_name boto_client = o_s3bucket.meta.client # ContentLength is an int self._size = boto_client.head_object( Bucket=o_s3bucket.name, Key=key_name)['ContentLength'] self._seen_so_far = 0 sys.stdout.write('\n') self.old_check = None def __call__(self, bytes_amount): self._seen_so_far += bytes_amount percentage = (float(self._seen_so_far) / float(self._size)) * 100 check = math.floor(percentage / 25) if self.old_check != check: # TERM_UP_ONE_LINE = '\033[A' # TERM_CLEAR_LINE = '\033[2K' # sys.stdout.write('\r' + TERM_UP_ONE_LINE + TERM_CLEAR_LINE) sys.stdout.write('{} {} MB/{} MB ({}%)\n'.format( self._key_name, str(round(self._seen_so_far / 1024 / 1024, 2)), str(round(self._size / 1024 / 1024, 2)), str(round(percentage, 2)))) sys.stdout.flush() self.old_check = check def get_s3_file_key(*args): """Get the s3 file key from env and filename.""" return os.path.join(*args).replace('\\', '/') @AWSRetry.backoff() def download_fileobj_backoff(key, path, bucket_name): """Download a key from S3 with backoff.""" # Create Bucket bucket_obj = s3_resource.Bucket(bucket_name) progress = ProgressPercentage(bucket_obj, key) # Use download_fileobj for speed and parallel processing with open(path, 'wb') as data: bucket_obj.download_fileobj( key, data, Callback=progress, Config=S3_CONFIG) @AWSRetry.backoff() def copy_from_backoff(src_bucket_name, src_key, dst_key=None, dst_bucket_name=None, **kwargs): """Copy from a key on S3 using S3 Resource with backoff.""" copy_source = { 'Bucket': src_bucket_name, 'Key': src_key } if not dst_bucket_name and not dst_key: return if not dst_key: dst_key = src_key if not dst_bucket_name: dst_bucket_name = src_bucket_name s3_resource.Object(dst_bucket_name, dst_key).copy_from( CopySource=copy_source, **kwargs) @AWSRetry.backoff() def copy_key_backoff( src_bucket_name, src_key, dst_bucket_name=None, dst_key=None, **kwargs): """Copy a key on S3 using S3 Client with backoff.""" copy_source = { 'Bucket': src_bucket_name, 'Key': src_key } if not dst_bucket_name and not dst_key: return if not dst_key: dst_key = src_key if not dst_bucket_name: dst_bucket_name = src_bucket_name # s3_resource.meta.client.copy( # copy_source, dst_bucket_name, dst_key, **kwargs) s3_client.copy(copy_source, dst_bucket_name, dst_key, **kwargs) @AWSRetry.backoff() def filter_backoff(key, bucket_name): """Filter S3 keys with backoff.""" # Create Bucket bucket_obj = s3_resource.Bucket(bucket_name) # Return filtered object filtered = bucket_obj.objects.filter(Prefix=key) return filtered @AWSRetry.backoff() def delete_backoff(key, bucket_name): """Delete a key with backoff.""" s3_resource.Object(bucket_name, key).delete() @AWSRetry.backoff() def head_object_backoff(key, bucket_name): """HEAD an object with backoff.""" response = s3_client.head_object(Bucket=bucket_name, Key=key) return response @AWSRetry.backoff() def get_object_backoff(key, bucket_name): """GET an object with backoff.""" response = s3_client.get_object(Bucket=bucket_name, Key=key) return response @AWSRetry.backoff() def s3_key_exists(key, bucket_name): """Check if an s3 key is in a given bucket.""" try: head_object_backoff(key, bucket_name) except ClientError: return False else: return True def filter_keys_by_ext(key_list, file_ext): """Return only keys from key_list which have extension file_ext.""" file_ext = file_ext.strip('.') file_ext = '.' + file_ext key_list = [f for f in key_list if os.path.splitext(f)[1] == file_ext] return key_list def filter_file_keys(source_key_folder, bucket_name, file_ext=None): """Get all files keys in a given folder.""" # Filter to get s3 leaf keys obj_list = filter_backoff(source_key_folder, bucket_name) # Process and strip key list files_only = [os.path.split(obj.key)[1] for obj in obj_list] files_only = [s for s in files_only if s and '/' not in s] # Limit matches by extension if file_ext: files_only = filter_keys_by_ext(files_only, file_ext) return files_only def move_all_files_in_folder( source_key_folder, output_key_folder, bucket_name, file_ext=None): """Move all files in one 'folder' location to another 'folder.'""" files_only = filter_file_keys(source_key_folder, bucket_name, file_ext) # Loop through leaf keys for k in files_only: # Create fully qualified keys source_key = get_s3_file_key(source_key_folder, k) target_key = get_s3_file_key(output_key_folder, k) try: # Copy key copy_from_backoff(bucket_name, source_key, target_key) except Exception as e: # log.log(logging.INFO, 'Copy failed: {}'.format(str(e))) raise e else: # Delete key delete_backoff(source_key, bucket_name) def delete_all_files_in_folder(source_key_folder, bucket_name, file_ext=None): """Delete all files in one 'folder'""" files_only = filter_file_keys(source_key_folder, bucket_name, file_ext) # Loop through leaf keys for k in files_only: # Create fully qualified keys source_key = get_s3_file_key(source_key_folder, k) try: # Delete key delete_backoff(source_key, bucket_name) except Exception as e: # log.log(logging.INFO, 'Copy failed: {}'.format(str(e))) raise e def delete_file_list(key_list, bucket_name): for key in key_list: delete_backoff(key, bucket_name) @AWSRetry.backoff() def put_string_to_s3(key, string, bucket_name): """Put an arbitrary string to an S3 Object.""" s3_resource.Object(bucket_name, key).put(Body=string) @AWSRetry.backoff() def upload_file_backoff(src_file_name, s3_path, bucket_name): """Upload a single file to S3 with backoff.""" s3_bucket = s3_resource.Bucket(bucket_name) s3_bucket.upload_file(src_file_name, s3_path) # TODO: Refactor out the logs to a return object @AWSRetry.backoff() def upload_files(file_name_list, local_path, s3_path, bucket_name): """Upload a list of files to an S3 Bucket. Args: :param file_name_list: (list) file names as strings :param local_path: (str) local storage path (file read source) :param s3_path: (str) S3 key path :param bucket_name: (str) The name of the bucket """ # Begin S3 transfer s3_bucket = s3_resource.Bucket(bucket_name) print('Uploading {} files to S3 bucket `{}`\n'.format( len(file_name_list), bucket_name)) for finished_file_name in file_name_list: # Create path to local file local_file_full_path = os.path.join( local_path, os.path.basename(finished_file_name)) # Create target path for uploaded file s3_file_key = get_s3_file_key(s3_path, finished_file_name) print('Uploading `{}` to `{}`...'.format( local_file_full_path, s3_file_key)) s3_bucket.upload_file(local_file_full_path, s3_file_key) def get_matching_s3_objects(bucket, prefix='', suffix='', delimeter=''): """ Generate objects in an S3 bucket. :param bucket: Name of the S3 bucket. :param prefix: Only fetch objects whose key starts with this prefix (optional). :param suffix: Only fetch objects whose keys end with this suffix (optional). """ paginator = s3_client.get_paginator("list_objects_v2") kwargs = {'Bucket': bucket} if delimeter: kwargs['Delimeter'] = delimeter # We can pass the prefix directly to the S3 API. If the user has passed # a tuple or list of prefixes, we go through them one by one. if isinstance(prefix, str): prefixes = (prefix, ) else: prefixes = prefix for key_prefix in prefixes: kwargs["Prefix"] = key_prefix for page in paginator.paginate(**kwargs): try: contents = page["Contents"] except KeyError: break for obj in contents: key = obj["Key"] if key.endswith(suffix): yield obj def get_matching_s3_keys(bucket, prefix="", suffix="", delimeter=''): """ Generate the keys in an S3 bucket. :param bucket: The S3 client object. :param bucket: Name of the S3 bucket. :param prefix: Only fetch keys that start with this prefix (optional). :param suffix: Only fetch keys that end with this suffix (optional). """ for obj in get_matching_s3_objects(bucket, prefix, suffix, delimeter): yield obj["Key"]