Azure Durable Function结合BlobTrigger处理大文件失效问题求助
问题描述
需上传大型CSV文件并将每条记录发送至第三方处理,第三方单条记录处理耗时数分钟且有并发限制,因此采用Durable Function实现。此前运行正常,现仅小文件可正常处理,大文件无法处理,日志输出如下:
2023-08-13T14:01:29Z [Verbose] Blob 'file_15000-3' will be skipped for function 'FileUploadBlobTrigger' because this blob with ETag '"0x8DB9C053CC8E3D8"' has already been processed. PollId: '64b8b6d2-5d7a-4904-b9bc-8b2258f02ab9'. Source: 'ContainerScan'.
相关代码及配置
Blob Trigger 代码
public class FileUploadBlobTrigger { private readonly IfileService _fileService; public FileUploadBlobTrigger(IfileService fileService) { _fileService = fileService; } [FunctionName(nameof(FileUploadBlobTrigger))] public async Task Run([BlobTrigger("%My:Container%/{name}", Connection = "BlobConnectionString")] Stream myBlob, string name, ILogger log, [DurableClient] IDurableOrchestrationClient starter) { // Business logic to get dataModel await starter.StartNewAsync(nameof(BatchProcess), dataModel); } }
编排及活动函数代码
public class BatchProcess { private readonly BlobTriggerConfiguration _blobTriggerConfiguration; private readonly IFileService _fileService; private readonly ILogger<BatchProcess> _logger; public BatchProcess(IOptions<BlobTriggerConfiguration> blobTriggerConfiguration, IFileService fileService, ILogger<BatchProcess> logger) { _blobTriggerConfiguration = blobTriggerConfiguration.Value; _fileService = fileService; _logger = logger; } [FunctionName(nameof(BatchProcess))] public async Task Run( [OrchestrationTrigger] IDurableOrchestrationContext context) { foreach (var batch in batchList) { var batchModel = new MatchBatchModel { batches = batch, Name = input.Name }; var batchResponse = await context.CallActivityAsync<BatchResponses>("Batch", batchModel); // Combine with another model } } } [FunctionName("Batch")] public async Task<BatchResponses> Batch( [ActivityTrigger] BatchModel context) { // Process file with the 3rd party return batchResponse; }
host.json 配置
"extensions": { "durableTask": { "maxConcurrentActivityFunctions": 1, "maxConcurrentOrchestratorFunctions": 1 }, "blobs": { "maxDegreeOfParallelism": 1 } }
问题分析
日志提示Blob已被处理的核心原因是:Blob Trigger默认租赁超时时间(30秒)不足以完成大型文件的读取、数据模型构建及Durable编排启动逻辑。当租赁超时后,触发器会重新尝试处理Blob,但此时Blob的ETag未发生变化,内置的重复处理检测机制会判定该Blob已处理并跳过。
解决方案
1. 延长Blob Trigger租赁超时时间
在host.json中调整Blob Trigger的租赁配置,设置足够长的超时时间适配大型文件初始化:
"extensions": { "blobs": { "maxDegreeOfParallelism": 1, "leaseDuration": "00:10:00", // 根据文件大小调整,示例为10分钟 "leaseRenewInterval": "00:01:00", "maximumLeaseRenewalInterval": "00:05:00" }, "durableTask": { "maxConcurrentActivityFunctions": 1, "maxConcurrentOrchestratorFunctions": 1 } }
2. 处理完成后标记Blob状态
在Blob Trigger完成编排启动后,将Blob移动至已处理容器或修改元数据,避免重复触发:
// 在FileUploadBlobTrigger的Run方法末尾添加 var sourceBlobClient = new BlobClient(Environment.GetEnvironmentVariable("BlobConnectionString"), "%My:Container%", name); var targetBlobClient = new BlobClient(Environment.GetEnvironmentVariable("BlobConnectionString"), "processed-container", name); await targetBlobClient.StartCopyFromUriAsync(sourceBlobClient.Uri); await sourceBlobClient.DeleteAsync();
3. 优化大型CSV读取逻辑
避免一次性加载整个文件到内存,改为流式分块读取,缩短初始化时间:
// 替换原数据读取逻辑 using var reader = new StreamReader(myBlob); List<string> batchLines = new List<string>(); int batchSize = 100; // 根据实际情况调整批次大小 while (!reader.EndOfStream) { var line = await reader.ReadLineAsync(); batchLines.Add(line); if (batchLines.Count >= batchSize) { // 构建批次模型并处理 batchLines.Clear(); } } // 处理剩余行 if (batchLines.Count > 0) { // 构建最后一个批次模型 }
4. 增加编排启动日志
添加日志记录Orchestration实例ID,确认编排是否成功启动:
var instanceId = await starter.StartNewAsync(nameof(BatchProcess), dataModel); log.LogInformation("Orchestration started with ID: {InstanceId}", instanceId);
内容的提问来源于stack exchange,提问作者Nipuna
相关产品推荐
相关产品推荐

