【问题标题】:Buffering byte data in C#在 C# 中缓冲字节数据
【发布时间】:2018-01-27 02:43:33
【问题描述】:

我的应用程序从 TCP 套接字读取字节并需要缓冲它们,以便我以后可以从中提取消息。由于 TCP 的性质,我可能会在一次读取中获得部分或多条消息,因此在每次读取后,我想检查缓冲区并提取尽可能多的完整消息。

因此,我想要一个允许我执行以下操作的课程:

  • 向其附加任意字节[]数据
  • 在不消费的情况下检查内容,特别是检查内容的数量,并搜索是否存在某个或多个字节
  • 提取和使用部分数据作为字节 [],而将其余部分保留在其中以供将来读取

我希望可以使用 .NET 库中的 1 个或多个现有类来完成我想要的操作,但我不确定哪些类。 System.IO.MemoryStream 看起来接近我想要的,但是(a)不清楚它是否适合用作缓冲区(读取的数据是否从容量中删除?)和(b) 读取和写入似乎发生在同一个地方 - “流的当前位置是下一次读取或写入操作可能发生的位置。” - 这不是我想。我需要写到最后,从前面读。

【问题讨论】:

  • 你想追加任意字节数组吗?还是任意字节?无论如何,也许 List 或 List 会为你工作
  • 字节数组,但它们被合并到一个连续的字节数组中。我认为 List 对于这个应用程序来说不够高效。
  • 我不知道有性能要求,jon 的建议看起来不错
  • 嗯,List其实可能没问题,检查了实现,但是接口有点麻烦。
  • 缓冲流呢?

标签: c# .net


【解决方案1】:

我建议你在后台使用MemoryStream,但将其封装在另一个存储类中:

  • MemoryStream
  • 当前“读取”位置
  • 当前“消耗”位置

然后它会暴露:

  • 写入:将流的位置设置到末尾,写入数据,将流的位置设置回读取位置
  • Read:读取数据,将读取位置设置为流的位置
  • 消费:更新消费位置(具体取决于您尝试消费的方式);如果消费位置高于某个阈值,则将现有缓冲数据复制到新的MemoryStream 并更新所有变量。 (您可能不想在每个消费请求上复制缓冲区。)

请注意,如果没有额外的同步,这些都不是线程安全的。

【讨论】:

  • 与普通的byte[] 数组相比,在底层使用MemoryStream 有什么优势? (好吧,没关系,明白了:它会自动调整大小)
  • @Groo:MemoryStream 具有可扩展的容量(如果需要)。
  • 将现有 MemoryStream 的剩余部分复制到新的有什么好的方法?
  • @Kylotan:最快的方法可能是调用MemoryStream.GetBuffer 来获取底层缓冲区,然后使用普通的Write 调用将一部分写入新流。
  • 我接受了下面发布的示例代码的答案,但也感谢这个答案,它看起来同样正确。
【解决方案2】:

只需使用大字节数组和 Array.Copy - 它应该可以解决问题。 如果没有,请使用List<byte>

如果你使用数组,你必须自己实现一个索引(复制额外数据的地方)(检查内容大小也是如此),但这很简单。

如果您有兴趣:这里是“循环缓冲区”的简单实现。测试应该运行(我对其进行了几个单元测试,但它没有检查所有关键路径):

public class ReadWriteBuffer
{
    private readonly byte[] _buffer;
    private int _startIndex, _endIndex;

    public ReadWriteBuffer(int capacity)
    {
        _buffer = new byte[capacity];
    }

    public int Count
    {
        get
        {
            if (_endIndex > _startIndex)
                return _endIndex - _startIndex;
            if (_endIndex < _startIndex)
                return (_buffer.Length - _startIndex) + _endIndex;
            return 0;
        }
    }

    public void Write(byte[] data)
    {
        if (Count + data.Length > _buffer.Length)
            throw new Exception("buffer overflow");
        if (_endIndex + data.Length >= _buffer.Length)
        {
            var endLen = _buffer.Length - _endIndex;
            var remainingLen = data.Length - endLen;

            Array.Copy(data, 0, _buffer, _endIndex, endLen);
            Array.Copy(data, endLen, _buffer, 0, remainingLen);
            _endIndex = remainingLen;
        }
        else
        {
            Array.Copy(data, 0, _buffer, _endIndex, data.Length);
            _endIndex += data.Length;
        }
    }

