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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:30:11