【问题标题】:How to send signal/data from a worker thread to main thread?如何将信号/数据从工作线程发送到主线程?
【发布时间】:2015-05-09 00:26:55
【问题描述】:

首先我会说我是第一次研究多线程。尽管阅读了很多关于并发和同步的文章,但我并没有很容易地看到满足我所提出要求的解决方案。

使用 C++11 和 Boost,我试图弄清楚如何将数据从工作线程发送到主线程。工作线程在应用程序启动时产生并持续监控无锁队列。对象以不同的时间间隔填充此队列。这部分工作正常。

一旦数据可用,它需要由主线程处理,因为另一个信号将被发送到应用程序的其余部分,它不能在工作线程上。这就是我的我遇到了麻烦。

如果我必须通过互斥锁或条件变量阻塞主线程,直到工作线程完成,这将如何提高响应能力?我还不如只使用一个线程,这样我就可以访问数据了。我一定在这里遗漏了一些东西。

我已经发布了几个问题,认为 Boost::Asio 是要走的路。有一个如何在线程之间发送信号和数据的示例,但正如响应所示,事情很快变得过于复杂并且无法完美运行:

How to connect signal to boost::asio::io_service when posting work on different thread?

Boost::Asio with Main/Workers threads - Can I start event loop before posting work?

与一些同事交谈后,建议使用两个队列——一个输入,一个输出。这将在共享空间中,并且输出队列将由工作线程填充。工作线程一直在运行,但可能需要一个计时器,可能在应用程序级别,它会强制主线程检查输出队列以查看是否有任何待处理的任务。

关于我应该把注意力放在哪里有什么想法吗?是否有任何技术或策略可能适用于我正在尝试做的事情?接下来我会看看 Timers。

谢谢。

编辑:这是用于后处理模拟结果的插件系统的生产代码。我们尽可能首先使用 C++11,然后是 Boost。我们正在使用 Boost 的 lockfree::queue。应用程序在单个线程上执行我们想要的操作,但现在我们正在尝试优化我们发现存在性能问题的地方(在这种情况下,计算是通过另一个库进行的)。主线程有很多职责,包括数据库访问,这就是为什么我想限制工作线程实际做的事情。

更新:我已经成功地使用 std::thread 启动了一个工作线程,该线程检查 Boost lock::free 队列并处理放置它的任务。这是@Pressacco 中的第 5 步我遇到问题的回应。当工作线程完成并通知主线程而不是简单地等待工作线程完成时,是否有任何示例向主线程返回值?

【问题讨论】:

  • 主线程还在做什么?

标签: multithreading c++11 boost boost-asio boost-signals2


【解决方案1】:

如果您的目标是从头开始开发解决方案(使用本机线程、队列等):

  1. 创建线程保存队列队列(围绕添加/删除的 Mutex/CriticalSection)
  2. 创建与队列关联的计数信号量
  3. 有一个或多个工作线程等待计数信号量(即线程将阻塞)
    • 信号量比让线程不断轮询队列更有效
  4. 随着消息/作业被添加到队列中,增加信号量
    • 线程将被唤醒
    • 线程应该删除一条消息
  5. 如果需要返回结果...
    • 设置另一个:Queue+Semaphore+WorkerThreads

附加说明

如果您决定从头开始实现线程安全队列,请查看:

话虽如此,我会再看看 BOOST。我没有使用过这个库,但据我所知,它很可能包含一些相关的数据结构(例如线程安全队列)。

MSDN 中我最喜欢的一句话:

“当你使用任何类型的多线程时,你可能会暴露 自己面对非常严重和复杂的错误”

边栏

由于您是第一次研究并发编程,您不妨考虑:

  • 您的目标是构建具有生产价值的代码,还是只是一个学习练习?
    • 生产?考虑我们现有的经过验证的库
    • 学习?考虑从头开始编写代码
  • 考虑使用带有异步回调的线程池而不是本机线程。
  • 更多线程!=更好
  • 真的需要线程吗?
  • 关注KISS principle

【讨论】:

  • 嗨,这是第 5 步,我很难绕开我的脑袋。结果需要返回给主线程,所以当主线程停止做所有其他事情时,我如何通知它“你现在必须现在检查这个队列”?它不是另一个工作线程,我显然不能长时间阻塞主线程,这会使系统无响应。
  • 如果主线程有其他工作要做,就让它吧。只要确保主线程定期检查它的消息队列以查看是否有任何消息。此检查可以是非阻塞的:如果 queue.count > 0 then queue.remove message
  • 这就是它的症结所在... 1) 我如何进行这种定期检查?我能想到的只是应用程序级别的全局计时器,它定期调用此对象上的方法来检查队列。 2)您确实说过“消息”队列——我假设这个“消息”可能只是需要处理的数据实例? 3)当主线程正在处理它时,我是否必须在输出队列上设置一个锁/互斥锁,以防工作线程试图添加额外的项目?
【解决方案2】:

