【问题标题】:producer/consumer get different data by Queue生产者/消费者通过队列获取不同的数据
【发布时间】:2015-10-08 01:04:49
【问题描述】:

解决了!!!谢谢

注意你把什么样的对象放入Queue。如果你放入一个值,比如int,那么enqueue会做一个副本,每个人都很高兴。如果你把一个引用,比如 byte[], string, enqueue 把这个引用放到队列中,那么问题就来了。如果在消费者读取它之前更改了此引用,则消费者将读取 changed 版本的数据。

为避免此问题,在入队后立即在 rxThread 中获取新版本的帧引用。代码:

public void rxThreadFunc()
        {
         byte[] data = new byte[datalen];//declare for the first iteration.
            int j = 0;
            while (true)
            {

                for (int i = 0; i < rxlen; i++)
                {
                    data[j] = (byte)i;
                    j++;
                    if (j >= datalen)
                    {
                        j = 0;
                        mQ.Add(data);
                        using (StreamWriter fwriter = new StreamWriter("C:\\testsave\\rxdata", true))
                        {
                            for (int k = 0; k < datalen; k++)
                            {
                                fwriter.Write(data[k]);
                                fwriter.Write(",");
                            }
                            fwriter.Write("\n");
                        }
                     data = new byte[datalen];//create new reference after enqueue/Add
                    }
                }

            }
        }//rxThreadFunc()

更新1 我刚刚写了另一个更简单的代码,这样每个人都可以在没有串口硬件的情况下测试它。单击按钮运行程序。 我认为线程优先级会导致这个问题,在不改变 rxthread 的线程优先级的情况下,dataProc 线程将获得正确的数据。但我还是不知道为什么。

rxThread.Priority=ThreadPriority.Hightest

using System;
using System.Collections.Generic;
using System.ComponentModel;
using System.Data;
using System.Drawing;
using System.Linq;
using System.Text;
using System.Windows.Forms;
using System.Collections.Concurrent;
using System.Threading;
using System.IO;
namespace WindowsFormsApplication1
{
    public partial class Form1 : Form
    {
        public Form1()
        {
            InitializeComponent();

        }

        private void button1_Click(object sender, EventArgs e)
        {
            Thread rxThread = new Thread(rxThreadFunc);
            rxThread.Priority = ThreadPriority.Highest;//this causes problem
            rxThread.Start();
            Thread procThread = new Thread(dataProc);
            procThread.Start();
        }
        BlockingCollection<byte[]> mQ = new BlockingCollection<byte[]>();
        int datalen = 30;         
        int rxlen = 200;
        public void rxThreadFunc()
        {
            int j = 0;
            while (true)
            {
               byte[] data = new byte[datalen];//is this in the right place?
                for (int i = 0; i < rxlen; i++)
                {
                    data[j] = (byte)i;
                    j++;
                    if (j >= datalen)
                    {
                        j = 0;
                        mQ.Add(data);
                        using (StreamWriter fwriter = new StreamWriter("C:\\testsave\\rxdata", true))
                        {
                            for (int k = 0; k < datalen; k++)
                            {
                                fwriter.Write(data[k]);
                                fwriter.Write(",");
                            }
                            fwriter.Write("\n");
                        }
                    }
                }

            }
        }//rxThreadFunc()
        public void dataProc()
        {
            byte[] outData = new byte[datalen];
            while (true)
            {
                if (mQ.Count > 1)
                {
                    outData=mQ.Take();
                    using(StreamWriter fwriter=new StreamWriter("C:\\testsave\\dataProc",true))
                    {
                        for (int i = 0; i < datalen; i++)
                        {
                            fwriter.Write(outData[i]);
                            fwriter.Write(",");
                        }
                        fwriter.Write("\n");
                    }
                }
            }
        }

    }
}


说明: 我编写了这个包含两个线程的应用程序。 RxThread从串口接收数据,按照header mark 0x55 0xaa排序,将后面30个字节放入一个FrameStruct类中,然后将此FrameStruct放入Queue。 dataProcess 线程从 Queue 中获取帧,然后将其存储到磁盘。

SerialPort=(rxbuff)=>RxThread=(rxFrame,Queue)=>dataProcess==>磁盘

