使用Python多进程异步上传S3分片时单进程耗时过长问题
我帮你分析下这个分片上传时单进程耗时过长的问题哈,结合你的代码和实际场景,可能的原因和优化方案如下:
问题背景
我正在使用Python的multiprocessing.Pool.apply_async结合S3分片上传功能上传2GB文件,步骤是:
- 生成分片上传ID(mpu_id)
- 将文件分割为512MB分片,通过多进程并行上传
- 完成分片上传
但遇到一个问题:我的处理器有4个核心,多数情况下3个进程能在1-2分钟内完成,但总有一个进程耗时长达10-15分钟。
可能的原因分析
- 磁盘IO瓶颈:多个进程同时读取同一个文件的不同分片时,机械硬盘(HDD)的随机IO性能会成为明显瓶颈——最后一个分片可能因为磁盘调度策略,需要等待其他进程的IO操作完成,导致等待时间过长。
- S3请求随机延迟:AWS S3的请求偶尔会出现抖动,某个分片的请求可能因为网络波动、区域负载过高或临时限流,被延迟处理。
- 进程调度与日志阻塞:默认进程池大小等于CPU核心数,可能导致CPU和磁盘IO过度竞争;同时多进程直接写日志会触发锁竞争,拖慢单个进程的执行速度。
- 缺少重试机制:如果某个分片上传遇到临时失败,没有重试逻辑会导致该进程停滞或耗时剧增。
针对性优化方案
1. 优化磁盘读取性能
用内存映射文件替代FileChunkIO,让所有进程共享同一块内存映射,避免重复磁盘读取,大幅降低IO竞争:
import mmap # 在上传前先映射整个文件 with open(local_path, 'rb') as f: mmapped_file = mmap.mmap(f.fileno(), length=0, access=mmap.ACCESS_READ) # 上传时直接切片内存映射 data = mmapped_file[offset:offset+bytes]
2. 调整进程池与日志配置
- 限制进程池大小为
CPU核心数-1,避免资源过度竞争:pool = Pool(processes=os.cpu_count()-1 if os.cpu_count() > 1 else 1) - 用进程锁保证日志输出安全,避免多进程写日志阻塞:
from multiprocessing import Lock log_lock = Lock() def safe_log(message): with log_lock: logger.info(message)
3. 优化S3客户端配置
- 添加超时与重试策略,避免临时网络问题导致的长时间等待:
from botocore.config import Config config = Config( connect_timeout=30, read_timeout=60, retries={'max_attempts': 5} ) client = boto3.client('s3', config=config) - 启用S3传输加速(如果适用),提升上传速度:
mpu = client.create_multipart_upload( Bucket=self.bucket, Key=self.key, UseAccelerateEndpoint=True )
4. 添加分片上传重试机制
用tenacity库给上传函数添加重试逻辑,处理临时失败:
from tenacity import retry, stop_after_attempt, wait_exponential @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) def _upload(mmapped_file, bucket, key, file_size, mpu_id, part_num, offset, bytes): # 原上传逻辑
整合优化后的完整代码
import os import math import mmap import logging from multiprocessing import Pool, Lock import boto3 from botocore.config import Config from tenacity import retry, stop_after_attempt, wait_exponential # 全局配置 success_status = 200 one_mb_bytes = 1024 * 1024 file_part_size = 512 * one_mb_bytes # 512MB region = 'us-east-1' # 进程安全日志锁 log_lock = Lock() logger = logging.getLogger(__name__) logging.basicConfig(level=logging.INFO) def as_percent(num, total): return (num / total) * 100 def safe_log(message): with log_lock: logger.info(message) @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) def _upload(mmapped_file, bucket, key, file_size, mpu_id, part_num, offset, bytes): data = mmapped_file[offset:offset+bytes] try: part = client.upload_part( Body=data, Bucket=bucket, Key=key, UploadId=mpu_id, PartNumber=part_num) safe_log("{0} of {1} uploaded ({2:.3f}%) successfully - Part number : {3} ETag: {4} Status code : {5}".format( offset + bytes, file_size, as_percent(offset + bytes, file_size), part_num, part['ETag'], part["ResponseMetadata"]["HTTPStatusCode"])) if part["ResponseMetadata"]["HTTPStatusCode"] != success_status: with log_lock: logger.warning(f"Part upload failed for below:\nPartNumber : {part_num}\nETag : {part['ETag']}\nUploadId : {mpu_id}\nBucket : {bucket}\nkey : {key}") raise Exception(f"Part {part_num} upload failed with status code {part['ResponseMetadata']['HTTPStatusCode']}") finally: del data # 释放内存 return part, part_num # 回调函数收集分片信息 parts = [] def mycallback(x): parts.append({"PartNumber": x[1], "ETag": x[0]['ETag']}) class FileUploads: def __init__(self, path, bucket): self.path = path self.bucket = bucket self.key = os.path.basename(path) class S3MultipartUpload(FileUploads): def __init__(self, path, bucket, part_size=file_part_size, profile_name=None, region_name=region, verbose=False): self.total_bytes = os.stat(path).st_size self.part_bytes = part_size super().__init__(path, bucket) # 初始化带超时配置的S3客户端 global client config = Config( connect_timeout=30, read_timeout=60, retries={'max_attempts': 5} ) if profile_name: session = boto3.Session(profile_name=profile_name, region_name=region_name) client = session.client('s3', config=config) else: client = boto3.client('s3', region_name=region_name, config=config) # 生成分片上传ID def create(self): safe_log("Executing create function to generate Multi part upload id.") try: mpu = client.create_multipart_upload( Bucket=self.bucket, Key=self.key, UseAccelerateEndpoint=True # 启用传输加速 ) mpu_id = mpu["UploadId"] except Exception as e: with log_lock: logger.error("An error occurred while creating the multipart upload id.", exc_info=True) raise e safe_log(f"Multipart upload id generated for key {self.key} in bucket {self.bucket} : {mpu_id} ") return mpu_id # 并行上传分片 def upload(self, mpu_id, local_path, file_size, bucket, key): safe_log(f"Executing upload function to divide the file into chunks of {self.part_bytes/one_mb_bytes} MB and then uploading it to S3.") try: safe_log(f"Size of the file : {file_size/one_mb_bytes} MB") bytes_per_chunk = file_part_size safe_log(f"Size of each part after dividing the file : {bytes_per_chunk/one_mb_bytes} MB") chunk_amount = int(math.ceil(file_size / float(bytes_per_chunk))) safe_log(f"Number of parts created after dividing the file : {chunk_amount}") # 内存映射文件 with open(local_path, 'rb') as f: mmapped_file = mmap.mmap(f.fileno(), length=0, access=mmap.ACCESS_READ) # 进程池大小设为CPU核心数-1 pool = Pool(processes=os.cpu_count()-1 if os.cpu_count() > 1 else 1) safe_log("Parallel execution started for uploading file parts to S3.") for i in range(chunk_amount): offset = i * bytes_per_chunk remaining_bytes = file_size - offset bytes = min([bytes_per_chunk, remaining_bytes]) part_num = i + 1 r = pool.apply_async(_upload, [mmapped_file, bucket, key, file_size, mpu_id, part_num, offset, bytes], callback=mycallback) pool.close() pool.join() mmapped_file.close() safe_log("Parallel execution completed and pool is closed.") except Exception as e: with log_lock: logger.error("An error occurred while uploading parts to S3.", exc_info=True) raise e safe_log("Upload process got completed.") return parts # 完成分片合并 def complete(self, mpu_id, parts): safe_log("Executing Multipart upload completion function to combine all the parts into one in S3.") try: parts_sorted = sorted(parts, key=lambda i: i['PartNumber']) result = client.complete_multipart_upload( Bucket=self.bucket, Key=self.key, UploadId=mpu_id, MultipartUpload={"Parts": parts_sorted}) except Exception as e: with log_lock: logger.error("An error occurred while completing multipart upload.", exc_info=True) raise e safe_log(f"Checking status code after completion of multipart upload : {result['ResponseMetadata']['HTTPStatusCode']}") return result # 完整上传流程 def multi_part_upload(self, local_path, file_size): mpu_id = self.create() try: parts = self.upload(mpu_id, local_path, file_size, self.bucket, self.key) result = self.complete(mpu_id, parts) status_code = result["ResponseMetadata"]["HTTPStatusCode"] if status_code == success_status: safe_log("Multipart upload was closed and completed successfully.") else: with log_lock: logger.warning("Multipart upload was not completed successfully.") return status_code except Exception as e: # 上传失败时终止分片上传,避免残留资源 client.abort_multipart_upload( Bucket=self.bucket, Key=self.key, UploadId=mpu_id ) with log_lock: logger.error(f"Multipart upload aborted due to error: {str(e)}") raise e
内容的提问来源于stack exchange,提问作者Sandeep Singh
相关产品推荐
相关产品推荐

