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
流程:
- 计算文件 MD5:将文件流复制到内存,计算完整文件的 MD5 值
- 小文件处理:如果文件 ≤ 1MB,直接调用小文件上传委托
- 大文件分片上传:
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/2 | Version = new Version(2, 0) 使用 HTTP/2 协议 |
潜在改进点
- 分片上传限流:当前并行发起所有分片,可能导致连接数过多
- 进度回调:缺少上传进度通知机制
- 取消令牌传递:部分方法未完全支持
CancellationToken - 失败重试策略:重试间隔固定为2秒,可考虑指数退避
Debug 特性
#if DEBUG
request.Headers.Add("X-Upyun-Meta-Ttl", "1"); // 测试环境1天后自动删除
#endif
这是一个生产级的对象存储上传组件,充分考虑了性能、可靠性和资源管理。
AI 正在分析代码…
评论加载中...