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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:30:43