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

ManualResetEventSlim.Wait()阻塞所有Task.Run启动的任务,导致生产者消费者内存控制逻辑停滞

ManualResetEventSlim.Wait()阻塞所有Task.Run启动的任务,导致生产者消费者内存控制逻辑停滞

我来帮你拆解问题根源,以及给出可落地的修复方案:


问题核心原因

你的代码里有3个关键问题共同导致了程序停滞:

  1. 阻塞式Wait()引发线程池饥饿
    你用Task.Run启动的生产者/消费者都运行在ThreadPool线程上。当生产者触发内存阈值后,调用ManualResetEventSlim.Wait()是同步阻塞调用,会持续占用这个ThreadPool线程不释放。如果ThreadPool可用线程被耗尽(或调度优先级失衡),消费者线程可能无法获得足够CPU时间处理队列、释放内存并触发Set(),最终形成死锁。

  2. 信号量触发逻辑不严谨
    消费者仅在MemoryThresholdLimit > currentMemory || DisposableFileBuffer.NumberOfProcessedItemsInQueue == 0时调用Set(),但如果消费者处理完一个元素后内存仍超过阈值,就不会触发唤醒信号,导致生产者一直阻塞。

  3. 同步文件操作加剧线程阻塞
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 11:38:05