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

.NET中HttpClient场景可重入分布式互斥锁实现方案咨询

分布式可重入硬件调用锁实现方案需求

我正在寻找一套锁实现方案:默认通过分布式锁保障对单个API的所有调用互斥执行;同时支持对包含多次该API调用的更大代码块(关键流程)加锁,这类场景下锁需具备可重入特性,既不会因代码块已持有锁阻塞内部的API调用,也支持多层嵌套加锁方法的重入逻辑。

使用示例

单API自动加锁场景

默认应自动注册锁(例如在HttpMessageHandler中实现),无需手动加锁即可保证互斥:

// 默认应自动注册锁(例如在HttpMessageHandler中实现)
await _deviceClient.PerformAction();

嵌套关键流程可重入场景

外层代码块持有锁时,内部API调用、嵌套加锁逻辑都应复用同一把锁,仅在最外层锁释放时才真正解锁:

async Task CriticalProcedure()
{
   // 嵌套代码中应复用同一把锁(可重入)
   await using (await _reentrantLockProvider.AcquireLockAsync())
   {
      await _deviceClient.TriggerAction();
      await SharedCriticalProcedure();
   }
   // 仅在此处释放锁
}

async Task SharedCriticalProcedure()
{
   await using (await _customLockProvider.AcquireLockAsync())
   {
      await _deviceClient.HardReset();
      await _deviceClient.Refresh();
   }
}

并发调用强制互斥场景

即使未直接await单个任务,批量发起的设备调用也需要强制顺序执行,保证互斥:

// 即使未直接await也应强制顺序执行(互斥)
var task1 = _deviceClient.PerformAction1();
var task2 = _deviceClient.PerformAction2();
await Task.WhenAll(task1, task2);

业务背景与约束

我的团队正在开发负责调用硬件设备的WebAPI,当我们的API端点被调用时,会从请求头获取目标硬件的标识,在启动阶段用于配置对应HttpClient的baseUrl,随后向该硬件API发起一次或多次调用。当前场景存在以下约束:

  • 设备正在处理请求时不可接收其他调用(互斥性要求)
  • 部分涉及设备的关键流程(包含多次调用的代码块)不可被其他设备调用打断(可重入特性的需求来源)
  • 用户可能同时向我们的API发起多个请求,因此锁需要跨请求生效
  • 我们的WebAPI可能多实例部署,因此锁需要支持分布式场景
  • 我们使用Refit定义硬件设备的API接口并生成HttpClient

现有实现方案

我目前实现了一套自认为可用的方案,但整体较为笨重、存在过度设计问题:核心原因是HttpMessageHandler的生命周期不可预测,不与请求范围绑定,因此我需要借助TraceIdentifier和字典来实现请求生命周期内的可重入逻辑,现有实现代码如下:

服务启动配置

// 启动阶段配置
service.AddSingleton<IReentrantLockProvider, ReentrantLockProvider>();

services
   .AddHttpClient(IDeviceClient)
   .AddTypedClient(client => RestService.For<IDeviceClient>(client, refitSettings))
   .ConfigureHttpClient((provider, client) => ConfigureHardwareBaseUrl())
   .AddHttpMessageHandler<HardwareMutexMessageHandler>();

锁消息拦截处理器

public class HardwareMutexMessageHandler : DelegatingHandler
{
   private readonly IReentrantLockProvider _reentrantPanelLockProvider;
   private readonly IHttpContextAccessor _httpContextAccessor;
   private readonly ConcurrentDictionary<string, object> _locks;

   public HardwareMutexMessageHandler(IReentrantLockProvider reentrantPanelLockProvider, IHttpContextAccessor httpContextAccessor)
   {
      _reentrantPanelLockProvider = reentrantPanelLockProvider;
      _httpContextAccessor = httpContextAccessor;
   }

   protected override async Task<HttpResponseMessage> SendAsync(HttpRequestMessage request, CancellationToken cancellationToken)
   {
      await using (await _reentrantPanelLockProvider.AcquireLockAsync(cancellationToken))
      {
         var hardwareId = _httpContextAccessor.HttpContext.Request.Headers["HardwareId"];

         var mutex = _locks.GetOrAdd(hardwareId, _ => new());

         // 仅用于处理开发者批量发起调用或遗漏await调用的场景
         lock (mutex)
         {
            return base.SendAsync(request, cancellationToken).Result;
         }
      }
   }
}

可重入锁提供器

public class ReentrantLockProvider : IReentrantLockProvider
{
   private readonly IDistributedLockProvider _distributedLockProvider;
   private readonly IHttpContextAccessor _httpContextAccessor;
   private readonly ConcurrentDictionary<string, ReferenceCountedDisposable> _lockDictionary;
   private readonly object _lockVar = new();

   public ReentrantLockProvider(IDistributedLockProvider distributedLockProvider, IHttpContextAccessor httpContextAccessor)
   {
      _distributedLockProvider = distributedLockProvider;
      _httpContextAccessor = httpContextAccessor;
      _lockDictionary = new ConcurrentDictionary<string, ReferenceCountedDisposable>();
   }

   public async Task<IAsyncDisposable> AcquireLockAsync(CancellationToken cancellationToken = default)
   {
      var hardwareId = _httpContextAccessor.HttpContext.Request.Headers["HardwareId"];
      var requestId = _httpContextAccessor.HttpContext.TraceIdentifier;

      lock (_lockVar)
      {
         if (_lockDictionary.TryGetValue(requestContext.CorrelationId, out referenceCountedLock))
         {
            referenceCountedLock.RegisterReference();
            return referenceCountedLock;
         }

         acquireLockTask = _distributedLockProvider.AcquireLockAsync(hardwareId, timeout: null, cancellationToken);
            
         referenceCountedLock = new ReferenceCountedDisposable(async () =>
            await RemoveLock(acquireLockTask.Result, requestContext.CorrelationId)
         );

         _lockDictionary.TryAdd(requestContext.CorrelationId, referenceCountedLock);
      }
   }

   private async Task RemoveLock(IDistributedSynchronizationHandle acquiredLock, string correlationId)
   {
      ValueTask disposeAsyncTask;

      lock (_lockVar)
      {
         disposeAsyncTask = acquiredLock.DisposeAsync();
         _ = _lockDictionary.TryRemove(correlationId, out _);
      }
        
      await disposeAsyncTask;
   }
}

引用计数可释放对象

public class ReferenceCountedDisposable : IAsyncDisposable
{
   private readonly Func<Task> _asyncDispose;
   private int _refCount;

   public ReferenceCountedDisposable(Func<Task> asyncDispose)
   {
      _asyncDispose = asyncDispose;
      _refCount = 1;
   }

   public void RegisterReference()
   {
      Interlocked.Increment(ref _refCount);
   }

   public async ValueTask DisposeAsync()
   {
      var references = Interlocked.Decrement(ref _refCount);

      if (references == 0)
      {
         await _asyncDispose();
      }
      else if (references < 0)
      {
         throw new InvalidOperationException("Can't dispose multiple times");
      }
      else
      {
         GC.SuppressFinalize(this);
      }
   }
}

内容的提问来源于stack exchange,提问作者Expotr

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 01:21:46