问题: dataProcess 线程接收并保存到磁盘的数据以某种方式损坏。

试过了: 这是我尝试过的,供您参考。

  1. 我尝试了 BlockedCollection,它自然是线程安全的,而不是队列,它仍然不起作用。所以我想这不是队列的问题。
  2. 我尝试在 FrameStruct 中添加另一个成员 int cnt,它在 RxThread 中的 messageQ.Enqueue() 之前自增。然后dataProcess线程可以正确获取它。所以我认为可能是 data[] 有问题,但是...
  3. 但我尝试将字节数据 [30] 而不是 FrameStruct 放入队列中,但不起作用。
  4. 另外我认为 dataProcess 收到的时间也是正确的。

    5.如果我在 RxThread 中的 Monitor.Pulse() 之后放置一个 Thread.sleep(20),问题就解决了,但我不明白为什么???如果我换到另一台电脑怎么办?

这是代码快照。

//declared:
//Queue<FrameStruct>messageQ=new Queue<FrameStruct>;
//object _LockerMQ=new object();
private void RxThread()
    {
        int bytestoread, i;
        bool f55 = false;//55 flag
        bool fs = false;//frame start flag
        int j=0;//data index in FrameStruct
        int m_lMaxFram=32;
        bytestoread = 0;
        FrameStruct rxFrame = new FrameStruct((int)m_lMaxFrame);
        while (true)
        {
            if (Serial_Port.IsOpen == true)
            {
                if ((bytestoread = Serial_Port.BytesToRead) > m_lMaxFrame*2)//get at least two frames
                {
                        rxbuff = new byte[bytestoread];  
                        Serial_Port.Read(rxbuff, 0, bytestoread);
                        for (i = 0; i < bytestoread; i++)
                        {
                            if (rxbuff[i] == 0x55)
                            {
                                f55 = true;
                                continue;
                            }
                            if (rxbuff[i] == 0xaa && f55)
                            {//frame header 0x55, 0xaa
                                //new frame start
                                fs = true;
                                f55 = false;
                                j = 0;//rxframe index;
                                rxFrame.time = DateTime.Now;//store the datetime when this thread gets this frame
                                continue;
                            }
                            if (fs && j < m_lMaxFrame - 2)
                            {//frame started but not ended
                                rxFrame.data[j] = rxbuff[i];
                                j++;
                            }
                            if (j >= (m_lMaxFrame - 2) && fs)
                            {//frame ended if j=30, reaches the end of rxFrame.data
                                fs = false;
                                lock(_LockerMQ)
                                {
                                    messageQ.Enqueue(rxFrame);
                                Monitor.Pulse(_LockerMQ);
                                }
                                 //Thread.Sleep(20);//if uncomment this sleep, problem solved
                              using (StreamWriter fWriter = new StreamWriter("c:\\testsave\\RXdata", true))//save rxThread result into a file rawdata
                    {
                        fWriter.Write(rxFrame.time.ToString("yyyy/MM/dd HH:mm:ss.fff"));
                        fWriter.Write(",");
                        for (int k = 0; k < m_lMaxFrame - 2; k++)
                        {
                            fWriter.Write(rxFrame.data[k]);
                            fWriter.Write(",");
                        }
                        fWriter.Write("\n");
                    }
                            }
                           }
                }//if ((bytestoread=Serial_Port.BytesToRead) > 0)
                rxbuff = null;
                Thread.Sleep(20);
            }//(Serial_Port.IsOpen==true)
            Thread.Sleep(100);
        }//while(true),RxThread sleep
    }//private void RxThread()

数据处理线程:

 public void dataProcess()
    {
      while (true)
        {
            lock (_LockerMQ)
            {
                while (messageQ.Count < 1) Monitor.Wait(_LockerMQ);//get at least one frame data
                f_NewFrame = messageQ.Count;
                if (f_NewFrame > 0)
                {
                    procFrame = messageQ.Dequeue();
                    using (StreamWriter fWriter = new StreamWriter("c:\\testsave\\dPdata", true))
                    {
                        fWriter.Write(procFrame.time.ToString("yyyy/MM/dd HH:mm:ss.fff"));
                        fWriter.Write(",");
                        for (int i = 0; i < m_lMaxFrame - 2; i++)
                        {
                            fWriter.Write(procFrame.data[i]);
                            fWriter.Write(",");
                        }
                        fWriter.Write("\n");
                    }
                }//if(f_NewFrame>0)

            }//lock(messageQ)
  }
}

