Azure Function(v3)服务总线主题触发器集成Polly熔断重试策略咨询
解答
1. 是否可在该场景下使用熔断重试策略?是否有替代方案?
完全可以在这个场景下使用Polly的熔断重试策略。Cosmos DB抛出的限流RequestRateTooLargeException、临时网络故障等属于可重试的瞬时异常,熔断重试能有效提升系统容错性:
- 重试:针对临时故障自动重试,避免单次失败导致消息处理失败
- 熔断:当Cosmos DB持续异常时,暂时停止请求,避免压垮下游服务,待恢复后再恢复调用
替代方案:
- Azure Cosmos DB SDK自带重试:默认处理限流类异常,但缺乏熔断能力,无法应对持续故障
- Azure Function内置重试:通过
host.json配置Service Bus触发的重试规则,但粒度较粗,无法针对Cosmos DB特定异常做精细化控制,也无熔断能力
2. 如何将Polly熔断重试集成到Azure Function中?
推荐将Polly策略封装为服务注入,在Repository调用层应用策略,具体步骤如下:
步骤1:安装Polly包
在项目中安装Polly NuGet包:
dotnet add package Polly
步骤2:在Startup中注册Polly策略
修改Startup.cs,定义并注册针对Cosmos DB异常的熔断+重试组合策略:
public override void Configure(IFunctionsHostBuilder builder) { builder.Services.AddLogging(); builder.Services.AddCosmosRepository(options => { options.SerializationOptions = new Microsoft.Azure.CosmosRepository.Options.RepositorySerializationOptions { PropertyNamingPolicy = CosmosPropertyNamingPolicy.Default }; }); // 定义针对Cosmos DB瞬时异常的重试策略 var retryPolicy = Policy .Handle<CosmosException>(ex => ex.StatusCode == HttpStatusCode.TooManyRequests || ex.StatusCode == HttpStatusCode.ServiceUnavailable || ex.StatusCode == HttpStatusCode.GatewayTimeout) .WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromSeconds(Math.Pow(2, retryAttempt)), (exception, timeSpan, retryCount, context) => { var logger = context.GetLogger<CardGroupEventProcessor>(); logger.LogWarning("Cosmos DB调用失败,第{retryCount}次重试,延迟{timeSpan}秒。异常:{message}", retryCount, timeSpan.TotalSeconds, exception.Message); }); // 定义熔断策略 var circuitBreakerPolicy = Policy .Handle<CosmosException>(ex => ex.StatusCode == HttpStatusCode.TooManyRequests || ex.StatusCode == HttpStatusCode.ServiceUnavailable) .CircuitBreakerAsync( exceptionsAllowedBeforeBreaking: 5, durationOfBreak: TimeSpan.FromMinutes(1), onBreak: (exception, breakDuration) => { var logger = builder.Services.BuildServiceProvider().GetRequiredService<ILogger<CardGroupEventProcessor>>(); logger.LogError("Cosmos DB调用失败次数过多,触发熔断,暂停{breakDuration}分钟。异常:{message}", breakDuration.TotalMinutes, exception.Message); }, onReset: () => { var logger = builder.Services.BuildServiceProvider().GetRequiredService<ILogger<CardGroupEventProcessor>>(); logger.LogInformation("Cosmos DB熔断已恢复"); }, onHalfOpen: () => { var logger = builder.Services.BuildServiceProvider().GetRequiredService<ILogger<CardGroupEventProcessor>>(); logger.LogInformation("Cosmos DB熔断进入半开状态,尝试恢复调用"); }); // 组合重试和熔断策略 var combinedPolicy = Policy.WrapAsync(retryPolicy, circuitBreakerPolicy); builder.Services.AddSingleton<IAsyncPolicy>(combinedPolicy); builder.Services.AddSingleton<IMessageBusFactory, AzureServiceBusFactory>(); builder.Services.AddTransient(typeof(ICardGroupEventProcessor), typeof(CardGroupEventProcessor)); builder.Services.AddTransient(typeof(IUINotificationPublisher), typeof(UINotificationPublisher)); builder.Services.AddTransient(typeof(IEmailNotificationPublisher), typeof(EmailNotificationPublisher)); }
步骤3:在Processor中应用Polly策略
修改CardGroupEventProcessor,注入Polly策略并在Cosmos DB调用时使用:
private readonly IAsyncPolicy _cosmosPolicy; // 修改构造函数,注入IAsyncPolicy public CardGroupEventProcessor(ILogger<CardGroupEventProcessor> logger, IEmailNotificationPublisher emailNotificationPublisher, IUINotificationPublisher uiNotificationPublisher, IRepositoryFactory factory, IAsyncPolicy cosmosPolicy) { _logger = logger; _emailNotificationPublisher = emailNotificationPublisher; _uiNotificationPublisher = uiNotificationPublisher; _accountSubscriptionRespository = factory.RepositoryOf<AccountSubscriptions>(); _userSubscriptionrepository = factory.RepositoryOf<UserSubscriptions>(); _cosmosPolicy = cosmosPolicy; } // 修改GetAccountSubscriptionsAsync方法 private async Task<IEnumerable<AccountSubscriptions>> GetAccountSubscriptionsAsync(CardGroupEvent @event) { return await _cosmosPolicy.ExecuteAsync(async () => { var result = await _accountSubscriptionRespository .GetByQueryAsync($"SELECT * FROM asub WHERE asub.Account ='{@event.AccountNumber}' " + $"AND asub.Payer = '{@event.PayerNumber}' " + $"AND asub.Active=1"); return result; }, new Context($"GetAccountSubscriptions-{@event.RequestID}", new Dictionary<string, object> { { "Logger", _logger } })); } // 修改GetUserSubscriptionsAsync方法 private async Task<IEnumerable<UserSubscriptions>> GetUserSubscriptionsAsync(string UserUUIDs, EventType eventType) { return await _cosmosPolicy.ExecuteAsync(async () => { var result = await _userSubscriptionrepository .GetByQueryAsync($"SELECT * FROM usub WHERE usub.Type='Subscription' " + $"AND usub.UUID IN('{UserUUIDs}') " + $"AND ARRAY_CONTAINS(usub.Events,{{'EventType':{(int)eventType}}},true)"); return result; }, new Context($"GetUserSubscriptions-{eventType}", new Dictionary<string, object> { { "Logger", _logger } })); }
步骤4:优化Function异常处理(可选)
在Function的RunAsync方法中,针对熔断异常单独处理:
catch (BrokenCircuitException ex) { _logger.LogError("Cosmos DB处于熔断状态,无法处理消息:{message}", ex.Message); // 抛出异常让Service Bus将消息放入死信队列,或根据业务需求调整 throw; } catch (Exception ex) { _logger.LogError($"处理卡片组事件失败:{sbMsg.AsJson()},异常信息:{ex.Message},堆栈:{ex.StackTrace}"); throw; }
内容的提问来源于stack exchange,提问作者Rakesh Kumar
相关产品推荐
相关产品推荐

