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

如何在DelegatingHandler中访问Durable Entity并解决DI注入问题?

问题描述

在Durable Functions编排任务中,需要对外部API请求实现全局限流,尝试通过Durable Entity跟踪特定时间范围内的剩余调用次数,但在AuthenticationHandler(自定义DelegatingHandler)中注入DurableTaskClient时应用崩溃,无法正确访问Durable Entity。

提供的代码

AuthenticationHandler 代码

public class AuthenticationHandler : DelegatingHandler
{
    private readonly IKeyVaultService _keyVaultService;
    private readonly ILogger<AuthenticationHandler> _logger;
    private readonly DurableTaskClient _durableClient;

    public AuthenticationHandler(IKeyVaultService keyVaultService, ILogger<AuthenticationHandler> logger, DurableTaskClient durableClient)
    {
        _keyVaultService = keyVaultService ?? throw new ArgumentNullException(nameof(keyVaultService));
        _logger = logger;
        _durableClient = durableClient;
    }
    protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken)
    {       
        string token = await _keyVaultService.GetAccessTokenAsync();
        request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", token);

        var entityId = new EntityInstanceId(nameof(RateLimiterEntity), "RateLimiter");
        EntityMetadata<RateLimiterEntity>? entity = await _durableClient.Entities.GetEntityAsync<RateLimiterEntity>(entityId, cancellation: cancellationToken);

        _logger.LogInformation("Remaining requests: {rem}", entity.State.RemainingRequests.ToString());

        if (entity != null && entity.State.RemainingRequests <= 25)
        {
            var delay = entity.State.ResetTime - DateTime.UtcNow;
            if (delay > TimeSpan.Zero)
            {
                _logger.LogWarning("Rate limit reached. Delaying for {Delay} seconds", delay.TotalSeconds);
                await Task.Delay(delay, cancellationToken);
            }
            await _durableClient.Entities.SignalEntityAsync(entityId, nameof(RateLimiterEntity.Reset), cancellationToken);
        }
        await _durableClient.Entities.SignalEntityAsync(entityId, nameof(RateLimiterEntity.Decrement), cancellationToken);
    
        var response = await base.SendAsync(request, cancellationToken);
        return response;
    }
}

RateLimiterEntity 代码(.NET 8孤立模式)

[DurableTask(nameof(RateLimiterEntity))]
public class RateLimiterEntity
{
    public int RemainingRequests { get; set; } = 100;
    public DateTime ResetTime { get; set; } = DateTime.UtcNow.AddSeconds(60);

    public void Reset()
    {
        RemainingRequests = 100;
        ResetTime = DateTime.UtcNow.AddSeconds(60);
    }

    public void Decrement()
    {
        RemainingRequests--;
    }

    [Function(nameof(RateLimiterEntity))]
    public static Task Run([EntityTrigger] TaskEntityDispatcher dispatcher)
        => dispatcher.DispatchAsync<RateLimiterEntity>();
}

核心问题:如何正确注册DurableTaskClient并在DelegatingHandler中访问该实体?


解决方案

1. 正确注册服务(Program.cs)

在.NET 8孤立模式下,必须显式注册DurableTaskClient和自定义DelegatingHandler,同时配置HttpClient关联该Handler。修改Program.cs的服务配置逻辑:

var host = new HostBuilder()
    .ConfigureFunctionsWebApplication()
    .ConfigureServices(services =>
    {
        // 注册Durable Task客户端,绑定Azure存储配置
        services.AddDurableTaskClient(clientBuilder =>
        {
            clientBuilder.UseAzureStorage(Environment.GetEnvironmentVariable("AzureWebJobsStorage"));
        });

        // 注册自定义AuthenticationHandler
        services.AddTransient<AuthenticationHandler>();

        // 配置HttpClient并添加限流Handler
        services.AddHttpClient("ExternalApiClient")
            .AddHttpMessageHandler<AuthenticationHandler>();

        // 注册KeyVault服务(根据你的实现调整)
        services.AddScoped<IKeyVaultService, KeyVaultService>();
    })
    .Build();

host.Run();

2. 修正Handler中的实体访问逻辑

GetEntityAsync可能返回null(实体未初始化),且直接读取状态后修改存在并发风险,需调整逻辑确保实体状态正确初始化,所有状态变更通过信号(Signal)发送:

protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken)
{       
    string token = await _keyVaultService.GetAccessTokenAsync();
    request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", token);

    var entityId = new EntityInstanceId(nameof(RateLimiterEntity), "RateLimiter");

    // 初始化实体(如果不存在)
    var entity = await _durableClient.Entities.GetEntityAsync<RateLimiterEntity>(entityId, cancellation: cancellationToken);
    if (entity == null)
    {
        await _durableClient.Entities.SignalEntityAsync(entityId, nameof(RateLimiterEntity.Reset), cancellationToken);
        entity = await _durableClient.Entities.GetEntityAsync<RateLimiterEntity>(entityId, cancellation: cancellationToken);
    }

    _logger.LogInformation("Remaining requests: {rem}", entity.State.RemainingRequests);

    // 处理限流逻辑
    if (entity.State.RemainingRequests <= 25)
    {
        var delay = entity.State.ResetTime - DateTime.UtcNow;
        if (delay > TimeSpan.Zero)
        {
            _logger.LogWarning("Rate limit reached. Delaying for {Delay} seconds", delay.TotalSeconds);
            await Task.Delay(delay, cancellationToken);
            // 延迟后重置计数
            await _durableClient.Entities.SignalEntityAsync(entityId, nameof(RateLimiterEntity.Reset), cancellationToken);
        }
    }

    // 递减请求计数
    await _durableClient.Entities.SignalEntityAsync(entityId, nameof(RateLimiterEntity.Decrement), cancellationToken);
    
    var response = await base.SendAsync(request, cancellationToken);
    return response;
}

3. 关键注意事项

  • 并发安全:Durable Entity是单线程处理命令的,所有状态修改必须通过SignalEntityAsync发送指令,禁止直接修改本地读取的状态,避免并发更新导致的数据不一致。
  • 实体初始化:首次访问实体时可能未创建,需先发送Reset信号完成初始化,否则访问entity.State会引发异常。
  • 依赖注入验证:确保DurableTaskClient、IKeyVaultService等服务都已正确注册,避免注入失败导致应用崩溃。

内容的提问来源于stack exchange,提问作者Dorian-B

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 15:28:14