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

自制Simple JobScheduler偶现死锁?求问题分析与规范实现

自定义JobScheduler死锁问题分析与常规实现说明

问题背景

为学习目的编写的JobScheduler核心代码如下:

internal long ItemCount;  // 待处理的作业数量
internal ManualResetEventSlim Event { get; set; }  // 通知工作线程有新作业的事件
internal ConcurrentQueue<JobMeta> Jobs { get; set; }  // 作业队列

private void Loop(CancellationToken token) {
    
    Loop:

    // 若收到取消请求则退出
    if (token.IsCancellationRequested) return;
    
    // 线程等待,事件触发表示有新作业到达
    Event.Wait(token);
    if (Jobs.TryDequeue(out var jobMeta)) {  // 并发场景下每次出队一个作业

        // 当没有更多作业时让其他线程等待
        if(Interlocked.Decrement(ref ItemCount) == 0) Event.Reset();

        jobMeta.Job.Execute();      // 执行作业
        jobMeta.JobHandle.Set(); // 通知主线程作业完成
    }
    goto Loop;
}

// 通知线程有新作业到达
public void NotifyThreads() {
    Interlocked.Exchange(ref ItemCount, Jobs.Count);  // 设置ItemCount
    Event.Set();  // 触发通知
}

// 入队新作业
public JobHandle Schedule(IJob job) {

    var handle = new ManualResetEvent(false);
    var jobMeta = new JobMeta{ JobHandle = handle, Job = job};
    Jobs.Enqueue(jobMeta);
    return handle;
}

在以下调用场景中偶现死锁:

var jobHandle = threadScheduler.Schedule(myJob);  // JobHandle是ManualResetEvent类型
threadScheduler.NotifyThreads();

for(var index = 0; index < 10000; index++){
   
   var otherJobHandle = threadScheduler.Schedule(otherJob);
   threadScheduler.NotifyThreads();
   otherJobHandle.Wait();
}

jobHandle.Wait();  // 偶尔会发生死锁... 

一、死锁原因分析

死锁的核心触发逻辑是:主线程在循环中反复等待otherJobHandle完成,而第一个提交的myJob可能因工作线程被错误挂起,始终未被执行,最终主线程等待jobHandle.Wait()时,没有任何线程会触发该事件,导致永久阻塞。

具体典型时序:

  1. 主线程提交myJob并调用NotifyThreads(),ItemCount设为1,Event触发,某工作线程被唤醒并出队myJob。
  2. 工作线程执行Interlocked.Decrement(ref ItemCount)后,ItemCount变为0,随即调用Event.Reset()。
  3. 主线程进入循环,提交otherJob并调用NotifyThreads(),ItemCount设为1,Event触发,工作线程被唤醒处理otherJob,完成后触发otherJobHandle,主线程继续循环。
  4. 某次循环中,主线程提交otherJob后调用NotifyThreads()时,Jobs.Count可能因工作线程已出队变为0,导致ItemCount被设为0;或工作线程处理otherJob时,Decrement后ItemCount为0,触发Event.Reset()。
  5. 此时myJob可能还未执行(比如工作线程被调度切换,或Event.Reset导致工作线程重新进入等待),而Event处于重置状态,后续无NotifyThreads()能唤醒工作线程处理myJob。
  6. 主线程循环结束后等待jobHandle.Wait(),但myJob永远不会被执行,jobHandle永远不会被触发,最终死锁。

二、代码存在的逻辑问题

  • 计数与队列状态不同步:NotifyThreads()中用Interlocked.Exchange(ref ItemCount, Jobs.Count)将计数设为队列当前长度,但Jobs.Count在多线程场景下不是原子值,读取Count后可能已有工作线程出队作业,导致ItemCount与实际待处理数不符。后续Decrement到0时错误重置事件,让工作线程提前进入等待。
  • 事件重置时机错误:单次作业出队后就判断计数是否为0并重置事件,忽略了队列中可能还有其他作业的情况(比如主线程在工作线程处理当前作业时又提交新作业)。此时重置事件会导致后续作业无法被处理,工作线程会卡在Event.Wait()。
  • 重复调用NotifyThreads()导致计数混乱:循环中每次提交作业都调用NotifyThreads(),会覆盖之前的ItemCount值。比如第一个作业的计数会被后续作业覆盖,后续作业处理完成后计数递减到0就会重置事件,导致第一个作业被“遗忘”在队列中无人处理。
  • 资源泄漏与低效:每次Schedule都创建ManualResetEvent这类内核对象,未手动释放会造成资源泄漏;且ManualResetEvent是重量级同步对象,性能远不如轻量的TaskCompletionSource。
  • 循环实现不规范:使用goto实现循环,代码可读性和可维护性差,且未正确处理Event.Wait()可能抛出的异常(如OperationCanceledException)。

三、常规JobScheduler的实现方式

1. 基于内置线程池封装

直接利用CLR提供的ThreadPool或Task框架,无需手动维护工作线程:

public Task Schedule(IJob job)
{
    return Task.Run(() => job.Execute());
}

这种方式依赖CLR的线程池管理,自动处理线程的创建、复用和销毁,避免手动维护的复杂度。

2. 轻量同步原语替代ManualResetEvent

使用TaskCompletionSource替代ManualResetEvent,更符合现代异步编程模型,且资源开销更低:

public Task Schedule(IJob job)
{
    var tcs = new TaskCompletionSource();
    _ = Task.Run(() => {
        try
        {
            job.Execute();
            tcs.SetResult();
        }
        catch (Exception ex)
        {
            tcs.SetException(ex);
        }
    });
    return tcs.Task;
}

3. 基于BlockingCollection的生产者-消费者模型

BlockingCollection内置阻塞取元素的逻辑,无需手动维护事件和计数:

private readonly BlockingCollection<JobMeta> _jobs = new BlockingCollection<JobMeta>();

public Task Schedule(IJob job)
{
    var tcs = new TaskCompletionSource();
    _jobs.Add(new JobMeta { Job = job, Tcs = tcs });
    return tcs.Task;
}

// 工作线程循环
private void WorkerLoop(CancellationToken token)
{
    foreach (var jobMeta in _jobs.GetConsumingEnumerable(token))
    {
        try
        {
            jobMeta.Job.Execute();
            jobMeta.Tcs.SetResult();
        }
        catch (Exception ex)
        {
            jobMeta.Tcs.SetException(ex);
        }
    }
}

GetConsumingEnumerable会在队列空时自动阻塞,有新元素时自动唤醒,无需手动处理事件。

4. 正确的线程退出与异常处理

使用CancellationToken配合规范的循环结构,处理线程退出和执行过程中的异常,避免资源泄漏和无限等待。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 04:15:40