.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
相关产品推荐
相关产品推荐

