大CSV文件拆分及Azure存储多线程低占用上传优化求助
大CSV文件拆分并行上传至Azure存储的性能优化
问题描述
需要将大CSV文件拆分为小文件并同步上传至Azure存储账户,核心目标是在低内存、低磁盘占用的前提下,实现CSV读写/解析与批量上传的高效并行,确保上传操作不影响原CSV的读写性能。当前实现中,原10万行CSV读写耗时45秒,加入上传逻辑后增至约60秒,并行效果未达预期。
当前实现代码
public async Task HandleEventAsync(EntityEventMessage message) { var executionGuid = Guid.NewGuid(); var uri = new Uri(message.MessageLocation); var containerName = uri.Segments[1].TrimEnd('/'); var blobName = string.Join("", uri.Segments[2..]); var entityName = message.EntityName; var blobFinalFileName = uri.Segments[^1].TrimEnd('/'); var blobServiceClient = new BlobServiceClient(uri, _tokenCredential); var blobContainerClient = blobServiceClient.GetBlobContainerClient(containerName); var sourceBlobClient = blobContainerClient.GetBlobClient(blobName); var blobExists = await sourceBlobClient.ExistsAsync(); if (!blobExists) { _logger.LogWarning("The message of the event does not exist. It might have been processed already."); throw new NoRetryException(); } var fileCounter = 1; var totalRecords = 0; SetUpBatchNames(executionGuid, blobFinalFileName, fileCounter, out var batchFileName, out var tempFilePath); var tasks = new Queue<Task>(); var processTimer = new Stopwatch(); var options = new BlobUploadOptions { TransferOptions = new StorageTransferOptions { MaximumConcurrency = 2 } }; processTimer.Start(); using var sourceBlobStream = await sourceBlobClient.OpenReadAsync(); var writer = new StreamWriter(tempFilePath, true, Encoding.UTF8); try { using var streamReader = new StreamReader(sourceBlobStream, Encoding.UTF8); using var csvParser = new CsvParser(streamReader, _csvConfiguration); var batchTimer = new Stopwatch(); batchTimer.Start(); while (await csvParser.ReadAsync()) { await writer.WriteLineAsync(csvParser.RawRecord); totalRecords++; if (totalRecords % _options.Value.RecordsPerFile == 0) { fileCounter = await PrepareBatchAndUpload(blobName, blobFinalFileName, blobContainerClient, options, fileCounter, batchFileName, tempFilePath, tasks, writer, batchTimer); SetUpBatchNames(executionGuid, blobFinalFileName, fileCounter, out batchFileName, out tempFilePath); writer = new StreamWriter(tempFilePath, true, Encoding.UTF8); batchTimer.Restart(); } } if (totalRecords % _options.Value.RecordsPerFile != 0) await PrepareBatchAndUpload(blobName, blobFinalFileName, blobContainerClient, options, fileCounter, batchFileName, tempFilePath, tasks, writer, batchTimer); await Task.WhenAll(tasks); processTimer.Stop(); _logger.LogDebug("Took '{TotalSplitTimeMs}' ms to split the original file with '{OriginalTotalRecords}' records.", processTimer.ElapsedMilliseconds, totalRecords); } finally { writer?.Dispose(); var dir = new DirectoryInfo(Path.GetTempPath()); var fileNamesToDelete = blobFinalFileName.Replace(".csv", $"_{executionGuid}_*.csv"); foreach (var file in dir.EnumerateFiles(fileNamesToDelete)) { file.Delete(); } } GC.Collect(); static void SetUpBatchNames(Guid executionGuid, string blobFinalFileName, int fileCounter, out string batchFileName, out string tempFilePath) { batchFileName = blobFinalFileName.Replace(".csv", $"_{executionGuid}_{fileCounter}.csv"); tempFilePath = $"{Path.GetTempPath()}{batchFileName}"; } async Task<int> PrepareBatchAndUpload(string blobName, string blobFinalFileName, BlobContainerClient blobContainerClient, BlobUploadOptions options, int fileCounter, string? batchFileName, string? tempFilePath, Queue<Task> tasks, StreamWriter writer, Stopwatch batchTimer) { batchTimer.Stop(); _logger.LogDebug($"Took '{batchTimer.Elapsed}' to parse the batch."); _logger.LogDebug($"Creating temp file '{tempFilePath}'."); await writer.FlushAsync(); writer.Close(); await writer.DisposeAsync(); var batchBlobName = blobName.Replace(blobFinalFileName, batchFileName); tasks.Enqueue(UploadBatchFileAsync(batchBlobName, tempFilePath, blobContainerClient, options)); fileCounter++; return fileCounter; } } private async Task UploadBatchFileAsync(string blobName, string filePath, BlobContainerClient containerClient, BlobUploadOptions options) { var blobClient = containerClient.GetBlobClient(blobName); _logger.LogDebug($"Uploading blob: '{blobClient.Uri}'."); using (var fileStream = File.OpenRead(filePath)) { await blobClient.UploadAsync(fileStream, options); } _logger.LogDebug($"Blob uploaded: '{blobClient.Uri}'."); await Task.Run(() => File.Delete(filePath)); }
已尝试的方案
- 使用
List<Task>添加上传任务并调用await Task.WhenAll - 将上传方法改为同步并通过
ThreadPool.QueueUserWorkItem实现多线程 - 参考TaskScheduler相关文档实现自定义调度
- 使用
Task.Factory.StartNew启动上传任务 - 调整
BlobUploadOptions.TransferOptions.MaximumConcurrency参数
优化方案建议
1. 控制上传并发数,避免资源竞争
当前代码会将所有上传任务加入队列等待最后统一执行,可能导致短时间内大量任务抢占磁盘、网络资源,影响CSV读写。用SemaphoreSlim限制同时运行的上传任务数量(建议3-5,根据服务器资源调整):
// 初始化信号量,限制最大并发上传数 var uploadSemaphore = new SemaphoreSlim(3); // 修改上传方法 private async Task UploadBatchFileAsync(string blobName, string filePath, BlobContainerClient containerClient, BlobUploadOptions options, SemaphoreSlim semaphore) { await semaphore.WaitAsync(); try { var blobClient = containerClient.GetBlobClient(blobName); _logger.LogDebug($"Uploading blob: '{blobClient.Uri}'."); using (var fileStream = File.OpenRead(filePath)) { await blobClient.UploadAsync(fileStream, options); } _logger.LogDebug($"Blob uploaded: '{blobClient.Uri}'."); File.Delete(filePath); // 流释放后直接删除,无需Task.Run } finally { semaphore.Release(); } }
调用时传入信号量:
tasks.Enqueue(UploadBatchFileAsync(batchBlobName, tempFilePath, blobContainerClient, options, uploadSemaphore));
2. 优化临时文件写入逻辑
- 新创建临时文件时,
StreamWriter的append参数设为false(默认值),每个批次都是新文件,无需追加模式:writer = new StreamWriter(tempFilePath, false, Encoding.UTF8); - 批量写入记录,减少磁盘IO次数:攒100条记录再一次性写入,替代逐行写:
var batchLines = new List<string>(); while (await csvParser.ReadAsync()) { batchLines.Add(csvParser.RawRecord); totalRecords++; if (batchLines.Count >= 100) { await writer.WriteAsync(string.Join(Environment.NewLine, batchLines) + Environment.NewLine); batchLines.Clear(); } if (totalRecords % _options.Value.RecordsPerFile == 0) { // 写入剩余未批量提交的行 if (batchLines.Count > 0) { await writer.WriteAsync(string.Join(Environment.NewLine, batchLines) + Environment.NewLine); batchLines.Clear(); } // 后续批次处理逻辑... } }
3. 提升Blob读取效率
设置Blob读取的缓冲区大小,增大缓冲区能减少网络IO次数:
var readOptions = new BlobOpenReadOptions(false) { BufferSize = 65536; // 64KB,可根据实际调整为128KB/256KB }; using var sourceBlobStream = await sourceBlobClient.OpenReadAsync(readOptions);
4. 移除不必要的操作
- 去掉手动调用
GC.Collect():手动触发垃圾回收会阻塞线程,CLR会自动处理内存回收。 - 临时文件在上传完成后立即删除(如优化方案1中修改的
UploadBatchFileAsync),无需等到finally阶段统一删除,减少磁盘占用同时避免文件残留风险。
5. 调整Azure上传配置
- 适当提高
MaximumConcurrency(建议4-6),需结合服务器网络带宽和CPU资源,避免过度占用; - 按需启用
TransferOptions.EnableCheckpointing,上传中断时可续传,适合大文件场景。
内容的提问来源于stack exchange,提问作者FEST
相关产品推荐
相关产品推荐