FrameStruct 包含时间和数据的成员[30]

class FrameStruct
{
        public FrameStruct(int m_lMaxFrame)
        {
            time = DateTime.Now;
            data = new byte[m_lMaxFrame - 2];
        }
        public DateTime time;
        public volatile byte[] data;//volatile doesn't help
}

RxThread保存的rxData是正确的,显示:

2015/07/18 18:40:26.125,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,111,51,204,
2015/07/18 18:40:26.177,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,112,51,204,
2015/07/18 18:40:26.177,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,113,51,204,
2015/07/18 18:40:26.297,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,114,51,204,
2015/07/18 18:40:26.298,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,115,51,204,
2015/07/18 18:40:26.298,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,116,51,204,
2015/07/18 18:40:26.299,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,117,51,204,
2015/07/18 18:40:26.420,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,118,51,204,
                                                                                     //^this columns is accumulated number

dataProcessThread保存的dPdata是错误的,显示:

2015/07/18 18:40:31.904,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,227,51,204,
2015/07/18 18:40:31.905,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,228,51,204,
2015/07/18 18:40:31.905,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,229,51,204,
2015/07/18 18:40:32.026,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,231,51,204,
2015/07/18 18:40:32.026,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,231,51,204,
2015/07/18 18:40:32.147,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,232,51,204,
2015/07/18 18:40:32.148,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,233,51,204,
2015/07/18 18:40:32.148,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,234,51,204,
2015/07/18 18:40:32.269,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,236,51,204,
2015/07/18 18:40:32.269,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,236,51,204,
2015/07/18 18:40:32.510,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,237,51,204,
2015/07/18 18:40:32.512,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,240,51,204,
2015/07/18 18:40:32.512,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,240,51,204,
2015/07/18 18:40:32.514,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,240,51,204,
2015/07/18 18:40:32.514,127,255,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,241,51,204,
2015/07/18 18:40:32.635,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,243,51,204,
2015/07/18 18:40:32.635,128,0,255,255,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,243,51,204,
                                                                                   //^this accumulated number is not correct

请帮忙!

谢谢!

【问题讨论】:

  • 您确认问题出在多线程上吗?您是否分别测试了每个组件?
  • 是的,我单独测试过。我在 rxThread 中保存了 rawdata,它显示 OK。但是dataProcess中的rawdata1不行。
  • 这个人可能有类似的问题,但他的问题没有答案。link
  • 只是想弄清楚您要做什么:您是否希望作者在队列中有一个元素时立即完成他的工作?或者你可以允许队列有多个元素?
  • 我可以允许队列有多个元素。

标签: c# multithreading queue pass-by-reference thread-priority


【解决方案1】:

FrameStruct 是一个类(而不是结构),当您将它排入队列时,您会一遍又一遍地使用相同的引用。

【讨论】:

  • 谢谢伙计,但我尝试使用 byte[] 帧而不是 FrameStruct 作为队列,(Queue messageQ=new Queue 和 messageQ.Enqueue(frame) ),它不起作用。其实一开始我把FrameStruct声明为Struct,但是看到有人说C#里有class比struct好,所以我改成class了。
  • 没关系,FrameStruct 是引用类型,那么你必须在每次迭代中初始化一个新实例。当您的消费者从队列中选择一个项目时,生产者会修改相同项目的数据(如果下一次迭代与消费者并行执行)。我在这里可能错了,但你应该仔细检查一下。
  • 是的,你是对的,我解决了。非常感谢!
  • 我希望你能接受这个答案,但重要的是一切都为你解决了:D
【解决方案2】:

更新 2

这是基本的生产者消费者,只需添加您的逻辑:

class Program
{

    static BlockingCollection<int> mQ = new BlockingCollection<int>();

