using System.Buffers;
using System.Security.Cryptography;
using Dpz.Core.Entity.Base.PublicStruct;
using Dpz.Core.Infrastructure;
using Dpz.Core.Shard.Service;
using Microsoft.Extensions.Logging;
using Microsoft.IO;

namespace Dpz.Core.Service.ObjectStorage;

public delegate Task<IFileMetadata> UploadAsync(
    ICollection<string> pathToFile,
    Func<HttpContent> setContent,
    long size = 0,
    string? hash = null,
    string? contentMd5 = null,
    CancellationToken cancellationToken = default
);

/// <summary>
/// 分片断点续传上传
/// </summary>
/// <param name="logger"></param>
/// <param name="httpClient"></param>
/// <param name="upyunOperator"></param>
public class ParallelChunkBreakpointUpload<TUpyunOperator>(
    ILogger<ParallelChunkBreakpointUpload<TUpyunOperator>> logger,
    HttpClient httpClient,
    TUpyunOperator upyunOperator
)
    where TUpyunOperator : UpyunOperator
{
    private readonly int _chunkSize = 1 << 20;

    /// <summary>
    ///
    /// </summary>
    /// <param name="file"></param>
    /// <param name="smallFileUpload"></param>
    /// <param name="checkMd5"></param>
    /// <param name="cancellationToken"></param>
    /// <returns></returns>
    /// <exception cref="Exception"></exception>
    public async Task<FileAddress?> UploadFileAsync(
        CloudFile file,
        UploadAsync smallFileUpload,
        bool checkMd5 = false,
        CancellationToken cancellationToken = default
    )
    {
        await using var stream = ApplicationTools.MemoryStreamManager.GetStream();
        await file.Stream.CopyToAsync(stream, cancellationToken);
        stream.Position = 0;

        using var md5 = MD5.Create();
        var hashBytes = await md5.ComputeHashAsync(stream, cancellationToken);
        var md5Value = BitConverter.ToString(hashBytes).Replace("-", "").ToLowerInvariant();
        logger.LogInformation("file:{@FileToPath},MD5 value:{MD5Value}", file.PathToFile, md5Value);

        stream.Position = 0;

        if (stream.Length <= _chunkSize)
        {
            var result = await smallFileUpload(
                file.PathToFile,
                () => new StreamContent(stream),
                stream.Length,
                md5Value,
                cancellationToken: cancellationToken
            );
            return new FileAddress(result.Url, md5Value);
        }

        var pathToFile = string.Join("/", file.PathToFile.Select(Uri.EscapeDataString));

        var taskId = await BeginUploadAsync(pathToFile, stream.Length);

        await SplitChunkAsync(stream, pathToFile, taskId);

        await CompleteUploadAsync(pathToFile, taskId);

        if (checkMd5)
        {
            var match = await CheckMd5Async(pathToFile, md5Value);
            if (!match)
            {
                logger.LogError(
                    "MD5 value are inconsistent,local calculate MD5 value:{MD5Value},path to file:{@PathToFile}",
                    md5Value,
                    pathToFile
                );
                throw new Exception("md5 value are inconsistent");
            }
        }

        var url =
            (
                upyunOperator.Host?.LastOrDefault() == '/'
                    ? upyunOperator.Host
                    : upyunOperator.Host + "/"
            ) + pathToFile;
        return new FileAddress(url, md5Value);
    }

    private HttpRequestMessage GetRequest(string pathToFile)
    {
        var request = new HttpRequestMessage(
            HttpMethod.Put,
            $"/{upyunOperator.Bucket}/{pathToFile}"
        )
        {
            Version = new Version(2, 0),
        };
        return request;
    }

    private async Task<string> BeginUploadAsync(string pathToFile, long fileSize)
    {
        return await ApplicationTools.RetryAsync(
            async () =>
            {
                var request = GetRequest(pathToFile);
                request.Headers.Add("X-Upyun-Multi-Disorder", "true");
                request.Headers.Add("X-Upyun-Multi-Stage", "initiate");
                request.Headers.Add("X-Upyun-Multi-Length", $"{fileSize}");
                request.Headers.Add("X-Upyun-Multi-Part-Size", $"{_chunkSize}");
#if DEBUG
                request.Headers.Add("X-Upyun-Meta-Ttl", "1");
#endif
                await request.SignatureAsync(upyunOperator);
                var response = await httpClient.SendAsync(request);

                var multiId = response.Headers.GetResponseHeaderValue("X-Upyun-Multi-Uuid");
                if (!response.IsSuccessStatusCode || multiId == null)
                {
                    logger.LogError("upload fail,status code:{StatusCode}", response.StatusCode);
                    throw new BusinessException("upload file initial multi disorder fail");
                }

                return multiId;
            },
            TimeSpan.FromSeconds(2)
        );
    }

    private async Task SplitChunkAsync(
        RecyclableMemoryStream stream,
        string pathToFile,
        string taskId
    )
    {
        var uploadTasks = new List<Task<bool>>();
        var buffer = ArrayPool<byte>.Shared.Rent(_chunkSize);
        var index = 0;
        try
        {
            int bytesRead;
            while ((bytesRead = await stream.ReadAsync(buffer.AsMemory(0, _chunkSize))) > 0)
            {
                var chunkCopy = new byte[bytesRead];
                buffer.AsSpan(0, bytesRead).CopyTo(chunkCopy);
                var memory = new Memory<byte>(chunkCopy, 0, bytesRead);
                var currentIndex = index;
                var chunkUploadTask = ApplicationTools.RetryAsync(
                    async () => await UploadChunkAsync(pathToFile, taskId, currentIndex, memory),
                    TimeSpan.FromSeconds(2)
                );
                uploadTasks.Add(chunkUploadTask);
                index++;
            }

            await Task.WhenAll(uploadTasks);
        }
        catch (Exception ex)
        {
            logger.LogError(ex, "Failed to upload chunks, path to file: {PathToFile}", pathToFile);
            // 清理已上传的任务
            foreach (var task in uploadTasks.Where(task => !task.IsCompletedSuccessfully))
            {
                try
                {
                    await task;
                }
                catch
                {
                    // Ignore exceptions for cleanup tasks
                }
            }
            throw;
        }
        finally
        {
            ArrayPool<byte>.Shared.Return(buffer);
        }
    }

    private async Task<bool> UploadChunkAsync(
        string pathToFile,
        string taskId,
        int currentIndex,
        Memory<byte> chunk
    )
    {
        return await ApplicationTools.RetryAsync(
            async () =>
            {
                var request = GetRequest(pathToFile);
                request.Headers.Add("X-Upyun-Multi-Stage", "upload");
                request.Headers.Add("X-Upyun-Multi-Uuid", taskId);
                request.Headers.Add("X-Upyun-Part-Id", $"{currentIndex}");
                request.Content = new ReadOnlyMemoryContent(chunk);
                await request.SignatureAsync(upyunOperator);
                var response = await httpClient.SendAsync(request);

                var multiId = response.Headers.GetResponseHeaderValue("X-Upyun-Multi-Uuid");
                if (!response.IsSuccessStatusCode || multiId == null)
                {
                    logger.LogError("upload fail,status code:{StatusCode}", response.StatusCode);
                    throw new BusinessException("upload file initial multi disorder fail");
                }

                return true;
            },
            TimeSpan.FromSeconds(2)
        );
    }

    private async Task CompleteUploadAsync(string pathToFile, string taskId)
    {
        await ApplicationTools.RetryAsync(
            async () =>
            {
                var completeRequest = GetRequest(pathToFile);
                completeRequest.Headers.Add("X-Upyun-Multi-Stage", "complete");
                completeRequest.Headers.Add("X-Upyun-Multi-Uuid", taskId);
                await completeRequest.SignatureAsync(upyunOperator);
                var completeResponse = await httpClient.SendAsync(completeRequest);

                var multiId = completeResponse.Headers.GetResponseHeaderValue("X-Upyun-Multi-Uuid");
                if (!completeResponse.IsSuccessStatusCode || multiId == null)
                {
                    logger.LogError(
                        "upload fail,status code:{StatusCode}",
                        completeResponse.StatusCode
                    );
                    throw new BusinessException("upload file initial multi disorder fail");
                }

                var mime = completeResponse.Headers.GetResponseHeaderValue("X-Upyun-Multi-Type");
                var length = completeResponse.Headers.GetResponseHeaderValue(
                    "X-Upyun-Multi-Length"
                );

                logger.LogInformation(
                    "upyun multi upload complete,multi id:{MultiId},MIME:{MIME},length:{Length}",
                    multiId,
                    mime,
                    length
                );

                return true;
            },
            TimeSpan.FromSeconds(2)
        );
    }

    private async Task<bool> CheckMd5Async(string pathToFile, string md5Value)
    {
        var request = new HttpRequestMessage(
            HttpMethod.Head,
            $"/{upyunOperator.Bucket}/{pathToFile}"
        )
        {
            Version = new Version(2, 0),
        };
        await request.SignatureAsync(upyunOperator);
        var response = await httpClient.SendAsync(request);

        var cloudMd5 = response.Headers.GetResponseHeaderValue("X-Upyun-Meta-Multi-MD5");

        return cloudMd5 == md5Value;
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码分析:分片断点续传上传服务

概述

这是一个用于又拍云(Upyun)对象存储的分片断点续传上传服务类,实现了大文件的并行分片上传功能。

核心功能

1. 类定义与依赖注入

public class ParallelChunkBreakpointUpload<TUpyunOperator>
  • 泛型类,TUpyunOperator 必须继承自 UpyunOperator
  • 构造函数注入:日志记录器、HttpClient、又拍云操作器
  • _chunkSize = 1 << 20:分片大小为 1MB (2^20 字节)

2. 主要上传方法 UploadFileAsync

流程:

  1. 计算文件 MD5:将文件流复制到内存,计算完整文件的 MD5 值
  2. 小文件处理:如果文件 ≤ 1MB,直接调用小文件上传委托
  3. 大文件分片上传
    • BeginUploadAsync:初始化多段上传任务
    • SplitChunkAsync:分片并行上传
    • CompleteUploadAsync:完成上传
    • CheckMd5Async:可选的 MD5 校验

3. 分片上传核心方法

BeginUploadAsync - 初始化上传

request.Headers.Add("X-Upyun-Multi-Disorder", "true");  // 允许乱序上传
request.Headers.Add("X-Upyun-Multi-Stage", "initiate");  // 初始化阶段
request.Headers.Add("X-Upyun-Multi-Length", $"{fileSize}");
  • 返回任务 UUID (X-Upyun-Multi-Uuid)
  • 支持重试机制(2秒间隔)

SplitChunkAsync - 分片并行上传

var buffer = ArrayPool<byte>.Shared.Rent(_chunkSize);  // 使用内存池
  • 使用 ArrayPool 优化内存分配
  • 读取流并分块,每块最大 1MB
  • 并行发起所有分片上传任务 (Task.WhenAll)
  • 异常处理:失败时清理已上传任务

UploadChunkAsync - 单个分片上传

request.Headers.Add("X-Upyun-Part-Id", $"{currentIndex}");
request.Content = new ReadOnlyMemoryContent(chunk);
  • 每个分片带索引 ID
  • 使用 ReadOnlyMemoryContent 避免多余复制
  • 支持自动重试

CompleteUploadAsync - 完成上传

completeRequest.Headers.Add("X-Upyun-Multi-Stage", "complete");
  • 通知服务器合并分片
  • 返回文件 MIME 类型和最终大小

4. MD5 校验 CheckMd5Async

  • 通过 HEAD 请求获取云端文件的 MD5
  • 对比本地计算的 MD5 值

技术亮点

特性实现方式
内存优化使用 ArrayPool<byte>RecyclableMemoryStream
并发上传Task.WhenAll 并行发起多个分片请求
断点续传支持 X-Upyun-Multi-Disorder 乱序上传
容错机制ApplicationTools.RetryAsync 自动重试
安全认证SignatureAsync 请求签名
HTTP/2Version = new Version(2, 0) 使用 HTTP/2 协议

潜在改进点

  1. 分片上传限流:当前并行发起所有分片,可能导致连接数过多
  2. 进度回调:缺少上传进度通知机制
  3. 取消令牌传递:部分方法未完全支持 CancellationToken
  4. 失败重试策略:重试间隔固定为2秒,可考虑指数退避

Debug 特性

#if DEBUG
request.Headers.Add("X-Upyun-Meta-Ttl", "1");  // 测试环境1天后自动删除
#endif

这是一个生产级的对象存储上传组件,充分考虑了性能、可靠性和资源管理。

评论加载中...