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

线程池任务未全部执行求助:如何让线程完成所有任务后终止

线程池任务执行异常求助

我正在开发一个简单的线程池,功能是接收数字数组,让线程并行计算每个数字的阶乘。现在遇到一个问题:比如用2个线程处理6个任务时,线程完成首个任务后程序就直接终止了。我需要的是线程完成当前任务后,自动获取并执行剩余任务,直到所有任务都完成后程序才终止。


现有代码

MyThreadPool.cs

using System;
using System.Collections.Generic;
using System.Threading;

namespace ThreadPool
{
    public class MyThreadPool : IDisposable
    {
        private bool _disposed = false;
        private object _lock = new object();
        private Queue<MyTask> tasks = new();

        public MyThreadPool(int numThreads = 0)
        {
            for (int i = 0; i < numThreads; i++)
            {
                new Thread(() =>
                {
                    while (_disposed == false)
                    {
                        MyTask? task;
                        lock (_lock)
                        {
                            while (!tasks.TryDequeue(out task) && _disposed == false)
                            {
                                Monitor.Wait(_lock);
                            }
                        }
                        if (task != null)
                        {
                            Console.WriteLine(@$"Thread {Thread.CurrentThread
                                .ManagedThreadId} calculated factorial for {task
                                .getNum()}. Result: {task.CalculateFactorial(task
                                .getNum())}");
                            lock (_lock)
                            {
                                if (tasks.Count > 0)
                                {
                                    Monitor.Pulse(_lock);
                                }
                            }
                        }
                    }
                }).Start();
            }
        }

        public MyTask Enqueue(int number)
        {
            MyTask task = new(number);
            lock (_lock)
            {
                tasks.Enqueue(task);
                Monitor.Pulse(_lock);
            }
            return task;
        }

        public void Dispose()
        {
            _disposed = true;
            lock (_lock)
            {
                Monitor.PulseAll(_lock);
            }
        }
    }
}

MyTask.cs

using System;

namespace ThreadPool
{
    public class MyTask
    {
        private int num;

        public void setNum(int num)
        {
            this.num = num;
        }

        public int getNum()
        {
            return this.num;
        }

        public MyTask(int number)
        {
            this.num = number;
        }

        public int CalculateFactorial(int number)
        {
            int result = 1;
            for (int i = number; i > 0; i--)
            { 
                result = result * i;
            }
            return result;
        }
    }
}

Program.cs

int[] array = { 1, 2, 3, 4, 5, 6 };
int numThread = 3;
MyThreadPool threadPool = new(numThread);
for (int i = 0; i < array.Length; i++)
{
    threadPool.Enqueue(array[i]);
}

threadPool.Dispose();

问题原因

  1. 主线程提前触发销毁:Program.cs中刚完成任务入队就立刻调用Dispose(),直接标记_disposed = true并唤醒所有线程,线程醒来后检测到销毁信号直接退出循环,不再处理剩余任务。
  2. 无任务完成等待机制:线程池没有跟踪任务执行状态,主线程无法知晓所有任务是否完成,导致程序提前终止。

修复方案

添加任务计数和主线程等待逻辑,确保所有任务执行完毕后再销毁线程池。

修改后的MyThreadPool.cs

using System;
using System.Collections.Generic;
using System.Threading;

namespace ThreadPool
{
    public class MyThreadPool : IDisposable
    {
        private bool _disposed = false;
        private object _lock = new object();
        private Queue<MyTask> tasks = new();
        // 跟踪待处理/处理中的任务数量
        private int _pendingTasks = 0;
        // 用于主线程等待所有任务完成的信号量
        private ManualResetEventSlim _allTasksCompleted = new ManualResetEventSlim(false);

        public MyThreadPool(int numThreads = 0)
        {
            for (int i = 0; i < numThreads; i++)
            {
                new Thread(() =>
                {
                    while (!_disposed)
                    {
                        MyTask? task;
                        lock (_lock)
                        {
                            // 等待有任务可执行或线程池被销毁
                            while (!tasks.TryDequeue(out task) && !_disposed)
                            {
                                Monitor.Wait(_lock);
                            }
                        }

                        if (task != null)
                        {
                            try
                            {
                                Console.WriteLine($"Thread {Thread.CurrentThread.ManagedThreadId} calculated factorial for {task.getNum()}. Result: {task.CalculateFactorial(task.getNum())}");
                            }
                            finally
                            {
                                lock (_lock)
                                {
                                    // 任务完成,计数减1
                                    _pendingTasks--;
                                    // 所有任务完成时触发信号
                                    if (_pendingTasks == 0)
                                    {
                                        _allTasksCompleted.Set();
                                    }
                                }
                            }
                        }
                    }
                }) { IsBackground = false }.Start(); // 设置为前台线程,避免主线程退出后被强制终止
            }
        }

        public MyTask Enqueue(int number)
        {
            MyTask task = new(number);
            lock (_lock)
            {
                tasks.Enqueue(task);
                _pendingTasks++;
                // 有新任务时重置完成信号
                _allTasksCompleted.Reset();
                Monitor.Pulse(_lock);
            }
            return task;
        }

        // 等待所有任务完成的方法
        public void WaitForAllTasks()
        {
            _allTasksCompleted.Wait();
        }

        public void Dispose()
        {
            // 先等待所有任务完成再销毁线程池
            WaitForAllTasks();
            
            lock (_lock)
            {
                _disposed = true;
                Monitor.PulseAll(_lock);
            }
            _allTasksCompleted.Dispose();
        }
    }
}

修改后的Program.cs

int[] array = { 1, 2, 3, 4, 5, 6 };
int numThread = 3;
using (MyThreadPool threadPool = new(numThread))
{
    for (int i = 0; i < array.Length; i++)
    {
        threadPool.Enqueue(array[i]);
    }
    // 等待所有任务执行完毕
    threadPool.WaitForAllTasks();
}
// using块会自动调用Dispose,确保线程池正确销毁

关键修改点

  • 新增_pendingTasks变量跟踪任务数量,入队时加1,任务完成时减1
  • 使用ManualResetEventSlim实现主线程对任务完成状态的等待
  • 将工作线程设置为前台线程,防止主线程提前退出导致任务中断
  • 调整Dispose逻辑,先等待所有任务完成再标记线程池销毁

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 10:39:54