    static void Main(string[] args)
    {

        Thread rxThread = new Thread(rxThreadFunc);
        rxThread.Priority = ThreadPriority.Highest;//this causes problem
        rxThread.Start();
        Thread procThread = new Thread(dataProc);
        procThread.Start();
        Console.ReadLine();
    }


    static public void rxThreadFunc()
    {
        for (int i = 0; i < 10; i++)
        {
            mQ.Add(i);
        }
    }


    static public void dataProc()
    {
        foreach (int outData in mQ.GetConsumingEnumerable())
        {
            Console.WriteLine(outData);
        }
    }


}

更新 1 的答案:

  • 是的,这是正确的地方

现在,由于您正在使用 BlockingCollection,消费者 (dataProc) 可能非常简单(只需删除循环和计数检查,它会为您完成所有这些同步):

 foreach (byte[] outData in _taskQ.GetConsumingEnumerable())
 {
     using(StreamWriter fwriter=new StreamWriter("C:\\testsave\\dataProc",true))
                {
                    for (int i = 0; i < datalen; i++)
                    {
                        fwriter.Write(outData[i]);
                        fwriter.Write(",");
                    }
                    fwriter.Write("\n");
                }
 }

现在,问题可能与原始帖子中的不同。也许是因为这里生产者也将数据写入文件?

原文:

它有很多解释,但解决方案很简单,只需将 FrameStruct 的初始化添加到您的第二个“if”语句中:

rxFrame = new FrameStruct((int)m_lMaxFrame);

如您保存的数据文件所示:在损坏的文件中,每次当您缺少值(例如 230)时 - 您有其他值重复(例如 231)。缺失值的计数等于重复值的计数。

这样做的原因是您将对同一对象实例的引用添加到队列中。 让我们看看下面的场景:RxThread 循环 N 次,在它上下文切换到 dataProcess 线程之前,它向队列中添加了 N 个对同一 FrameStruct 实例的引用。此实例中的数据将是上下文切换之前最后一次读取循环迭代的数据。 现在发生了上下文切换:dataProcess 在上下文切换回 RxThread 之前循环了 M

现在,为什么 Thread.Sleep 有帮助。 简短的回答:每次 RxThread 将 1 个元素添加到队列中时,上下文切换到 dataProcess 线程的概率非常高。所以它实际上是:读一个 --> 上下文切换 --> 写一个......然后再做同样的事情。

长答案是:在 dataProcess 线程执行 Monitor.Wait 之后,它进入等待队列,上下文切换调度 RxThread。现在线程将第一个元素添加到队列并执行 Monitor.Pulse。这会将 dataProcess 线程移动到就绪队列。但不一定会安排它立即运行,因此 RxThread 可以进行另一次迭代。但是,如果您执行 Thread.Sleep - 现在很有可能会发生上下文切换和 ataProcess 线程。

【讨论】:

  • 谢谢我的朋友。我想我明白你的解释,请参阅 Update1, In rxThreadFunc(), while(true){byte[]data=new byte[datalen];....} 我把它放在正确的地点?还是不行。
  • 我根据与我的应用程序相同的逻辑编写了 Update1 代码,在我测试它们时它们共享相同的错误和权利。所以我像你说的那样将 byte[]data=new byte[datalen] 放在 while(true)loop 中(没有“if”语句)
  • Update 1 是原始应用程序的简化版本,我认为它们存在相同的问题,在 Update1 中我使用 Queue(byte[]) 其中 byte[] 也是一个参考。但是我像你说的那样把 data=new bytes[datalen] 放在 while(true) 循环中,它仍然不起作用。您只需将此代码复制到新的 c# WinForm 应用程序中即可。 Update1 的生产者和原始应用程序的生产者都将数据写入文件,以进行调试。
  • 这段代码的同步没有问题,我会再贴一次“minical prducer consumer case”的更新,添加你的逻辑即可。
  • 谢谢你的朋友,你是对的“将参考放入队列”。虽然我必须在入队后立即将 byte [] data=new data[len](或 rxFrame=new FrameStruct())放在该位置,但它可以确保下一帧在下一次迭代中获得新的新数据/rxFrame。
猜你喜欢
  • 2012-01-28
  • 2011-04-22
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多