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

单生产者单消费者模型死锁问题排查及实现合理性咨询

问题分析与解决方案

你的实现存在多处问题,直接导致了消费者线程死锁;但仅用两个线程(生产者+消费者)实现生产者/消费者模型是完全可行的,只要逻辑严谨。

核心问题点

1. 生产者结束后未通知消费者,导致无限等待

生产者退出循环后仅设置_continueProducingData = false,但如果消费者此时正处于Monitor.Wait(bufferLock)的阻塞状态,根本无法感知生产已结束,会一直卡在等待逻辑中。

2. _continueProducingData存在内存可见性问题

该变量未标记为volatile,也没有始终在锁内访问,消费者线程可能无法及时读取到生产者对它的修改,导致循环无法正常退出。

3. 生产者循环逻辑缺陷

生产者仅在缓冲区为空时才生成数据,否则等待,但完成最后一次生产后,没有额外的通知告知消费者“生产已终止”;同时原循环条件可能导致最后一次生产的数据超出目标文件大小。

4. 消费者未处理剩余数据

即使消费者感知到生产结束,当前逻辑也会直接退出,不会处理缓冲区中可能剩余的未消费数据。

修正后的实现方案

以下是修复后的核心代码,解决了死锁和数据完整性问题:

public class FileGenerator
{
    private static readonly string _fileLocation = Path.Combine("Resources", "files");
    private static Faker<Employee> _employeeGenerator = null!;
    private static Random _random = new Random(1);

    private long _maxBufferSize = 28;
    private long _maxFileSizeInBytes = 0;
    private long _totalBytesWritten = 0;

    private Queue<byte[]> _buffer = null;
    // 用volatile保证跨线程内存可见性
    private volatile bool _continueProducingData = false;
    static object bufferLock = new object();

    public FileGenerator()
    {
        _buffer = new Queue<byte[]>();
        _continueProducingData = true;

        _employeeGenerator = new Faker<Employee>()
            .RuleFor(e => e.RandomNumberDelimator, 100)
            .RuleFor(e => e.AboutMe, f => f.Lorem.Paragraph(1));
    }

    public async Task Execute(long fileSize, string location = null)
    {
        _maxFileSizeInBytes = fileSize * Convert.ToInt64(100);
        var filename = $"unsorted.{fileSize}.txt";
        var path = string.IsNullOrWhiteSpace(location) ? Path.Combine(_fileLocation, filename) : Path.Combine(location, filename);            

        var producerTask = Task.Run(() => GenerateTestDataSetInternal());
        var consumerTask = Task.Run(() => ConsumeTestDataSet(new StreamWriter(File.Open(path, FileMode.Create), bufferSize: 65535)));

        await Task.WhenAll(producerTask, consumerTask);
    }

    private void GenerateTestDataSetInternal()
    {
        try
        {
            while (_totalBytesWritten <= _maxFileSizeInBytes)
            {
                lock (bufferLock)
                {
                    // 用while循环处理虚假唤醒,确保缓冲区为空才生产
                    while (_buffer.Count > 0)
                    {
                        Monitor.Wait(bufferLock);
                    }

                    long noOfbytesWrittenPerDataBytes = 0;
                    long remainingBytes = _maxFileSizeInBytes - _totalBytesWritten;
                    // 控制生产数据不超过缓冲区上限和剩余需要的字节数
                    while (noOfbytesWrittenPerDataBytes <= _maxBufferSize && noOfbytesWrittenPerDataBytes < remainingBytes)
                    {
                        var dataBytes = Encoding.UTF8.GetBytes(_employeeGenerator.Generate().ToString());
                        // 避免最后一次生产超出总大小
                        if (noOfbytesWrittenPerDataBytes + dataBytes.Length > remainingBytes)
                        {
                            break;
                        }
                        noOfbytesWrittenPerDataBytes += dataBytes.Length;
                        _buffer.Enqueue(dataBytes);
                    }

                    _totalBytesWritten += noOfbytesWrittenPerDataBytes;
                    // 通知消费者有数据可处理
                    Monitor.Pulse(bufferLock);
                }
            }
        }
        finally
        {
            // 生产结束后主动通知消费者
            _continueProducingData = false;
            lock (bufferLock)
            {
                Monitor.Pulse(bufferLock);
            }
            Console.WriteLine("Producer finished...");
        }
    }

    private void ConsumeTestDataSet(StreamWriter writer)
    {
        try
        {
            while (true)
            {
                byte[] dataBytes = null;
                lock (bufferLock)
                {
                    // 等待条件:缓冲区为空且生产未结束
                    while (_buffer.Count == 0 && _continueProducingData)
                    {
                        Monitor.Wait(bufferLock);
                    }

                    // 生产结束且缓冲区为空,退出循环
                    if (!_continueProducingData && _buffer.Count == 0)
                    {
                        break;
                    }

                    if (_buffer.Count > 0)
                    {
                        dataBytes = _buffer.Dequeue();
                    }
                }

                if (dataBytes != null)
                {
                    writer.BaseStream.Write(dataBytes, 0, dataBytes.Length);
                    writer.Flush();
                    // 通知生产者缓冲区有空间
                    lock (bufferLock)
                    {
                        Monitor.Pulse(bufferLock);
                    }
                }
            }
        }
        finally
        {
            writer.Close();
            Console.WriteLine("Consumer finished...");
        }
    }
}

关键修复说明

  1. _continueProducingData标记为volatile:确保消费者能及时读取到生产者对该变量的修改。
  2. 生产者结束时主动唤醒消费者:在finally块中设置终止状态后,调用Monitor.Pulse唤醒可能阻塞的消费者。
  3. 优化等待逻辑:用while循环替代if判断等待条件,避免虚假唤醒导致的逻辑错误。
  4. 控制生产数据量:计算剩余需要生产的字节数,防止最后一次生产超出目标文件大小。
  5. 消费者处理剩余数据:循环条件同时检查生产状态和缓冲区状态,确保生产结束后处理完所有剩余数据再退出。

关于双线程生产者/消费者模型的可行性

完全正确且高效。双线程模型是生产者/消费者的经典场景之一:一个线程负责生成数据,另一个负责消费,通过锁和等待/唤醒机制协调节奏,避免资源浪费和数据竞争,只要逻辑严谨就不会出现问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 03:25:10