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

Azure函数通过EventGrid监听Blob上传时重复触发事件问题

问题:EventGrid驱动的BlobTrigger多次触发同一事件

我正在把应用升级为用EventGrid监听Blob上传,替换Azure函数里旧的Blob轮询机制。但用Azure函数写入Blob时,同一个事件会被触发2-6次,实际只创建了一个Blob,希望事件只触发一次。

这个问题只有代码写入Blob时才会出现,通过Azure门户直接上传文件就没问题,我怀疑是写入Blob的方式不对,重写上传逻辑也没解决。


写入Blob的函数代码

[FunctionName("CreateRequests")]
public static async Task<IActionResult> Run(
    [HttpTrigger(AuthorizationLevel.Anonymous, "get", "post", Route = "request")] HttpRequest req,
    IBinder binder,
    ILogger log)
{
    try
    {
        string data = await new StreamReader(req.Body).ReadToEndAsync();
        string schema = req.GetQueryParameterDictionary()["schema"];
        string processType = req.GetQueryParameterDictionary()["processType"];

        var arrayOfThings = JsonConvert.DeserializeObject<List<object>>(data);
        int countOfContracts = contracts.Count();

        foreach (var item in arrayOfThings)
        {
            var blobId = Guid.NewGuid();
            BlobAttribute dynamicBlobAttribute = new BlobAttribute($"{GetAzureBlobFolder(processType)}/{schema}/instance-{Guid.NewGuid()}/{blobId}");
            dynamicBlobAttribute.Connection = "RequestStorage";

            using (var writer = binder.Bind<TextWriter>(dynamicBlobAttribute))
            {
                writer.Write(_contract);
            }
        }

        return new OkResult(); 
    }
    catch (Exception ex)
    {
        HttpError err = new HttpError(ex, false);
        return new BadRequestErrorMessageResult(err.ExceptionMessage);
    }
}

监听器代码

[FunctionName("DidAddReevalRequestToBlob")]
public async Task Run([BlobTrigger("reevaluationrequests/{name}", Source = BlobTriggerSource.EventGrid,Connection = "RequestStorage")] Stream myBlob,
    string name, ILogger log)
{
    string contents = string.Empty;
    using (var reader = new StreamReader(myBlob))
    {
        contents = reader.ReadToEnd();
    }

    // await "do other stuff"

}

相关日志

[2022-11-15T16:50:01.738Z] Executing 'DidAddReevalRequestToBlob' (Reason='New blob detected(EventGrid): azure-webjobs-hosts/timers/eth000893dvale-2057002712/FIPA.Functions.AutomatedRequest.TimerProcessNextAutomatedRequests.Run/status', Id=6c1f6544-e682-4ea2-b6eb-9bdf55ff8662)
[2022-11-15T16:50:01.740Z] Trigger Details: MessageId: 578038d0-42b8-4194-9f43-b63715e63384, DequeueCount: 1, InsertedOn: 2022-11-15T16:50:01.000+00:00, BlobCreated: 2022-11-14T20:43:09.000+00:00, BlobLastModified: 2022-11-15T16:50:00.000+00:00
[2022-11-15T16:50:01.935Z] Executed 'DidAddReevalRequestToBlob' (Failed, Id=6c1f6544-e682-4ea2-b6eb-9bdf55ff8662, Duration=362ms)
[2022-11-15T16:50:02.742Z] Executing 'DidAddReevalRequestToBlob' (Reason='New blob detected(EventGrid): azure-webjobs-hosts/timers/eth000893dvale-2057002712/FIPA.Functions.AutomatedRequest.TimerProcessNextAutomatedRequests.Run/status', Id=9320d805-eb39-414b-a752-39ef13b52faa)
[2022-11-15T16:50:02.744Z] Trigger Details: MessageId: 578038d0-42b8-4194-9f43-b63715e63384, DequeueCount: 2, InsertedOn: 2022-11-15T16:50:01.000+00:00, BlobCreated: 2022-11-14T20:43:09.000+00:00, BlobLastModified: 2022-11-15T16:50:00.000+00:00
[2022-11-15T16:50:02.911Z] Executed 'DidAddReevalRequestToBlob' (Failed, Id=9320d805-eb39-414b-a752-39ef13b52faa, Duration=330ms)
[2022-11-15T16:50:03.354Z] Executing 'DidAddReevalRequestToBlob' (Reason='New blob detected(EventGrid): azure-webjobs-hosts/timers/eth000893dvale-2057002712/FIPA.Functions.AutomatedRequest.TimerProcessNextAutomatedRequests.Run/status', Id=0b3b082c-ef47-4d3a-9b29-f3a405b7369e)
[2022-11-15T16:50:03.355Z] Trigger Details: MessageId: 578038d0-42b8-4194-9f43-b63715e63384, DequeueCount: 3, InsertedOn: 2022-11-15T16:50:01.000+00:00, BlobCreated: 2022-11-14T20:43:09.000+00:00, BlobLastModified: 2022-11-15T16:50:00.000+00:00
[2022-11-15T16:50:03.529Z] Executed 'DidAddReevalRequestToBlob' (Failed, Id=0b3b082c-ef47-4d3a-9b29-f3a405b7369e, Duration=330ms)
[2022-11-15T16:50:03.966Z] Executing 'DidAddReevalRequestToBlob' (Reason='New blob detected(EventGrid): azure-webjobs-hosts/timers/eth000893dvale-2057002712/FIPA.Functions.AutomatedRequest.TimerProcessNextAutomatedRequests.Run/status', Id=347c98c0-89e9-4a5d-b895-f63448bddfa0)
[2022-11-15T16:50:03.967Z] Trigger Details: MessageId: 578038d0-42b8-4194-9f43-b63715e63384, DequeueCount: 4, InsertedOn: 2022-11-15T16:50:01.000+00:00, BlobCreated: 2022-11-14T20:43:09.000+00:00, BlobLastModified: 2022-11-15T16:50:00.000+00:00
[2022-11-15T16:50:04.132Z] Executed 'DidAddReevalRequestToBlob' (Failed, Id=347c98c0-89e9-4a5d-b895-f63448bddfa0, Duration=340ms)
[2022-11-15T16:50:04.592Z] Executing 'DidAddReevalRequestToBlob' (Reason='New blob detected(EventGrid): azure-webjobs-hosts/timers/eth000893dvale-2057002712/FIPA.Functions.AutomatedRequest.TimerProcessNextAutomatedRequests.Run/status', Id=92086bb8-8723-42c3-abe8-ef68f551b72e)
[2022-11-15T16:50:04.594Z] Trigger Details: MessageId: 578038d0-42b8-4194-9f43-b63715e63384, DequeueCount: 5, InsertedOn: 2022-11-15T16:50:01.000+00:00, BlobCreated: 2022-11-14T20:43:09.000+00:00, BlobLastModified: 2022-11-15T16:50:00.000+00:00
[2022-11-15T16:50:04.757Z] Executed 'DidAddReevalRequestToBlob' (Failed, Id=92086bb8-8723-42c3-abe8-ef68f551b72e, Duration=311ms)
[2022-11-15T16:50:04.761Z] Message has reached MaxDequeueCount of 5. Moving message to queue 'webjobs-blobtrigger-poison'.

