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

大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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:29:51