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'.
问题分析
从日志能看出两个核心问题:
- 监听函数
DidAddReevalRequestToBlob每次执行都失败,导致消息被重复重试(DequeueCount从1递增到5),最终进入死信队列。 - 触发函数的是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
相关产品推荐
相关产品推荐

