Azure Durable Functions处理大任务量时执行时间剧增,传递大文件列表即使实际处理量少仍性能下降
看起来你遇到了Durable Functions里一个挺头疼的性能怪问题——明明实际处理的文件量没变,只是Orchestrator手里攥着的文件名列表从1k变成5k,执行时间就直接炸了,哪怕你只让活动函数处理其中1k个文件也没用。结合你给出的所有细节(本地/Azure都慢、文件全一致、无内存问题、Flex Consumption计划),我来帮你拆解核心原因,以及不用大改代码的优化方向。
核心原因:Orchestrator的状态序列化/重放开销
这几乎可以确定是Durable Functions的重放机制+大状态序列化导致的问题,和你实际处理多少文件无关,只和Orchestrator持有的状态大小有关:
- Durable Functions的Orchestrator是基于重放执行的,它会把整个执行状态(包括你手里的5k文件名列表)序列化后存储到Azure Storage(本地跑用Emulator)。每次遇到
await、需要checkpoint,或者因为环境重启/缩放被重触发时,都要把整个状态反序列化加载进来。 - 你说5k个文件名每个不到50字符,算下来序列化后的实际大小(含结构信息)远不止表面的250KB,而且Orchestrator的重放不是只做一次——每次执行到await点都会触发状态序列化/存储,每次重放都会反序列化整个列表。哪怕你实际没用到全部列表,只要它在Orchestrator的状态里,这个开销就跑不掉。
- 本地跑也慢,说明这不是Azure环境的问题,是Durable Functions框架本身的状态处理逻辑导致的CPU/IO开销。
不用大重构的优化方案
完全不用推翻现有代码,只要把大状态数据从Orchestrator移到外部存储,就能解决这个问题,以下是具体可落地的步骤:
1. 把文件名列表从Orchestrator状态移到外部存储(Blob Storage优先)
不要让Orchestrator直接持有5k文件名列表,而是把列表存在Blob Storage里,Orchestrator只存一个指向这个列表的“指针”(比如Blob的存储容器+Blob名称)。这样Orchestrator的状态里只有几个短字符串,序列化成本直接降到可以忽略的程度。
- 改动点:
- 原来的
GetExcelTemplate活动函数,不要返回List<string>,而是把文件名列表写到Blob里,返回Blob的访问键/路径。 - Orchestrator里不再持有整个列表,而是通过Blob路径,让活动函数自己去读取需要处理的文件范围。
- 原来的
2. 让活动函数按需读取文件列表,而非Orchestrator拆分
原来你在Orchestrator里用GetRange拆分大列表,现在改成Orchestrator给活动函数传起始索引+批量大小,活动函数自己从Blob里读取对应范围的文件名。这样Orchestrator完全不碰大列表,状态大小和列表长度彻底无关。
示例代码改动(最小侵入式)
修改Orchestrator函数
[Function(nameof(ProcessDocuments))] public async Task<string> ProcessDocuments( [OrchestrationTrigger] TaskOrchestrationContext context, int uploadId) { bool isValid = await context.CallActivityAsync<bool>(nameof(VerifyTemplateData), uploadId); if (isValid) { // 现在GetExcelTemplate返回Blob存储的文件名列表路径,而非列表本身 string fileListBlobKey = await context.CallActivityAsync<string>(nameof(GetExcelTemplate), uploadId); // 可以在GetExcelTemplate里顺便返回总文件数,或单独调用小活动获取 int totalFileCount = await context.CallActivityAsync<int>(nameof(GetTotalFileCount), uploadId); List<Task> parallelTasks = new List<Task>(); try { int maxConcurrency = _serviceConfig.MaxProcessingConcurrency; int batchSize = Math.Max(totalFileCount / maxConcurrency, 1); for (int i = 0; i < totalFileCount; i += batchSize) { int cappedBatchSize = Math.Min(batchSize, totalFileCount - i); // 给活动函数传Blob键、起始索引、批量大小,完全不传大列表 HandleDocumentProcessingPayload payload = new HandleDocumentProcessingPayload() { UploadId = uploadId, FileListBlobKey = fileListBlobKey, StartIndex = i, BatchSize = cappedBatchSize }; parallelTasks.Add(context.CallActivityAsync(nameof(HandleDocumentProcessing), payload)); } await Task.WhenAll(parallelTasks); // 后续的SendMail等逻辑保持不变... } catch (Exception ex) { _logger.LogError(ex, "An error occurred while uploading documents for upload {uploadId}.", uploadId); } } return context.InstanceId; }
修改活动函数
// 先更新Payload类,去掉DocumentNames,换成Blob键和索引参数 public class HandleDocumentProcessingPayload { public int UploadId { get; set; } public string FileListBlobKey { get; set; } public int StartIndex { get; set; } public int BatchSize { get; set; } } [Function(nameof(HandleDocumentProcessing))] public async Task HandleDocumentProcessing([ActivityTrigger] HandleDocumentProcessingPayload payload, FunctionContext executionContext) { try { // 从Blob读取完整文件名列表,或优化为按范围读取(如果列表按行存储) List<string> allFileNames = await _blobStorageService.GetFileListFromBlobAsync(payload.FileListBlobKey); var processingBatch = allFileNames.Skip(payload.StartIndex).Take(payload.BatchSize).ToList(); foreach (string documentData in processingBatch) { await UpdateDocument(documentData, payload.UploadId); } } catch (Exception ex) { executionContext.GetLogger(nameof(HandleDocumentProcessing)).LogError(ex, "An error occurred while uploading a document batch."); } }
3. 额外小优化:减少Orchestrator的重放次数
如果你想再进一步,可以优化Orchestrator的checkpoint频率:
- 原来的Orchestrator是一次性创建所有活动任务然后
await Task.WhenAll,可以改成把活动任务分批创建,每创建10-20个就await一次小的Task.WhenAll,这样每次checkpoint的状态增量更小,进一步减少状态存储的开销。
验证方案
你可以先在本地做个小测试:
- 写个简单的控制台程序,序列化1k和5k的文件名列表,分别计时,看看两者的序列化/反序列化时间差——大概率5k的耗时是1k的4-5倍,这就能实锤是序列化开销导致的问题。
- 把5k列表存到本地文件(模拟Blob),Orchestrator只存文件路径,然后跑一次,看看执行时间是不是和1k列表的情况一致。
这样改完之后,不管你的文件名列表是1k还是10k,Orchestrator的状态大小都一样,性能就不会再因为列表长度而下降了,完全符合你“不用大重构、支持大规模”的需求。
内容来源于stack exchange

