如何在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
相关产品推荐
相关产品推荐

