You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Python多进程异步上传S3分片时单进程耗时过长问题

我帮你分析下这个分片上传时单进程耗时过长的问题哈,结合你的代码和实际场景,可能的原因和优化方案如下:


问题背景

我正在使用Python的multiprocessing.Pool.apply_async结合S3分片上传功能上传2GB文件,步骤是:

  1. 生成分片上传ID(mpu_id)
  2. 将文件分割为512MB分片,通过多进程并行上传
  3. 完成分片上传

但遇到一个问题:我的处理器有4个核心,多数情况下3个进程能在1-2分钟内完成,但总有一个进程耗时长达10-15分钟。


可能的原因分析

  1. 磁盘IO瓶颈:多个进程同时读取同一个文件的不同分片时,机械硬盘(HDD)的随机IO性能会成为明显瓶颈——最后一个分片可能因为磁盘调度策略,需要等待其他进程的IO操作完成,导致等待时间过长。
  2. S3请求随机延迟:AWS S3的请求偶尔会出现抖动,某个分片的请求可能因为网络波动、区域负载过高或临时限流,被延迟处理。
  3. 进程调度与日志阻塞:默认进程池大小等于CPU核心数,可能导致CPU和磁盘IO过度竞争;同时多进程直接写日志会触发锁竞争,拖慢单个进程的执行速度。
  4. 缺少重试机制:如果某个分片上传遇到临时失败,没有重试逻辑会导致该进程停滞或耗时剧增。

针对性优化方案

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 08:12:19