单生产者单消费者模型死锁问题排查及实现合理性咨询
问题分析与解决方案
你的实现存在多处问题,直接导致了消费者线程死锁;但仅用两个线程(生产者+消费者)实现生产者/消费者模型是完全可行的,只要逻辑严谨。
核心问题点
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..."); } } }
关键修复说明
_continueProducingData标记为volatile:确保消费者能及时读取到生产者对该变量的修改。- 生产者结束时主动唤醒消费者:在
finally块中设置终止状态后,调用Monitor.Pulse唤醒可能阻塞的消费者。 - 优化等待逻辑:用
while循环替代if判断等待条件,避免虚假唤醒导致的逻辑错误。 - 控制生产数据量:计算剩余需要生产的字节数,防止最后一次生产超出目标文件大小。
- 消费者处理剩余数据:循环条件同时检查生产状态和缓冲区状态,确保生产结束后处理完所有剩余数据再退出。
关于双线程生产者/消费者模型的可行性
完全正确且高效。双线程模型是生产者/消费者的经典场景之一:一个线程负责生成数据,另一个负责消费,通过锁和等待/唤醒机制协调节奏,避免资源浪费和数据竞争,只要逻辑严谨就不会出现问题。
内容的提问来源于stack exchange,提问作者Usman
相关产品推荐
相关产品推荐

