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

异步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同时进入锁定状态,不符合仅单任务锁定的预期。

问题分析与修复方案

问题原因

  1. 非原子性的锁检查与获取:最初的if (obj.Locked)判断和后续的Monitor.TryEnter操作不是原子的。当多个任务同时通过obj.Locked == false的检查后,都尝试获取锁,此时Monitor.TryEnter可能让多个任务成功获取锁(因为获取后立即释放了Monitor,仅依赖obj.Locked标记,但这个标记的设置没有被持续保护)。
  2. 错误的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();
        }
    }
}

关键改进点

  1. SemaphoreSlim异步锁:支持异步等待,避免阻塞线程,适配async/await场景。
  2. 原子化操作:将锁的获取和obj.Locked的检查放在同一个锁保护块内,确保操作的原子性。
  3. 简化场景分支:通过waitForCompletion参数直接区分两种场景,逻辑更清晰。
  4. 移除冗余依赖:无需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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 11:50:25