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

Azure Blob Storage V12 BlockBlobClient分块上传实现技术咨询

我来帮你搞定这个Azure Blob分块上传的Stream包装类实现,你的思路方向是对的,但代码里有几个需要调整的地方,我给你梳理清楚并提供完整的可运行方案:

Azure Blob分块上传的Stream包装类实现

首先,你选BlockBlobClient的StageBlockAsync和CommitBlockListAsync来实现分块上传是完全正确的——这两个方法就是Azure Storage专门为大文件分块上传设计的API。我先给你修正逻辑并提供完整代码,再逐一解释关键细节:

核心实现思路

  • 用内存流缓冲写入的数据,每次写入后检查缓冲大小是否达到4MB(410241024字节)
  • 达到阈值时,将缓冲流内容作为一个分块上传,生成唯一块ID并记录
  • 上传完成后清空缓冲流,继续接收新数据
  • 最后需要主动调用提交方法,把剩余缓冲数据作为最后一块上传,再提交所有块ID列表完成Blob创建

完整可运行代码

using Azure.Storage.Blobs.Specialized;
using System;
using System.Collections.Generic;
using System.IO;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

public class BlockBlobUploadStream : Stream
{
    private readonly BlockBlobClient _blockBlobClient;
    private readonly MemoryStream _bufferStream;
    private readonly List<string> _blockIds;
    private bool _isDisposed;
    private bool _isCommitted;
    private int _blockCounter;
    private const int _maxBlockSizeInBytes = 4 * 1024 * 1024; // 4MB分块大小

    public BlockBlobUploadStream(BlockBlobClient blockBlobClient)
    {
        _blockBlobClient = blockBlobClient ?? throw new ArgumentNullException(nameof(blockBlobClient));
        _bufferStream = new MemoryStream(_maxBlockSizeInBytes); // 预分配内存提升性能
        _blockIds = new List<string>();
        _blockCounter = 0;
    }

    // 重写Stream的基础属性,只保留写入能力
    public override bool CanRead => false;
    public override bool CanSeek => false;
    public override bool CanWrite => true;
    public override long Length => throw new NotSupportedException();
    public override long Position 
    { 
        get => throw new NotSupportedException(); 
        set => throw new NotSupportedException(); 
    }

    public override void Flush() => _bufferStream.Flush();
    public override Task FlushAsync(CancellationToken cancellationToken) => _bufferStream.FlushAsync(cancellationToken);
    public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException("此流仅支持写入操作");
    public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException("不支持Seek操作");
    public override void SetLength(long value) => throw new NotSupportedException("不支持设置长度");

    public override void Write(byte[] buffer, int offset, int count)
    {
        WriteAsync(buffer, offset, count).GetAwaiter().GetResult();
    }

    public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
    {
        if (_isCommitted)
            throw new InvalidOperationException("无法向已提交的流写入数据");

        // 先将数据写入缓冲流
        await _bufferStream.WriteAsync(buffer, offset, count, cancellationToken);

        // 循环检查并上传分块(处理一次写入超过4MB的情况)
        while (_bufferStream.Length >= _maxBlockSizeInBytes)
        {
            await UploadCurrentBlockAsync(cancellationToken);
        }
    }

    /// <summary>
    /// 提交所有分块,完成Blob上传
    /// </summary>
    public async Task CommitAsync(CancellationToken cancellationToken = default)
    {
        if (_isCommitted) return;

        // 上传剩余的缓冲数据(如果存在)
        if (_bufferStream.Length > 0)
        {
            await UploadCurrentBlockAsync(cancellationToken);
        }

        // 提交块列表,Azure会按顺序拼接所有分块
        await _blockBlobClient.CommitBlockListAsync(_blockIds, cancellationToken: cancellationToken);
        _isCommitted = true;
    }

    private async Task UploadCurrentBlockAsync(CancellationToken cancellationToken)
    {
        // 将缓冲流指针重置到开头,准备读取上传
        _bufferStream.Position = 0;

        // 生成唯一块ID:用递增数字转Base64,确保每个块ID唯一且长度一致
        string blockId = Convert.ToBase64String(Encoding.UTF8.GetBytes(_blockCounter.ToString("D6")));
        _blockCounter++;

        // 上传当前分块
        await _blockBlobClient.StageBlockAsync(blockId, _bufferStream, cancellationToken: cancellationToken);
        _blockIds.Add(blockId);

        // 清空缓冲流,准备接收下一批数据
        _bufferStream.SetLength(0);
    }

    protected override void Dispose(bool disposing)
    {
        if (_isDisposed) return;

        if (disposing)
        {
            _bufferStream.Dispose();
            // 可选:如果流未提交,可清理已上传的无效块
            // await _blockBlobClient.DeleteIfExistsAsync(cancellationToken: CancellationToken.None);
        }

        _isDisposed = true;
        base.Dispose(disposing);
    }
}

关键细节解释

  1. 块ID生成逻辑:
    不能直接用写入的buffer生成Base64作为块ID——不同分块可能有重复内容,会导致块ID冲突覆盖。这里用6位格式化的递增数字转Base64,既保证唯一性,也符合Azure对块ID长度的要求。

  2. 缓冲流预分配:
    初始化MemoryStream时传入_maxBlockSizeInBytes,预分配内存可以减少内存碎片,提升写入性能。

  3. 写入循环处理:
    用while循环检查缓冲大小,避免一次写入超过4MB时漏传分块的情况。

  4. 提交操作的必要性:
    必须调用CommitAsync才能完成Blob创建——它会处理剩余缓冲数据,并告知Azure按照块ID列表的顺序拼接所有分块内容。

使用示例

// 假设你已经初始化了BlobContainerClient
var blobClient = containerClient.GetBlockBlobClient("my-large-file.dat");

// 使用自定义上传流
using var uploadStream = new BlockBlobUploadStream(blobClient);

// 模拟写入数据(可以替换为从文件/网络读取的实际数据)
var chunkData = new byte[1024 * 1024]; // 1MB测试数据
for (int i = 0; i < 10; i++) // 写入10次,总大小10MB
{
    await uploadStream.WriteAsync(chunkData, 0, chunkData.Length);
}

// 最后提交所有分块,完成上传
await uploadStream.CommitAsync();

注意事项

  • Azure Blob最多支持50000个分块,所以如果你的文件超过200GB(4MB*50000),需要调整分块大小或者使用其他上传方案。
  • 如果上传中断,已上传的分块会保留在Azure Storage中,你可以选择重试提交,或者调用DeleteIfExistsAsync清理无效块。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 21:02:48