【问题标题】:Producer/Consumer, Stream buffer problem生产者/消费者,流缓冲区问题
【发布时间】:2011-03-21 01:47:51
【问题描述】:

我正在尝试编写一个管理 3 个流的缓冲区管理器。典型的用法是慢速生产者和快速消费者。三个缓冲区背后的想法是,生产者总是有一个缓冲区要写入,而消费者总是得到最新的数据。

现在我已经有了这个,它可以正常工作了。

namespace YariIfStream
{

    /// <summary>
    /// A class that manages three buffers used for IF data streams
    /// </summary>
    public class YariIFStream
    {
        private Stream writebuf; ///<value>The stream used for writing</value>
        private Stream readbuf; ///<value>The stream used for reading</value>
        private Stream swapbuf; ///<value>The stream used for swapping</value>
        private bool firsttime; ///<value>Boolean used for checking if it is the first time a writebuffers is asked</value>
        private Object sync; ///<value>Object used for syncing</value>

        /// <summary>
        /// Initializes a new instance of the Yari.YariIFStream class with expandable buffers
        /// </summary>
        public YariIFStream()
        {
            sync = new Object();
            eerste = true;

            writebuf = new MemoryStream();
            readbuf = new MemoryStream();
            swapbuf = new MemoryStream();
        }

        /// <summary>
        /// Returns the stream with the buffer with new data ready to be read
        /// </summary>
        /// <returns>Stream</returns>
        public Stream GetReadBuffer()
        {
            lock (sync)
            {
                Monitor.Wait(sync);
                Stream tempbuf = swapbuf;
                swapbuf = readbuf;
                readbuf = tempbuf;
            }
            return readbuf;
        }

        /// <summary>
        /// Returns the stream with the buffer ready to be written with data
        /// </summary>
        /// <returns>Stream</returns>
        public Stream GetWriteBuffer()
        {
            lock (sync)
            {
                Stream tempbuf = swapbuf;
                swapbuf = writebuf;
                writebuf = tempbuf;
                if (!firsttime)
                {
                    Monitor.Pulse(sync);
                }
                else
                {
                    firsttime = false;

                }
            }
            //Thread.Sleep(1);
            return writebuf;
        }

    }
}

使用第一次检查是因为第一次请求写入缓冲区时,它不能脉冲消费者,因为缓冲区仍然必须写入数据。当第二次询问 writebuffer 时,我们可以确定前一个缓冲区包含数据。

我有两个线程,一个生产者和一个消费者。 这是我的输出:

prod: uv_hjd`alv   cons: N/<]g[)8fV
prod: N/<]g[)8fV   cons: 5Ud*tJ-Qkv
prod: 5Ud*tJ-Qkv   cons: 4Lx&Z7qqjA
prod: 4Lx&Z7qqjA   cons: kjUuVyCa.B
prod: kjUuVyCa.B

现在消费者落后了也没关系,它应该这样做。 如您所见,我丢失了第一串数据,这是我的主要问题。

其他问题是这样的:

  • 如果我删除第一次检查,它可以工作。但我认为不应该...
  • 如果我添加一个 Thread.Sleep(1);在 GetWriteBuffer() 中它也可以工作。我不明白的东西。

提前感谢您的任何启发。

【问题讨论】:

    标签: c# stream buffer


    【解决方案1】:

    我已经解决了我的问题。我用 byte[] 替换了所有 Stream 实例。现在它工作正常。 不知道为什么 Stream 不起作用,不想花更多时间解决这个问题。

    这是为遇到相同问题的任何人提供的新代码。

    /// <summary>
    /// This namespace provides a crossthread-, concurrentproof buffer manager. 
    /// </summary>
    namespace YariIfStream
    {
    
        /// <summary>
        /// A class that manages three buffers used for IF data streams
        /// </summary>
        public class YariIFStream
        {
            private byte[] writebuf; ///<value>The buffer used for writing</value>
            private byte[] readbuf; ///<value>The buffer used for reading</value>
            private byte[] swapbuf; ///<value>The buffer used for swapping</value>
            private bool firsttime; ///<value>Boolean used for checking if it is the first time a writebuffers is asked</value>
            private Object sync; ///<value>Object used for syncing</value>
    
            /// <summary>
            /// Initializes a new instance of the Yari.YariIFStream class with expandable buffers with a initial capacity as specified
            /// </summary>
            /// <param name="capacity">Initial capacity of the buffers</param>
            public YariIFStream(int capacity)
            {
                sync = new Object();
                firsttime = true;
    
                writebuf = new byte[capacity];
                readbuf = new byte[capacity];
                swapbuf = new byte[capacity];
            }
    
            /// <summary>
            /// Returns the buffer with new data ready to be read
            /// </summary>
            /// <returns>byte[]</returns>
            public byte[] GetReadBuffer()
            {
                byte[] tempbuf;
                lock (sync)
                {
                    Monitor.Wait(sync);
                    tempbuf = swapbuf;
                    swapbuf = readbuf;
                }
                readbuf = tempbuf;
    
                return readbuf;
            }
    
            /// <summary>
            /// Returns the buffer ready to be written with data
            /// </summary>
            /// <returns>byte[]</returns>
            public byte[] GetWriteBuffer()
            {
                byte[] tempbuf;
                lock (sync)
                {
                    tempbuf = swapbuf;
                    swapbuf = writebuf;
    
                    writebuf = tempbuf;
    
                    if (!firsttime)
                    {
                        Monitor.Pulse(sync);
                    }
                    else
                    {
                        firsttime = false;
                    }
                }
                return writebuf;
            }
        }
    }
    

    【讨论】:

    • 为什么是 4 个缓冲区?要交换,您只需要一个读取缓冲区、一个写入缓冲区和一个临时缓冲区。
    • Pulse 和 Wait 之间没有同步。 Monitor 类不维护指示 Pulse 方法已被调用的状态。因此,如果您在没有线程等待时调用 Pulse,则调用 Wait 的下一个线程会阻塞,就好像 Pulse 从未被调用过一样。
    猜你喜欢
    • 2020-10-24
    • 1970-01-01
    • 1970-01-01
    • 2017-04-07
    • 2016-06-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多