    public byte[] Read(int len, bool keepData = false)
    {
        if (len > Count)
            throw new Exception("not enough data in buffer");
        var result = new byte[len];
        if (_startIndex + len < _buffer.Length)
        {
            Array.Copy(_buffer, _startIndex, result, 0, len);
            if (!keepData)
                _startIndex += len;
            return result;
        }
        else
        {
            var endLen = _buffer.Length - _startIndex;
            var remainingLen = len - endLen;
            Array.Copy(_buffer, _startIndex, result, 0, endLen);
            Array.Copy(_buffer, 0, result, endLen, remainingLen);
            if (!keepData)
                _startIndex = remainingLen;
            return result;
        }
    }

    public byte this[int index]
    {
        get
        {
            if (index >= Count)
                throw new ArgumentOutOfRangeException();
            return _buffer[(_startIndex + index) % _buffer.Length];
        }
    }

    public IEnumerable<byte> Bytes
    {
        get
        {
            for (var i = 0; i < Count; i++)
                yield return _buffer[(_startIndex + i) % _buffer.Length];
        }
    }
}

请注意:代码在读取时会“消耗” - 如果您不想这样做,只需删除“_startIndex = ...”部分(或制作重载可选参数并检查或其他)。

【讨论】:

  • 看起来很不错,一个缺点是在决定是否调用 Read() 时,我真的可以使用整个缓冲区的连续视图来做。有没有没有副本的好方法来实现它?也许使用迭代两个部分的 IEnumerable?
  • 我更新了代码,以便您可以选择保持读取数据 - 可以吗? (您可以使用 obj.Read(obj.Count, true) 或更具可读性的 obj.Read(obj.Count, keepData: true) 预览所有数据
  • 是的,我知道我可以这样做,但我希望每次都能避免这些副本。这是我可以做的很多额外的分配和释放。 (我是一名游戏程序员,这会被很多人称为。)
  • 哦 - Array.Copy ...对不起。我添加了一个索引器和一个 IEnumerable ... 完成
  • 我可能无法使用它,因为为每个字节调用该函数太慢了。不过不知道我会做什么!
【解决方案3】:

我认为BufferedStream 是解决问题的方法。也可以通过调用 Seek 去读取 len 字节的数据。

BufferdStream buffer = new BufferedStream(tcpStream, size); // we have a buffer of size
...
...
while(...)
{
    buffer.Read(...);
    // do my staff
    // I read too much, I want to put back len bytes
    buffer.Seek(-len, SeekOrigin.End);
    // I shall come back and read later
}

不断增长的记忆

与最初指定sizeBufferedStream 相反,MemoryStream 可以增长。

记住流数据

MemoryStream 一直保存所有的 dara-read,而BufferedStream 只保存一段流数据。

源流与字节数组

MemoryStream 允许在Write() 方法中添加输入字节,将来可以是Read()。而BufferedSteam 从构造函数中指定的另一个源流中获取输入字节。

【讨论】:

    【解决方案4】:

    这是我不久前写的一个缓冲区的another implementation

    • Resizeable:允许对数据进行排队,不抛出缓冲区溢出异常;
    • 高效:使用单个缓冲区和 Buffer.Copy 操作将数据入队/出队

    【讨论】:

    • 如果它允许任意长度的 peek 操作,那对我的目的来说会很棒。否则,要知道我是否可以安全地使用数据可能会有点棘手。
    【解决方案5】:

    迟到了,但为了后代:

    当我过去这样做时,我采取了稍微不同的方法。 如果您的消息具有固定的标头大小(告诉您正文中有多少字节),并且记住网络流已经缓冲,我执行操作两个阶段:

    • 在流​​上读取标头的字节
    • 随后根据标头读取流中的正文字节
    • 重复

    这利用了以下事实 - 对于流 - 当您请求 'n' 字节时,您将永远不会得到 more 回来,因此您可以忽略许多“我读了太多”的 opps,让我把这些放在一边,直到下一次。

    公平地说,这并不是故事的全部。我在流上有一个底层包装类来处理碎片问题(即,如果要求提供 4 个字节,则在收到 4 个字节或流关闭之前不要返回)。但这一点相当容易。

    在我看来,关键是将消息处理与流机制解耦,如果您不再尝试将消息作为单个 ReadBytes() 从流中使用,生活就会变得简单得多。

    [无论您的读取是阻塞的还是异步的(APM/await),所有这些都是正确的]

    【讨论】:

    • 我会遇到的问题是,“如果要求 4 个字节,在收到 4 个字节之前不要返回”的方法对于我需要运行的那种系统是不实用的,它可以t 等待数据时挂起线程。我需要能够将部分数据放在一边并稍后返回,这需要一个 FIFO 缓冲区,正如其他答案所建议的那样。 (我的主要问题真的是'为什么标准库中没有这样的东西?)
    【解决方案6】:

    听起来您想从套接字读取到 MemoryStream 缓冲区,然后从缓冲区中“弹出”数据并在每次遇到某个字节时重置它。它看起来像这样:

    void ReceiveAllMessages(Action<byte[]> messageReceived, Socket socket)
    {
        var currentMessage = new MemoryStream();
        var buffer = new byte[128];
    
        while (true)
        {
            var read = socket.Receive(buffer, 0, buffer.Length);
            if (read == 0)
                break;     // Connection closed
    
            for (var i = 0; i < read; i++)
            {
                var currentByte = buffer[i];
                if (currentByte == END_OF_MESSAGE)
                {
                    var message = currentMessage.ToByteArray();
                    messageReceived(message);
    
                    currentMessage = new MemoryStream();
                }
                else
                {
                    currentMessage.Write(currentByte);
                }
            }
        }
    }
    

    【讨论】:

    • 不,我不需要解析方面的帮助,只需要实际缓冲,但谢谢。 :)
    • 如果消息被分成几个数据包,这将不起作用。正如 OP 所要求的,它必须使用 FIFO 缓冲区来实现。
    • 确实如此 - 它会一直附加到同一个 MemoryStream 直到遇到 END_OF_MESSAGE 字节,这可能在第一个或第 85 个数据包中
    • edit 知道了,没看到无限循环。但是,您仍然在同一方法中混合了两种职责:接收和解析数据。如果没有“END_OF_MESSAGE”字节,但每条消息的长度取决于其内容(例如,特定的起始 cookie,然后是消息中编码的长度信息)怎么办?这个问题必须通过两个步骤来解决:1. 类只将字节排入 FIFO 缓冲区和 2. 解析数据并将其出列的类。
    【解决方案7】:

    你可以用 Stream 包裹 ConcurrentQueue&lt;ArraySegment&lt;byte&gt;&gt; 来做到这一点(请记住,这只会让它向前)。但是,我真的不喜欢在使用数据之前将数据保存在内存中的想法。它使您面临一堆关于消息大小的攻击(有意或无意)。您可能还想Google "circular buffer"

    您实际上应该编写的代码在收到数据后立即对数据执行一些有意义的操作:“推送解析”(例如,SAX 支持的内容)。例如,您将如何处理文本:

    private Encoding _encoding;
    private Decoder _decoder;
    private char[] _charData = new char[4];
    
    public PushTextReader(Encoding encoding)
    {
        _encoding = encoding;
        _decoder = _encoding.GetDecoder();
    }
    
    // A single connection requires its own decoder
    // and charData. That connection should never
    // call this method from multiple threads
    // simultaneously.
    // If you are using the ReadAsyncLoop you
    // don't need to worry about it.
    public void ReceiveData(ArraySegment<byte> data)
    {
        // The two false parameters cause the decoder
        // to accept 'partial' characters.
        var charCount = _decoder.GetCharCount(data.Array, data.Offset, data.Count, false);
        charCount = _decoder.GetChars(data.Array, data.Offset, data.Count, _charData, 0, false);
        OnCharacterData(new ArraySegment<char>(_charData, 0, charCount));
    }
    

    如果您必须能够在反序列化之前接受完整的消息,则可以使用MemoryMappedFile,其优点是发送实体不会出现内存不足的情况你的服务器。棘手的是将文件重置为零。因为这可能会带来很多问题。解决此问题的一种方法是:

    TCP 接收端

    1. 写入当前流。
    2. 如果流超过特定长度,则移动到新的。

    反序列化结束

    1. 从当前流中读取。
    2. 清空流后将其销毁。

    TCP 接收端非常简单。解串器端将需要一些基本的缓冲区拼接逻辑(请记住使用Buffer.BlockCopy 而不是Array.Copy)。

    旁注:听起来是一个有趣的项目,如果我有时间并记得我可能会继续实施这个系统。

    【讨论】:

    • Jonathan,我看不出解析部分消息如何解决您提到的安全问题,因为解析的对象可能比原始字节大 - 攻击实际上变得更有效。事实上,无论如何,我的消息都有大小限制,并且大小在标头中,因此可以尽早检测到无效消息。
    • @Kylotan - 向服务器发送昂贵的消息也可以正常工作,即使您限制消息的总大小,也可以发送一条消息,例如20 秒执行,然后只是排队一堆其他的:有意或无意再次。
    • 当然,但我看不出有任何内在原因可以解释为什么这些消息通过零碎解析会降低成本。
    • 这正在成为一个讨论。如果您不戴安全帽,不这样做也没关系。
    • 我戴上了安全帽,但我不同意你的方式更安全,抱歉。将数据解析为更大的表示会使您更有可能达到资源限制,而不是更少,并且需要额外的 CPU 和内存资源来解析可能最终不会完全形成的消息。
    【解决方案8】:

    这里只有三个提供代码的答案。 其中一个很笨拙,其他人不回答问题。

    这是一个你可以复制和粘贴的类:

    /// <summary>
    /// This class is a very fast and threadsafe FIFO buffer
    /// </summary>
    public class FastFifo
    {
        private List<Byte> mi_FifoData = new List<Byte>();
    
        /// <summary>
        /// Get the count of bytes in the Fifo buffer
        /// </summary>
        public int Count
        {
            get 
            { 
                lock (mi_FifoData)
                {
                    return mi_FifoData.Count; 
                }
            }
        }
    
        /// <summary>
        /// Clears the Fifo buffer
        /// </summary>
        public void Clear()
        {
            lock (mi_FifoData)
            {
                mi_FifoData.Clear();
            }
        }
    
        /// <summary>
        /// Append data to the end of the fifo
        /// </summary>
        public void Push(Byte[] u8_Data)
        {
            lock (mi_FifoData)
            {
                // Internally the .NET framework uses Array.Copy() which is extremely fast
                mi_FifoData.AddRange(u8_Data);
            }
        }
    
        /// <summary>
        /// Get data from the beginning of the fifo.
        /// returns null if s32_Count bytes are not yet available.
        /// </summary>
        public Byte[] Pop(int s32_Count)
        {
            lock (mi_FifoData)
            {
                if (mi_FifoData.Count < s32_Count)
                    return null;
    
                // Internally the .NET framework uses Array.Copy() which is extremely fast
                Byte[] u8_PopData = new Byte[s32_Count];
                mi_FifoData.CopyTo(0, u8_PopData, 0, s32_Count);
                mi_FifoData.RemoveRange(0, s32_Count);
                return u8_PopData;
            }
        }
    
        /// <summary>
        /// Gets a byte without removing it from the Fifo buffer
        /// returns -1 if the index is invalid
        /// </summary>
        public int PeekAt(int s32_Index)
        {
            lock (mi_FifoData)
            {
                if (s32_Index < 0 || s32_Index >= mi_FifoData.Count)
                    return -1;
    
                return mi_FifoData[s32_Index];
            }
        }
    }
    

    【讨论】:

    • 为了后代,这个实现将写入它的所有字节装箱到对象中,以便可以将它们添加到集合 mi_FifoData。考虑到这一点,这个类的性能和内存占用将是可怕的,如果不是致命的话。
    猜你喜欢
    • 1970-01-01
    • 2013-03-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多