ManualResetEventSlim.Wait()阻塞所有Task.Run启动的任务,导致生产者消费者内存控制逻辑停滞
我来帮你拆解问题根源,以及给出可落地的修复方案:
问题核心原因
你的代码里有3个关键问题共同导致了程序停滞:
阻塞式Wait()引发线程池饥饿
你用Task.Run启动的生产者/消费者都运行在ThreadPool线程上。当生产者触发内存阈值后,调用ManualResetEventSlim.Wait()是同步阻塞调用,会持续占用这个ThreadPool线程不释放。如果ThreadPool可用线程被耗尽(或调度优先级失衡),消费者线程可能无法获得足够CPU时间处理队列、释放内存并触发Set(),最终形成死锁。信号量触发逻辑不严谨
消费者仅在MemoryThresholdLimit > currentMemory || DisposableFileBuffer.NumberOfProcessedItemsInQueue == 0时调用Set(),但如果消费者处理完一个元素后内存仍超过阈值,就不会触发唤醒信号,导致生产者一直阻塞。同步文件操作加剧线程阻塞
File.ReadAllBytes和File.WriteAllBytes是同步阻塞调用,会进一步占用ThreadPool线程,让线程池饥饿问题雪上加霜。
修复方案
我们需要用异步等待替代阻塞等待,优化内存控制逻辑,同时把文件操作改为异步模式,彻底避免线程池资源被占用。
1. 替换ManualResetEventSlim为SemaphoreSlim
SemaphoreSlim支持WaitAsync()异步等待方法,不会阻塞ThreadPool线程,完美适配异步场景:
// 替换原来的静态ManualResetEventSlim private readonly SemaphoreSlim _memorySemaphore = new SemaphoreSlim(1, 1); private readonly BlockingCollection<DisposableFileBuffer> m_fileQueue = new BlockingCollection<DisposableFileBuffer>(); private const long MemoryThresholdLimit = 100L * 1024 * 1024;
2. 重构生产者为全异步逻辑
把同步文件读取改为异步,用await _memorySemaphore.WaitAsync()替代阻塞等待:
private async void button1_Click(object sender, EventArgs e) { string readPath = @"C:\projects"; string savePath = @"C:\thrash"; Directory.CreateDirectory(savePath); // 确保目标目录存在 // 直接调用异步方法,无需Task.Run包装 await Task.WhenAll( ReadFromHddAsync(readPath), SaveToHddAsync(savePath)); MessageBox.Show("Done"); } // 全异步生产者 private async Task ReadFromHddAsync(string path) { foreach (var fileName in Directory.EnumerateFiles(path, "*.*", SearchOption.AllDirectories)) { try { if (!File.Exists(fileName)) continue; // 异步读取文件,避免阻塞线程池 byte[] byteArray; using (var stream = new FileStream(fileName, FileMode.Open, FileAccess.Read, FileShare.Read, 4096, FileOptions.Asynchronous)) { byteArray = new byte[stream.Length]; await stream.ReadAsync(byteArray, 0, byteArray.Length); } var buffer = new DisposableFileBuffer(byteArray); m_fileQueue.Add(buffer); // 循环检查内存阈值,异步等待直到内存回落 long currentMemory = Process.GetCurrentProcess().PrivateMemorySize64; while (currentMemory > MemoryThresholdLimit && m_fileQueue.Count > 0) { Debug.WriteLine( $"Wait: 已用内存 {currentMemory / 1024 / 1024}Mb, 队列长度 {m_fileQueue.Count}"); await _memorySemaphore.WaitAsync(); // 异步等待,不占用线程 // 重新检查内存,避免虚假唤醒 currentMemory = Process.GetCurrentProcess().PrivateMemorySize64; } } catch (Exception ex) { long currentMemory = Process.GetCurrentProcess().PrivateMemorySize64; Debug.WriteLine($"读取文件失败: {ex.Message}, 已用内存 {currentMemory / 1024 / 1024}Mb"); } } m_fileQueue.CompleteAdding(); // 生产者完成,标记队列结束 }
3. 优化消费者的唤醒逻辑
消费者处理完元素后,只要内存低于阈值就释放信号量,确保生产者能及时被唤醒:
// 全异步消费者 private async Task SaveToHddAsync(string path) { foreach (var buffer in m_fileQueue.GetConsumingEnumerable()) { try { string saveFileName = Path.Combine(path, Guid.NewGuid() + ".bin"); // 异步写入文件 using (var stream = new FileStream(saveFileName, FileMode.Create, FileAccess.Write, FileShare.None, 4096, FileOptions.Asynchronous)) { await stream.WriteAsync(buffer.Data, 0, buffer.Data.Length); } } catch (Exception ex) { Debug.WriteLine($"写入文件失败: {ex.Message}"); } finally { // 确保资源被释放 buffer.Dispose(); // 检查内存,低于阈值则唤醒生产者 long currentMemory = Process.GetCurrentProcess().PrivateMemorySize64; if (currentMemory <= MemoryThresholdLimit) { _memorySemaphore.Release(); Debug.WriteLine( $"唤醒生产者: 已用内存 {currentMemory / 1024 / 1024}Mb, 队列长度 {m_fileQueue.Count}"); } } } }
4. 修正DisposableFileBuffer的线程安全与内存回收
确保静态计数的线程安全,优化内存释放逻辑:
public class DisposableFileBuffer : IDisposable { // 用Interlocked保证多线程下计数准确 private static int _processedItemCount = 0; public static int NumberOfProcessedItemsInQueue => Volatile.Read(ref _processedItemCount); public byte[] Data { get; private set; } private bool _disposed; public DisposableFileBuffer(byte[] data) { Data = data ?? throw new ArgumentNullException(nameof(data)); Interlocked.Increment(ref _processedItemCount); } public void Dispose() { if (_disposed) return; // 立即释放数组引用,帮助GC回收 Data = null; Interlocked.Decrement(ref _processedItemCount); _disposed = true; GC.SuppressFinalize(this); } ~DisposableFileBuffer() { Dispose(); } }
额外优化建议
- 内存指标选择:用
PrivateMemorySize64替代WorkingSet64,它更准确反映进程实际占用的私有内存(不包含系统共享内存)。 - BlockingCollection容量限制:可以给
BlockingCollection设置一个固定容量,比如new BlockingCollection<DisposableFileBuffer>(10),配合信号量进一步避免队列无限膨胀。 - 避免手动GC:生产环境尽量不要调用
GC.Collect(),通过合理的资源释放逻辑让GC自动处理内存回收。
内容来源于stack exchange

