使用ConcurrentQueue等并发组件实现多线程文件写入的问题排查
问题分析与解决方案
核心问题:共享缓冲区引用导致数据覆盖
你的代码中最严重的问题是所有DataBuffer实例共享同一个_ReceivedRawDataBuffer数组引用。在数据接收事件中,每次新数据到来都会覆盖这个数组的内容,而消费者线程读取时,拿到的是数组的最新状态(已被后续数据覆盖),不是入队时的原始数据。这直接导致了消费者端数据错乱、重复或丢失。
修复方法:为每个DataBuffer创建独立缓冲区
修改DataBuffer的构造逻辑,将接收到的数据复制到新的字节数组中,而不是传递共享数组的引用:
// 在数据接收事件中 byte[] copiedBuffer = new byte[(int)_ReceivedNumberOfBytes]; Array.Copy(e.Bytes, 0, copiedBuffer, 0, (int)_ReceivedNumberOfBytes); producerDataBuffer = new DataBuffer(copiedBuffer, (int)_ReceivedNumberOfBytes); collection.Add(producerDataBuffer);
或者直接修改DataBuffer构造函数自动复制:
public class DataBuffer { public byte[] Buffer { get; set; } public int Length { get; set; } public DataBuffer(byte[] buffer, int length) { // 创建新数组并复制数据,避免引用共享 Buffer = new byte[length]; Array.Copy(buffer, 0, Buffer, 0, length); Length = length; } }
各方案具体问题与修正
1. ConcurrentQueue 完全无法写入
原消费者代码仅尝试一次出队就退出,没有循环监听新数据。修正后的消费者需持续监听队列:
private static ConcurrentQueue<DataBuffer> queue = new ConcurrentQueue<DataBuffer>(); private static AutoResetEvent queueEvent = new AutoResetEvent(false); private static bool isRunning = true; static void Consumer() { using (var fileStream = File.OpenWrite("received_data.bin")) { while (isRunning) { // 等待数据信号 queueEvent.WaitOne(); // 一次性处理队列中所有待处理数据(避免多次触发事件) while (queue.TryDequeue(out var dataBuffer)) { fileStream.Write(dataBuffer.Buffer, 0, dataBuffer.Length); } } } } // 停止消费时调用 public void StopConsumer() { isRunning = false; queueEvent.Set(); // 唤醒等待的线程 }
2. AutoResetEvent 仅写入一次
消费者处理完一个数据后就退出循环,没有持续等待新的信号。修正方法同上,添加外层循环保持线程运行,直到主动停止。
3. BlockingCollection 数据错乱
除了共享缓冲区问题,消费者代码存在冗余循环:GetConsumingEnumerable()会自动阻塞并持续读取数据,直到调用CompleteAdding()。修正后的消费者:
private static BlockingCollection<DataBuffer> collection = new BlockingCollection<DataBuffer>(); static void Consumer() { using (var fileStream = File.OpenWrite("received_data.bin")) { // GetConsumingEnumerable() 会自动阻塞等待数据,直到集合标记为完成 foreach (DataBuffer dataBuffer in collection.GetConsumingEnumerable()) { fileStream.Write(dataBuffer.Buffer, 0, dataBuffer.Length); } } } // 停止数据接收时调用,通知消费者结束 public void StopCapture() { collection.CompleteAdding(); _FileStream?.Close(); _FileStreamBeforeQueue?.Close(); }
4. Channel 数据异常
同样是共享缓冲区问题,修正后结合Channel的异步特性优化消费者:
private static Channel<DataBuffer> channel = Channel.CreateUnbounded<DataBuffer>(); static async Task Consumer() { using (var fileStream = File.OpenWrite("received_data.bin")) { // 异步读取所有数据直到通道关闭 await foreach (var dataBuffer in channel.Reader.ReadAllAsync()) { System.Diagnostics.Debug.WriteLine("Consumer: " + BitConverter.ToString(dataBuffer.Buffer.Take(8).ToArray())); fileStream.Write(dataBuffer.Buffer, 0, dataBuffer.Length); } } } // 停止时关闭通道写入端 public void StopCapture() { channel.Writer.Complete(); _FileStream?.Close(); _FileStreamBeforeQueue?.Close(); }
额外注意事项
- 文件流关闭:停止数据捕获时需显式关闭
_FileStream和_FileStreamBeforeQueue,避免文件占用或数据未完全写入。 - 线程安全:确保
_OperationMode等状态变量的访问是线程安全的(可添加lock或使用volatile)。 - 资源清理:在应用退出或设备断开时,停止消费者任务并释放所有资源。
内容的提问来源于stack exchange,提问作者bLAZ
相关产品推荐
相关产品推荐

