import gzip import logging from collections import deque from functools import reduce from io import BytesIO from typing import Iterable logger = logging.getLogger(__name__) def stick_last_two(source: Iterable): deck = deque() for item in source: deck.append(item) if len(deck) >= 3: yield deck.popleft() yield reduce(lambda a, b: a + b, deck) def gunzip_head(data: bytes, limit_lines): io_in = BytesIO(data) with gzip.GzipFile(fileobj=io_in, mode='rb') as in_file: data_out = b'' for line_number, line in enumerate(in_file): if line_number >= limit_lines: return data_out data_out += line def human_readable_size(size, obj='B', decimal_places=2): for unit in ['', 'K', 'M', 'G', 'T', 'P']: if size < 1024.0 or unit == 'P': break size /= 1024.0 return f"{size:.{decimal_places}f} {unit}{obj}" class S3MultipartUpload(object): # AWS throws EntityTooSmall error for parts smaller than 5 MB # An error occurred (EntityTooSmall) when calling the CompleteMultipartUpload operation: Your proposed upload is smaller than the minimum allowed size PART_MINIMUM = 1024 ** 2 * 5 def __init__(self, bucket, key, s3_client): self.bucket = bucket self.key = key self.s3_client = s3_client self.mpu_id = None self.parts = [] def create(self): mpu = self.s3_client.create_multipart_upload(Bucket=self.bucket, Key=self.key) self.mpu_id = mpu["UploadId"] # insecure info exposed TODO: reduce exposed info logger.info(f'Created MultipartUpload {mpu}') return mpu def upload(self, data_chunks: Iterable): if not self.mpu_id: self.create() uploaded_bytes = 0 part_number = 1 for data in data_chunks: part_size = len(data) logger.info(f"Uploading part #{part_number} {human_readable_size(part_size)}") if part_size < self.PART_MINIMUM: raise ValueError(f'Chunk size {part_size} is smaller than the minimum allowed ({self.PART_MINIMUM}) size by S3 MultipartUpload') part = self.s3_client.upload_part( Body=data, Bucket=self.bucket, Key=self.key, UploadId=self.mpu_id, PartNumber=part_number) self.parts.append({"PartNumber": part_number, "ETag": part["ETag"]}) uploaded_bytes += part_size logger.info(f"Uploaded total {human_readable_size(uploaded_bytes)}") part_number += 1 return self.complete() def complete(self): result = self.s3_client.complete_multipart_upload( Bucket=self.bucket, Key=self.key, UploadId=self.mpu_id, MultipartUpload={"Parts": self.parts}) return result