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

Azure Blob Storage文件上传的暂停与恢复实现方案问询

Azure Blob Storage异步上传实现暂停与恢复方案

当然可以实现,核心是利用Azure Blob的分块上传(Block Blob)机制,结合异步暂停控制和本地进度记录来完成。下面是具体的实现思路和代码示例:

实现核心要点

  • 分块上传基础:Block Blob支持将文件拆分为多个独立块(单块最大100MB,推荐4-8MB),每个块上传后可单独记录ID,最终提交所有块ID完成文件组装,天然支持断点续传。
  • 异步暂停控制:用asyncio.Event实现异步环境下的线程安全暂停/恢复,上传每个块前检查状态,暂停时阻塞等待。
  • 进度持久化:将已上传的块ID、文件读取偏移量保存到本地文件,恢复时读取状态跳过已完成块,从断点处继续上传。
  • 异步适配:使用Azure Blob异步客户端AsyncBlobServiceClient+aiofiles异步读文件,避免同步操作阻塞事件循环。

完整代码实现

import asyncio
import aiofiles
import json
import os
from azure.storage.blob.aio import AsyncBlobServiceClient

class BlobUploader:
    def __init__(self, connection_string, container_name):
        self.blob_service_client = AsyncBlobServiceClient.from_connection_string(connection_string)
        self.container_name = container_name
        self.pause_event = asyncio.Event()
        self.pause_event.set()  # 初始状态:未暂停

    async def pause_upload(self):
        """暂停上传"""
        self.pause_event.clear()

    async def resume_upload(self):
        """恢复上传"""
        self.pause_event.set()

    def _get_progress_file(self, file_path):
        """生成对应文件的进度文件名,避免多文件冲突"""
        file_hash = hash(file_path)
        return f"upload_progress_{file_hash}.json"

    async def _save_progress(self, file_path, uploaded_blocks, offset):
        """保存当前上传进度到本地"""
        progress_file = self._get_progress_file(file_path)
        progress_data = {
            "file_path": file_path,
            "uploaded_blocks": uploaded_blocks,
            "offset": offset
        }
        async with aiofiles.open(progress_file, 'w', encoding='utf-8') as f:
            await f.write(json.dumps(progress_data))

    async def _load_progress(self, file_path):
        """加载已保存的上传进度"""
        progress_file = self._get_progress_file(file_path)
        try:
            async with aiofiles.open(progress_file, 'r', encoding='utf-8') as f:
                progress_data = json.loads(await f.read())
                if progress_data["file_path"] == file_path:
                    return progress_data["uploaded_blocks"], progress_data["offset"]
        except (FileNotFoundError, json.JSONDecodeError):
            pass
        return [], 0

    async def upload_single_file(self, file_path, chunk_size=4*1024*1024):
        """单文件异步分块上传,支持暂停恢复"""
        blob_client = self.blob_service_client.get_blob_client(
            container=self.container_name,
            blob=os.path.basename(file_path)  # 可选:用文件名而非完整路径作为Blob名称
        )
        uploaded_blocks, offset = await self._load_progress(file_path)

        # 校验已上传块的有效性
        if uploaded_blocks:
            try:
                # 尝试提交块列表,若成功说明之前已完成上传
                await blob_client.commit_block_list(uploaded_blocks)
                print(f"{file_path} 已完成上传,无需恢复")
                self._clean_progress_file(file_path)
                return
            except Exception:
                # 提交失败,重新获取已提交的块列表
                block_list = await blob_client.get_block_list('committed')
                uploaded_blocks = [block.name for block in block_list.committed_blocks]
                offset = len(uploaded_blocks) * chunk_size

        # 开始/恢复上传
        async with aiofiles.open(file_path, 'rb') as f:
            await f.seek(offset)
            while True:
                # 等待暂停状态解除
                await self.pause_event.wait()

                chunk = await f.read(chunk_size)
                if not chunk:
                    break

                # 生成唯一块ID(保证顺序)
                block_id = f"block_{len(uploaded_blocks):08d}".encode('utf-8')
                # 上传单个块
                await blob_client.stage_block(block_id=block_id, data=chunk)
                uploaded_blocks.append(block_id.decode('utf-8'))
                offset += len(chunk)

                # 实时保存进度
                await self._save_progress(file_path, uploaded_blocks, offset)
                print(f"{file_path} 已上传:{offset / (1024*1024):.2f} MB")

            # 所有块上传完成,提交组装文件
            await blob_client.commit_block_list(uploaded_blocks)
            print(f"{file_path} 上传完成")
            self._clean_progress_file(file_path)

    def _clean_progress_file(self, file_path):
        """上传完成后删除进度文件"""
        progress_file = self._get_progress_file(file_path)
        try:
            os.remove(progress_file)
        except FileNotFoundError:
            pass

    async def upload_multiple_files(self, file_paths):
        """批量上传多个文件"""
        for file_path in file_paths:
            await self.upload_single_file(file_path)

使用示例

async def main():
    # 替换为你的Azure存储连接字符串和容器名
    CONNECTION_STRING = "your_azure_storage_connection_string"
    CONTAINER_NAME = "your_container_name"
    FILE_PATHS = ["file1.txt", "large_file.zip", "data.csv"]

    uploader = BlobUploader(CONNECTION_STRING, CONTAINER_NAME)

    # 启动上传任务
    upload_task = asyncio.create_task(uploader.upload_multiple_files(FILE_PATHS))

    # 模拟用户触发暂停/恢复(实际场景可绑定到API、命令行输入等)
    async def simulate_pause_resume():
        await asyncio.sleep(15)  # 上传15秒后暂停
        await uploader.pause_upload()
        print("\n=== 上传已暂停 ===")
        await asyncio.sleep(8)  # 暂停8秒后恢复
        await uploader.resume_upload()
        print("\n=== 上传已恢复 ===")

    asyncio.create_task(simulate_pause_resume())

    await upload_task

if __name__ == "__main__":
    asyncio.run(main())

注意事项

  • 块ID必须保证唯一且顺序正确,示例中用递增数字格式化的ID确保顺序。
  • 进度文件按单个文件独立生成,避免多文件上传时互相干扰。
  • 若程序意外崩溃,下次启动会自动读取进度文件,从断点处恢复上传。
  • 需确保Azure Blob容器具有Write权限,允许分块上传操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 17:05:21