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

C#中使用BlockingCollection并发消息反序列化的线程安全问题

问题描述

我实现了一个基于BlockingCollection的消息队列,用于Socket接收消息的生产消费:

internal class Program
{
    public BlockingCollection<ArraySegment<byte>> MessageQueue { get; set; } = [];

    static void Main(string[] args)
    {
        // Task.Run(Consumer)
        // Callback Func for socket ReceiveAsync calls Producer
    }

    public void Producer(ArraySegment<byte> person)
    {
        MessageQueue.Add(person);
    }

    public void Consumer(CancellationToken token)
    {
        do
        {
            var person = MessageQueue.Take();
            var update = Deserialize(person);

            if (update != null)
            {
                DoWork(update);
            }

        } while (!token.IsCancellationRequested);
    }

    public static Person? Deserialize(ArraySegment<byte> socketMessage)
    {
        try
        {
            Person? person = JsonSerializer.Deserialize<Person>(socketMessage);

            return person;
        }
        catch (Exception e)
        {
            Console.WriteLine(e);

            return null;
        }
    }
}

异常情况

反序列化时抛出如下异常:

System.Text.Json.JsonException: '5' is invalid after a single JSON value. Expected end of data. Path: $ | LineNumber: 0 | BytePositionInLine: 656.  ---> System.Text.Json.JsonReaderException: '5' is invalid after a single JSON value. Expected end of data.

查看Encoding.UTF8.GetString(socketMessage)发现字符串已损坏,字节被截断或追加,导致JSON无效。

奇怪现象

如果把反序列化逻辑移到Producer中,将BlockingCollection改为BlockingCollection<Person>,消费者直接取出已反序列化的对象,后续DoWork可正常执行。我怀疑每次Take后反序列化的字节被其他消息的部分内容覆盖,但暂未找到证据。目前将反序列化移至Producer可解决问题,但并非理想方案。

Socket相关代码

var rentedBuffer = ArrayPool<byte>.Shared.Rent(receiveBufferSize);
var buffer = new ArraySegment<byte>(rentedBuffer, 0, receiveBufferSize);

result = await handler.ReceiveAsync(buffer, cancellationToken);
var message = new ArraySegment<byte>(buffer.Array, buffer.Offset, result.Count);
onMessage.ForEach(onmessage => onmessage(message));
ArrayPool<byte>.Shared.Return(rentedBuffer);
问题原因与解决方案

核心原因:ArrayPool缓冲区提前归还导致数据覆盖

Socket代码中,你在触发onMessage回调(也就是让Producer把ArraySegment加入队列)之后,立刻将租来的缓冲区还给了ArrayPool<byte>.Shared。而ArraySegment<byte>本质只是对原数组的引用+范围标记,并没有复制字节内容。

当后续其他Socket操作从池中租用到同一个数组时,会直接覆盖数组里的旧数据。等到Consumer线程去反序列化这个ArraySegment时,原数组的内容已经被新的Socket接收操作修改,自然会出现JSON损坏、截断或追加的情况。

为什么移到Producer就正常?

因为反序列化是在归还缓冲区之前完成的,此时原数组的内容还没被覆盖,能正确解析出Person对象。队列里存储的是对象实例(值类型或引用类型的副本),不受后续缓冲区复用的影响。

解决方案

1. 复制字节内容(推荐)

在Producer中,把ArraySegment的字节复制到独立的新数组里,再将新数组的ArraySegment加入队列。这样队列里的引用指向专属数组,不会被ArrayPool复用覆盖:

public void Producer(ArraySegment<byte> person)
{
    byte[] copy = new byte[person.Count];
    Array.Copy(person.Array, person.Offset, copy, 0, person.Count);
    MessageQueue.Add(new ArraySegment<byte>(copy));
}

2. 延迟归还缓冲区

确保所有使用该ArraySegment的操作完成后,再归还缓冲区。比如用Task.WhenAll等待所有onmessage回调执行完毕:

var tasks = onMessage.Select(onmessage => Task.Run(() => onmessage(message)));
await Task.WhenAll(tasks);
ArrayPool<byte>.Shared.Return(rentedBuffer);

这种方式需要跟踪缓冲区的生命周期,复杂度较高,适合对内存开销敏感的场景。

3. 保留当前方案

如果反序列化的性能开销不大,继续在Producer中完成反序列化,直接存储Person对象到队列。这种方式最简单,也能彻底避免缓冲区复用的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 21:47:16