问题分析

从日志能看出两个核心问题:

  1. 监听函数DidAddReevalRequestToBlob每次执行都失败,导致消息被重复重试(DequeueCount从1递增到5),最终进入死信队列。
  2. 触发函数的是WebJobs内部的计时器状态Blob(azure-webjobs-hosts/timers/.../status),而非业务的reevaluationrequests容器下的Blob,说明EventGrid订阅或BlobTrigger过滤规则存在问题。

解决方案

1. 修复监听函数执行失败问题

你的监听函数标记为async,但代码中没有实际的await操作(注释的异步逻辑必须真正await)。如果函数未正确等待异步操作完成,运行时会判定执行失败,触发重试机制。

修改后的监听函数示例:

[FunctionName("DidAddReevalRequestToBlob")]
public async Task Run(
    [BlobTrigger("reevaluationrequests/{name}", Source = BlobTriggerSource.EventGrid, Connection = "RequestStorage")] Stream myBlob,
    string name, 
    ILogger log)
{
    string contents = string.Empty;
    using (var reader = new StreamReader(myBlob))
    {
        contents = await reader.ReadToEndAsync(); // 使用异步读取方法
    }

    // 确保所有异步业务逻辑都被await
    await ProcessBlobContentAsync(contents, name, log);
}

// 示例异步业务处理方法
private async Task ProcessBlobContentAsync(string contents, string blobName, ILogger log)
{
    log.LogInformation($"开始处理Blob: {blobName}");
    // 替换为你的实际业务逻辑
    await Task.Delay(100);
    log.LogInformation($"Blob处理完成: {blobName}");
}

2. 调整EventGrid订阅的过滤规则

当前EventGrid订阅可能包含了整个存储账户的事件,导致WebJobs内部Blob的事件也被发送到函数。需要修改订阅,只监听reevaluationrequests容器下的BlobCreated事件:

  • 进入Azure门户的目标存储账户,找到对应的EventGrid订阅
  • 在筛选器设置中添加:
    • 主题后缀匹配:/blobServices/default/containers/reevaluationrequests
    • 事件类型仅选择:Microsoft.Storage.BlobCreated

3. 优化Blob写入逻辑,确保原子性

当前用TextWriter通过IBinder写入Blob的方式,可能导致Blob被多次修改(比如先创建空Blob再写入内容),触发多次BlobCreated或BlobModified事件。建议使用Azure.Storage.Blobs SDK直接上传,确保Blob原子创建:

修改CreateRequests函数的Blob写入部分:

// 通过构造函数注入BlobServiceClient
private readonly BlobServiceClient _blobServiceClient;

public CreateRequestsFunction(BlobServiceClient blobServiceClient)
{
    _blobServiceClient = blobServiceClient;
}

// 在Run方法中替换原写入逻辑
foreach (var item in arrayOfThings)
{
    var blobId = Guid.NewGuid();
    var containerName = GetAzureBlobFolder(processType); // 确保返回"reevaluationrequests"
    var blobPath = $"{schema}/instance-{Guid.NewGuid()}/{blobId}";
    
    var blobContainerClient = _blobServiceClient.GetBlobContainerClient(containerName);
    await blobContainerClient.CreateIfNotExistsAsync();
    
    var blobClient = blobContainerClient.GetBlobClient(blobPath);
    // 原子上传字符串,避免多次修改Blob
    await blobClient.UploadAsync(BinaryData.FromString(_contract), overwrite: false);
}

这种方式会一次性完成Blob创建,只会触发一次BlobCreated事件,避免因多次修改导致的重复事件。

4. 添加幂等处理(可选)

即使解决了上述问题,分布式系统仍可能出现重复事件,建议在业务逻辑中添加幂等处理:

  • 用Blob的name或内容中的唯一标识作为去重键
  • 处理前检查该Blob是否已被处理(比如记录到数据库或Blob元数据中)

内容的提问来源于stack exchange,提问作者dvalentine314

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 17:50:35