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

使用Python多进程异步上传文件至Azure Blob时程序卡住

问题:多进程上传Azure Blob后程序卡在结束阶段,无法输出最终打印语句

我有一个图片目录需要上传至Azure Blob,使用了如下Python代码:

from multiprocessing import Pool
def upload_file_to_blob(file_path, dest):
     file_name = os.path.basename(file_path)
     blob_client = BlobClient.from_connection_string(connection_string, container_name, os.path.join(dest, file_name))
     with open(file_path, "rb") as data:        
         blob_client.upload_blob(data, overwrite=True)
   
with Pool(processes=10) as pool:
     for file_ind, file_path in enumerate(file_list):
         pool.apply_async(upload_file_to_blob, args=(file_path, file_dest))
     pool.close()
     pool.join()
print(f"done")

所有文件似乎已成功上传,但程序卡在结束阶段,终端未输出最终的print(f"done")语句,请问哪里操作有误?


解决方案

核心问题分析

  1. 未处理子进程异常:apply_async提交的异步任务如果抛出异常,异常会被封装在AsyncResult对象中,若未主动获取并处理,进程池可能因异常子进程未正常退出而挂起,导致pool.join()一直等待。
  2. 多进程规范缺失:未在if __name__ == "__main__"块中执行多进程逻辑,Windows系统下会触发子进程重复初始化的问题,可能导致进程池异常。
  3. 全局变量依赖:直接在子进程函数中使用全局的connection_string和container_name,多进程环境下全局变量的传递不可靠,可能引发隐性错误。

修改后的代码示例

import os
from multiprocessing import Pool
from azure.storage.blob import BlobClient

def upload_file_to_blob(file_path, dest, connection_string, container_name):
    try:
        file_name = os.path.basename(file_path)
        blob_client = BlobClient.from_connection_string(connection_string, container_name, os.path.join(dest, file_name))
        with open(file_path, "rb") as data:        
            blob_client.upload_blob(data, overwrite=True)
        print(f"上传完成: {file_path}")
    except Exception as e:
        print(f"上传失败 [{file_path}]: {str(e)}")

if __name__ == "__main__":
    # 替换为你的实际配置
    connection_string = "你的Azure存储连接字符串"
    container_name = "目标容器名称"
    file_list = ["/path/to/image1.jpg", "/path/to/image2.png"]  # 你的文件列表
    file_dest = "blob中的目标目录"

    with Pool(processes=10) as pool:
        # 保存所有异步任务的结果对象
        task_results = []
        for file_path in file_list:
            task = pool.apply_async(upload_file_to_blob, args=(file_path, file_dest, connection_string, container_name))
            task_results.append(task)
        
        # 遍历检查所有任务,捕获并处理异常
        for task in task_results:
            try:
                task.get()  # 获取任务结果,触发异常抛出
            except Exception as e:
                print(f"任务执行异常: {str(e)}")
        
        pool.close()
        pool.join()
    
    print(f"done")

关键修改点说明

  • 强制多进程规范:所有多进程逻辑放在if __name__ == "__main__"块中,避免子进程重复执行主模块代码。
  • 显式传递参数:将connection_string和container_name作为函数参数传入,避免全局变量引发的问题。
  • 异常双重捕获:
    • 子进程函数内部捕获上传过程的异常,定位具体出错文件。
    • 主进程遍历AsyncResult对象,调用get()主动触发并处理任务异常,确保进程池能正常退出。
  • 任务结果跟踪:保存所有异步任务的结果对象,确保所有任务都被正确监控,避免进程池因隐性异常挂起。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 20:18:10