Azure Durable Function读取Blob后删除出现数据缺失问题排查
问题分析:Azure Durable Function删除Blob导致读取数据不完整
问题描述
我有一个Azure Durable Function,需要读取“blobs”路径下的若干Blob,上传新Blob完成后,删除“blobs”路径下的Blob。当前代码在调用UploadFunction后调用DeletFunction,但启用DeletFunction时,读取到的Blob数量减少,上传的新Blob中对象数量不足;注释DeletFunction则能正常读取所有Blob并上传完整数据,请问我忽略了什么?
相关代码
filePath = await context.CallActivityAsync<string>(nameof(UpoadFunction)); await context.CallActivityAsync<string>(nameof(DeletFunction)); // UploadFunction calls SendAsync() public async Task<string> SendAsync() { string prefix = "blobs"; var blobResult = await ReadAsync(prefix); var fileName = $"new.json"; await UploadAsync(fileName, JsonSerializer.Serialize(blobResult)); return fileName; } public async Task<List<AzureADUser>> ReadAsync(string path) { var blobs = _containerClient.GetBlobsAsync(prefix: path); var blobReaderTasks = new List<Task<Stream>>(); await foreach (BlobItem blobItem in blobs) { blobReaderTasks.Add(_containerClient.GetBlobClient(blobItem.Name).OpenReadAsync()); } var blobStreams = await Task.WhenAll(blobReaderTasks); var blobDeserializationTasks = blobStreams .Select(async stream => { var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; var users = await JsonSerializer.DeserializeAsync<List<AzureADUser>>(stream, options); if (users == null) { throw new Exception("Failed to deserialize blob"); } return users; }) .ToArray(); var adUsers = await Task.WhenAll(blobDeserializationTasks); return adUsers.SelectMany(userList => userList).ToList(); } public async Task UploadAsync(string path, string content, Dictionary<string, string> metadata = null) { var blobClient = _containerClient.GetBlobClient(path); await blobClient.UploadAsync(BinaryData.FromString(content), overwrite: true); if (metadata != null && metadata.Count > 0) blobClient.SetMetadata(metadata); } // DeletFunction calls DeleteAsync() public async Task DeleteAsync(string path) { var blobItems = _containerClient.GetBlobsAsync(prefix: "blobs"); var deleteTasks = new List<Task>(); await foreach (BlobItem blobItem in blobItems) { var blobClient = _containerClient.GetBlobClient(blobItem.Name); deleteTasks.Add(blobClient.DeleteIfExistsAsync()); } await Task.WhenAll(deleteTasks); }
原因分析
问题核心在于ReadAsync方法中使用的GetBlobsAsync是异步流式枚举,而Azure Blob存储的列表操作是最终一致性的。当DeletFunction在UploadFunction执行完成后立即启动删除操作时,ReadAsync里的流式枚举可能还没完全遍历完所有Blob,删除操作就已经开始修改Blob列表,导致枚举过程中丢失部分Blob项。
另外,GetBlobsAsync默认的分页机制会因为删除操作的干扰,导致后续分页无法正确获取剩余的Blob,最终读取到的数量比实际存在的少。
解决方案
1. 先缓存所有Blob信息再执行读取
修改ReadAsync方法,先把所有BlobItem一次性读取到内存列表中,避免流式枚举过程中Blob列表被删除操作干扰:
public async Task<List<AzureADUser>> ReadAsync(string path) { // 先将所有BlobItem缓存到内存列表,避免流式枚举时被删除操作影响 var blobs = await _containerClient.GetBlobsAsync(prefix: path).ToListAsync(); var blobReaderTasks = new List<Task<Stream>>(); foreach (BlobItem blobItem in blobs) { blobReaderTasks.Add(_containerClient.GetBlobClient(blobItem.Name).OpenReadAsync()); } var blobStreams = await Task.WhenAll(blobReaderTasks); var blobDeserializationTasks = blobStreams .Select(async stream => { var options = new JsonSerializerOptions { PropertyNameCaseInsensitive = true }; var users = await JsonSerializer.DeserializeAsync<List<AzureADUser>>(stream, options); if (users == null) { throw new Exception("Failed to deserialize blob"); } return users; }) .ToArray(); var adUsers = await Task.WhenAll(blobDeserializationTasks); return adUsers.SelectMany(userList => userList).ToList(); }
2. 修正DeleteAsync的参数冗余问题
DeleteAsync方法中传入的path参数未被使用,直接硬编码了prefix: "blobs",可以改成使用传入参数,保持代码一致性:
public async Task DeleteAsync(string path) { var blobItems = _containerClient.GetBlobsAsync(prefix: path); var deleteTasks = new List<Task>(); await foreach (BlobItem blobItem in blobItems) { var blobClient = _containerClient.GetBlobClient(blobItem.Name); deleteTasks.Add(blobClient.DeleteIfExistsAsync()); } await Task.WhenAll(deleteTasks); }
内容的提问来源于stack exchange,提问作者user989988
相关产品推荐
相关产品推荐

