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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 01:57:19