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

Python多线程上传S3偶发漏传 故障安全上传机制设计

S3批量上传漏传问题根因分析

问题场景

Python业务流程单任务执行数据处理后会生成超过20000个小文件,后续流程将文件批量上传至S3时偶发漏传,已做文件存在性校验仍无法定位问题,S3操作基于boto3 resource实例实现,核心代码如下:

count_lock = threading.Lock()
obj_count = 0

def __upload(object_path_pair):
    global obj_count
    sleep_time = 5
    num_retries = 10
    for x in range(0, num_retries):
        try:
            libera_resource.upload_file(*object_path_pair)
            sleep(random.uniform(1, 5))
            with count_lock:
                libera_resource_status = libera_resource.Object(object_path_pair[1]).get()['ResponseMetadata'].get('HTTPStatusCode')
                if libera_resource_status == 200 and obj_count > 0:
                    print(f'Item: {file_name} - HLS segment {obj_count} / {len(segment_upload_list)} uploaded successfully.')
                elif libera_resource_status != 200:
                    print(f'Item: {file_name} - HLS segment {obj_count} / {len(segment_upload_list)} uploaded failed, will be tried again.')
                obj_count += 1
                upload_error = None
        except Exception as upload_error:
            pass
        if upload_error or libera_resource_status != 200:
            sleep(sleep_time)  # wait before trying to fetch the data again
            sleep_time *= 2
        else:
            break

def upload_segments(segment_upload_list):
    global obj_count
    obj_count = 0
    with ThreadPoolExecutor(max_workers=100) as executor:
        executor.map(__upload, segment_upload_list)

upload_segments(segment_upload_list)

漏传根因定位

代码存在5处直接导致漏传的硬伤:

  • 重试逻辑变量未初始化,触发未定义错误直接终止任务
    libera_resource_status变量仅在count_lock持有时的代码块内赋值,如果upload_file调用直接抛出异常(比如网络错误、文件不存在、权限错误),会直接进入except块,跳过变量赋值逻辑。后续执行重试判定if upload_error or libera_resource_status != 200时,会直接抛出UnboundLocalError,该异常不在try捕获范围内,会直接终止当前文件的上传循环,且异常不会被主线程感知,对应文件直接漏传。
  • 异常全量静默吞掉,失败无任何感知
    except块直接用pass吞掉所有异常,既不打印异常栈,也不记录失败文件路径,上传过程中出现的任何错误(连接失败、S3限流、本地文件丢失、权限不足)都无日志可查。同时executor.map返回的是惰性迭代器,代码未遍历该迭代器,子线程抛出的所有异常都不会传递到主线程,主线程会默认所有任务执行成功。
  • 跨线程共享非线程安全的boto3 resource实例,高并发下请求异常
    boto3官方明确说明S3 resource实例非线程安全,100个并发线程共享同一个全局实例时,会出现连接状态错乱、请求上下文串扰、请求静默失败的问题。同时默认urllib3连接池大小仅为10,远小于100的并发数,连接耗尽时会触发连接错误,这类错误又被异常捕获逻辑吞掉,直接导致上传失败。
  • 重试耗尽无兜底逻辑,失败任务直接丢弃
    单文件最多重试10次,如果10次重试全部失败,循环会直接退出,没有任何失败标记、没有收集失败文件做二次补偿,主线程完全感知不到失败。另外计数逻辑存在错误:无论上传是否成功,只要进入锁代码块就会累加obj_count,最终打印的上传进度完全不准,会误导判断认为全部文件上传完成。
  • 上传后额外GET校验逻辑不合理,放大失败概率
    上传完成后主动发GET请求校验状态的逻辑完全多余:boto3的upload_file方法内部已经做了内容MD5校验,调用无异常就代表上传成功。额外的GET请求一方面增加了S3请求量,容易触发限流;另一方面GET请求被放在全局锁的临界区内,所有校验请求串行执行,大量线程阻塞在锁等待上,进一步放大并发稳定性问题。

修复方案

  • 修复变量作用域问题:在__upload函数的重试循环前,初始化upload_error = None、libera_resource_status = 0,避免未定义错误。
  • 禁止无日志吞异常:except块必须打印异常信息、对应失败的本地文件路径和S3对象键,方便排查问题。
  • 正确处理线程池任务结果:遍历executor.map返回的迭代器,或改用concurrent.futures.as_completed捕获每个任务的执行结果和异常,所有任务执行完成后统计失败文件列表,做二次兜底上传。
  • 避免跨线程共享boto3实例:每个工作线程内部初始化独立的S3 client/resource实例,或调整boto3的连接池配置匹配并发数;优先用官方推荐的自适应重试模式,替代手写的重试逻辑,自动处理网络抖动、S3限流场景。
  • 去掉冗余的GET校验逻辑,移除不必要的随机sleep,降低S3请求压力。
  • 调整并发数到20-30区间,避免过高并发触发S3限流或本地端口耗尽问题。

内容的提问来源于stack exchange,提问作者user18726875

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 04:45:38