"""Tasks for working with AWS S3.""" import re import time import boto import boto3 from garcon import task s3_client = boto3.client('s3') s3_resource = boto3.resource('s3') @task.decorate(timeout=12000) def wait_until_empty_s3_path(activity, s3_root): """Wait until a S3 path is empty. Args: activity (ActivityWorker): The activity worker. s3_root (str): The S3 path (e.g. s3://bucket/path/). """ source = re.match('s3://(.+?)/(.*)$', s3_root) bucket_name, prefix = source.groups() bucket = boto.connect_s3().get_bucket(bucket_name) activity.logger.info('Waiting for empty S3 path...') while True: try: files = iter(bucket.list(prefix=prefix)) current = next(files) # S3 includes the parent path an object. If the prefix is the same # just checking against the next file. if current.name == prefix: next(files) time.sleep(10) except StopIteration: activity.logger.info('All files were processed, exit') break def _s3_objects(bucket, prefix, limit=None): bucket_obj = s3_resource.Bucket(bucket) if limit: return bucket_obj.objects.filter(Prefix=prefix).limit(limit) return bucket_obj.objects.filter(Prefix=prefix) def copy_key(source_bucket, source_key, target_bucket, target_key): """Copy S3 key. Args: source_bucket (str): S3 source bucket name. source_key (str): S3 source key name. target_bucket (str): S3 target bucket name. target_key (str): S3 target key name. """ bucket = s3_resource.Bucket(target_bucket) bucket.copy({'Bucket': source_bucket, 'Key': source_key}, target_key) def delete_key(bucket, key): """Delete S3 key. Args: bucket (str): S3 bucket name. key (str): S3 key name. """ s3_client.delete_object(Bucket=bucket, Key=key) def copy_keys(source_folder, target_folder, amount): """Copy s3 keys. Args: source_folder (dict): Source S3 folder bucket and prefix. target_folder (dict): Target S3 folder bucket and prefix. amount (int): Amount of files for copy. """ source_bucket = source_folder['bucket'] target_bucket = target_folder['bucket'] s3_objects = _s3_objects(source_bucket, source_folder['prefix'], amount) for s3_object in s3_objects: source_key = s3_object.key key_name = source_key.split('/')[-1] target_key = '{}/{}'.format(target_folder['prefix'], key_name) copy_key(source_bucket, source_key, target_bucket, target_key) delete_key(source_bucket, source_key) def get_number_of_files_in_s3_folder(bucket, prefix): """Calculate number of files in S3 folder. Args: bucket (str): S3 bucket name. prefix (str): S3 folder prefix. Returns: int: number of files in S3 folder. """ s3_objects = _s3_objects(bucket, prefix) return len(list(s3_objects)) @task.decorate(timeout=21600) def granular_folder_copy( activity, source_folder, target_folder, max_processed_files): """Granularly copy files for source path to target path. Function assumes that copied to the target path files will be processed and removed. Function copies required amount of files to the target path and keep this number until all files will be processed. """ while True: source_number_of_files = get_number_of_files_in_s3_folder( source_folder['bucket'], source_folder['prefix']) activity.logger.info( '{} files in source {}.'.format( source_number_of_files, source_folder)) target_number_of_files = get_number_of_files_in_s3_folder( target_folder['bucket'], target_folder['prefix']) if source_number_of_files == 0 and target_number_of_files == 0: break activity.logger.info( '{} files in target {}.'.format( target_number_of_files, target_folder)) if target_number_of_files < max_processed_files: for_copy = min( max_processed_files - target_number_of_files, source_number_of_files) activity.logger.info('{} files will be copied'.format(for_copy)) copy_keys(source_folder, target_folder, for_copy) time.sleep(2)