异步Task扩展方法线程锁一致性异常排查
执行逻辑说明
有多个任务尝试执行同一作业,需满足两种场景:
- 场景一:作业正在执行时,其他任务直接跳过,无需等待
- 场景二:作业正在执行时,其他任务等待作业完成后再继续,但不会重复执行作业
需求描述
需要实现一个async Task扩展方法,达到以下执行效果:
单任务执行输出:
[Task 1] LOCKED [Task 1] WORK-START [Task 1] WORK-END [Task 1] UNLOCKED
**五任务并发(场景二)**输出:
[Task 1] LOCKED [Task 1] WORK-START [Task 2] WAIT-TO-SKIP [Task 3] WAIT-TO-SKIP [Task 4] WAIT-TO-SKIP [Task 5] WAIT-TO-SKIP [Task 1] WORK-END [Task 1] UNLOCKED [Task 2] SKIP [Task 3] SKIP [Task 4] SKIP [Task 5] SKIP
场景一执行效果:
[Task 1] LOCKED [Task 1] WORK-START [Task 2] SKIP [Task 3] SKIP [Task 4] SKIP [Task 5] SKIP [Task 1] WORK-END [Task 1] UNLOCKED
当前实现代码
扩展方法
public static async Task FirstThreadAsync(IFirstThread obj, Func<Task> action, TaskCompletionSource? waitTaskSource, string threadName = "") { if (obj.Locked) { if (waitTaskSource != null && !waitTaskSource.Task.IsCompleted) { Log.Debug(Logger, $"[{threadName}] WAIT-TO-SKIP"); await waitTaskSource.Task; } Log.Debug(Logger, $"[{threadName}] SKIP-1"); return; } var lockWasTaken = false; var temp = obj; try { if (waitTaskSource == null || waitTaskSource.Task.IsCompleted == false) { Monitor.TryEnter(temp, ref lockWasTaken); if (lockWasTaken) obj.Locked = true; } } finally { if (lockWasTaken) Monitor.Exit(temp); } if (waitTaskSource?.Task.IsCompleted == true) { Log.Debug(Logger, $"[{threadName}] SKIP-3"); return; } if (waitTaskSource != null && !lockWasTaken) { if (!waitTaskSource.Task.IsCompleted) { Log.Debug(Logger, $"[{threadName}] WAIT-TO-SKIP (LOCKED)"); await waitTaskSource.Task; } Log.Debug(Logger, $"[{threadName}] SKIP-2"); return; } Log.Debug(Logger, $"[{threadName}] LOCKED"); try { Log.Debug(Logger, $"[{threadName}] WORK-START"); await action.Invoke().ConfigureAwait(false); Log.Debug(Logger, $"[{threadName}] WORK-END"); } catch (Exception ex) { waitTaskSource?.TrySetException(ex); throw; } finally { obj.Locked = false; Log.Debug(Logger, $"[{threadName}] UNLOCKED"); waitTaskSource?.TrySetResult(); } }
接口与示例类
public interface IFirstThread { bool Locked { get; set; } } public class Example : IFirstThread { public bool Locked { get; set; } public async Task DoWorkAsync(string taskName) { for (var i = 0; i < 10; i++) { await Task.Delay(5); } } }
单元测试代码
[TestMethod] public async Task DoWorkOnce_AsyncX() { var waitTaskSource = new TaskCompletionSource(); var example = new Methods.Example(); var tasks = new List<Task> { FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task 1"), FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task 2"), FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task 3"), FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task 4"), FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task 5"), }; await Task.WhenAll(tasks); await FirstThreadAsync(example, example.DoWorkAsync, waitTaskSource, "Task End"); //代码获取并对比输出 }
遇到的并发问题
测试中偶尔出现多个任务同时获取锁并执行作业的异常,示例日志:
[Task 5] LOCKED, [Task 4] WAIT-TO-SKIP, [Task 2] LOCKED, [Task 3] WAIT-TO-SKIP, [Task 1] WAIT-TO-SKIP, [Task 5] WORK-START, [Task 2] WORK-START, [Task 2] WORK-END, [Task 5] WORK-END, [Task 5] UNLOCKED, [Task 2] UNLOCKED, [Task 1] SKIP-1, [Task 4] SKIP-1, [Task 3] SKIP-1, [Task End] LOCKED, [Task End] WORK-START, [Task End] WORK-END, [Task End] UNLOCKED
如日志所示,Task 5和Task 2同时进入锁定状态,不符合仅单任务锁定的预期。
问题分析与修复方案
问题原因
- 非原子性的锁检查与获取:最初的
if (obj.Locked)判断和后续的Monitor.TryEnter操作不是原子的。当多个任务同时通过obj.Locked == false的检查后,都尝试获取锁,此时Monitor.TryEnter可能让多个任务成功获取锁(因为获取后立即释放了Monitor,仅依赖obj.Locked标记,但这个标记的设置没有被持续保护)。 - 错误的Monitor使用方式:获取
Monitor锁后立即释放,仅用它来设置obj.Locked,无法保证后续逻辑的线程安全。释放锁后,其他线程可以立即修改obj.Locked,导致多个线程进入执行逻辑。
修复方案
重新设计锁逻辑,使用异步兼容的锁机制(SemaphoreSlim),同时将锁检查、获取和标记设置改为原子操作,简化场景分支逻辑:
优化后的扩展方法
public static async Task ExecuteOnceAsync(this IFirstThread obj, Func<Task> action, bool waitForCompletion = false, string taskName = "") { // 异步锁:确保仅一个线程进入执行区,兼容async/await场景 static readonly SemaphoreSlim _asyncLock = new SemaphoreSlim(1, 1); bool acquiredLock = false; try { if (waitForCompletion) { // 场景二:等待锁释放,再检查是否已在执行 await _asyncLock.WaitAsync(); acquiredLock = true; if (obj.Locked) { Log.Debug(Logger, $"[{taskName}] WAIT-TO-SKIP"); // 等待当前执行任务完成 while (obj.Locked) { await Task.Delay(10); } Log.Debug(Logger, $"[{taskName}] SKIP"); return; } } else { // 场景一:尝试获取锁,失败直接跳过 if (!_asyncLock.Wait(0)) { Log.Debug(Logger, $"[{taskName}] SKIP"); return; } acquiredLock = true; if (obj.Locked) { Log.Debug(Logger, $"[{taskName}] SKIP"); return; } } // 标记为锁定状态 obj.Locked = true; Log.Debug(Logger, $"[{taskName}] LOCKED"); // 执行作业 Log.Debug(Logger, $"[{taskName}] WORK-START"); await action().ConfigureAwait(false); Log.Debug(Logger, $"[{taskName}] WORK-END"); } catch (Exception ex) { Log.Error(Logger, ex, $"[{taskName}] WORK-FAILED"); throw; } finally { if (acquiredLock) { // 解锁并标记为未锁定 obj.Locked = false; Log.Debug(Logger, $"[{taskName}] UNLOCKED"); _asyncLock.Release(); } } }
关键改进点
- SemaphoreSlim异步锁:支持异步等待,避免阻塞线程,适配async/await场景。
- 原子化操作:将锁的获取和
obj.Locked的检查放在同一个锁保护块内,确保操作的原子性。 - 简化场景分支:通过
waitForCompletion参数直接区分两种场景,逻辑更清晰。 - 移除冗余依赖:无需
TaskCompletionSource,直接通过锁和状态标记实现等待逻辑,降低复杂度。
使用示例
// 场景一:直接跳过 await example.ExecuteOnceAsync(() => example.DoWorkAsync("Task 1"), waitForCompletion: false, taskName: "Task 1"); // 场景二:等待完成后跳过 await example.ExecuteOnceAsync(() => example.DoWorkAsync("Task 1"), waitForCompletion: true, taskName: "Task 1");
内容的提问来源于stack exchange,提问作者LorneCash
相关产品推荐
相关产品推荐