上面的反馈让我找到了我需要的正确方向。该解决方案绝对比我之前尝试过的必须使用信号/插槽或 Boost::Asio 更简单。我有两个无锁队列,一个用于输入(在工作线程上),一个用于输出(在主线程上,由工作线程填充)。我使用计时器来安排处理输出队列的时间。代码如下;也许它对某人有用:

//Task.h

#include <iostream>
#include <thread>


class Task
{
public:
   Task(bool shutdown = false) : _shutdown(shutdown) {};
   virtual ~Task() {};

   bool IsShutdownRequest() { return _shutdown; }

   virtual int Execute() = 0;

private:
   bool _shutdown;
};


class ShutdownTask : public Task
{
public:
   ShutdownTask() : Task(true) {}

   virtual int Execute() { return -1; }
};


class TimeSeriesTask : public Task
{
public:
   TimeSeriesTask(int value) : _value(value) {};

   virtual int Execute()
   {
      std::cout << "Calculating on thread " << std::this_thread::get_id() << std::endl;
      return _value * 2;
   }

private:
   int _value;
};


// Main.cpp : Defines the entry point for the console application.

#include "stdafx.h"
#include "afxwin.h"

#include <boost/lockfree/spsc_queue.hpp>

#include "Task.h"

static UINT_PTR ProcessDataCheckTimerID = 0;
static const int ProcessDataCheckPeriodInMilliseconds = 100;


class Manager
{
public:
   Manager() 
   {
      //Worker Thread with application lifetime that processes a lock free queue
      _workerThread = std::thread(&Manager::ProcessInputData, this);
   };

   virtual ~Manager() 
   {
      _workerThread.join();
   };

   void QueueData(int x)
   {
      if (x > 0)
      {
         _inputQueue.push(std::make_shared<TimeSeriesTask>(x));
      }
      else
      {
         _inputQueue.push(std::make_shared<ShutdownTask>());
      }
   }

   void ProcessOutputData()
   {
      //process output data on the Main Thread
      _outputQueue.consume_one([&](int value)
      {
         if (value < 0)
         {
            PostQuitMessage(WM_QUIT);
         }
         else
         {
            int result = value - 1;
            std::cout << "Final result is " << result << " on thread " << std::this_thread::get_id() << std::endl;
         }
      });
   }

private:
   void ProcessInputData()
   {
      bool shutdown = false;

      //Worker Thread processes input data indefinitely
      do
      {
         _inputQueue.consume_one([&](std::shared_ptr<Task> task)
         {    
            std::cout << "Getting element from input queue on thread " << std::this_thread::get_id() << std::endl;           

            if (task->IsShutdownRequest()) { shutdown = true; }

            int result = task->Execute();
            _outputQueue.push(result);
         });

      } while (shutdown == false);
   }

   std::thread _workerThread;
   boost::lockfree::spsc_queue<std::shared_ptr<Task>,   boost::lockfree::capacity<1024>> _inputQueue;
   boost::lockfree::spsc_queue<int, boost::lockfree::capacity<1024>> _outputQueue;
};


std::shared_ptr<Manager> g_pMgr;


//timer to force Main Thread to process Manager's output queue
void CALLBACK TimerCallback(HWND hWnd, UINT nMsg, UINT nIDEvent, DWORD dwTime)
{
   if (nIDEvent == ProcessDataCheckTimerID)
   {
      KillTimer(NULL, ProcessDataCheckPeriodInMilliseconds);
      ProcessDataCheckTimerID = 0;

      //call function to process data
      g_pMgr->ProcessOutputData();

      //reset timer
      ProcessDataCheckTimerID = SetTimer(NULL, ProcessDataCheckTimerID, ProcessDataCheckPeriodInMilliseconds, (TIMERPROC)&TimerCallback);
   }
}


int main()
{
   std::cout << "Main thread is " << std::this_thread::get_id() << std::endl;

   g_pMgr = std::make_shared<Manager>();

   ProcessDataCheckTimerID = SetTimer(NULL, ProcessDataCheckTimerID, ProcessDataCheckPeriodInMilliseconds, (TIMERPROC)&TimerCallback);

   //queue up some dummy data
   for (int i = 1; i <= 10; i++)
   {
      g_pMgr->QueueData(i);
   }

   //queue a shutdown request
   g_pMgr->QueueData(-1);

   //fake the application's message loop
   MSG msg;
   bool shutdown = false;
   while (shutdown == false)
   {
      if (GetMessage(&msg, NULL, 0, 0))
      {
         TranslateMessage(&msg);
         DispatchMessage(&msg);
      }
      else   
      {
         shutdown = true;
      }
   }

   return 0;
}

【讨论】:

  • 我很高兴我的 cmets 提供了帮助。仅供参考:等待事件比 while(true) 更有效,我会确保您的队列是线程安全的。否则,坏事会意外发生。 (我从未使用过提升队列。)
猜你喜欢
  • 1970-01-01
  • 2018-04-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-30
  • 1970-01-01
  • 2014-08-05
  • 1970-01-01
相关资源
最近更新 更多