【问题标题】:Issue with Concurrent Access of a Queue (Multiple Producers and Consumers) - C++, Boost队列的并发访问问题(多个生产者和消费者) - C++,Boost
【发布时间】:2012-11-03 03:44:41
【问题描述】:

我正在编写一个具有事件队列的应用程序。我的意图是以这样一种方式创建它,即多个线程可以写入并且一个线程可以从队列中读取,并将弹出元素的处理交给另一个线程,以便后续再次弹出不会被阻塞。我使用了一个锁和一个条件变量来从队列中推送和弹出项目:

void Publisher::popEvent(boost::shared_ptr<Event>& event) {

    boost::mutex::scoped_lock lock(queueMutex);
    while(eventQueue.empty())
    {
        queueConditionVariable.wait(lock);
    }
    event = eventQueue.front();
    eventQueue.pop();
    lock.unlock(); 
}

void Publisher::pushEvent(boost::shared_ptr<Event> event) {

    boost::mutex::scoped_lock lock(queueMutex);
    eventQueue.push(event);
    lock.unlock();
    queueConditionVariable.notify_one();

}

在 Publisher 类的构造函数中(仅创建一个实例),我正在启动一个线程,该线程将遍历一个循环,直到捕获 notify_one(),然后启动另一个线程来处理弹出的事件从队列中:

在构造函数中:

publishthreadGroup = boost::shared_ptr<boost::thread_group> (new boost::thread_group());
publishthreadGroup->create_thread(boost::bind(queueProcessor, this));

queueProcessor 方法:

void queueProcessor(Publisher* agent) {

while(true) {
    boost::shared_ptr<Event> event;
    agent->getEvent(event);
    agent->publishthreadGroup->create_thread(boost::bind(dispatcher, agent, event));

    }
}

在dispatcher方法中,相关处理完成,处理后的信息通过thrift发布到服务器。在程序存在之前调用的另一个方法中,即在主线程中,我调用 join_all() 以便主线程等待线程完成。

在这个实现中,在为dispatcher创建线程之后,在上面的while循环中,我遇到了死锁/挂起。运行代码似乎卡住了。这个实现有什么问题?有没有更清洁、更好的方法来做我想做的事情? (多个生产者和一个消费者线程遍历队列并将一个元素的处理移交给不同的线程)

谢谢!

【问题讨论】:

  • 你能发布完整的例子吗?
  • 我假设在您的 queueProcessor 方法中您打算调用 agent-&gt;popEvent(event) 而不是 agent-&gt;getEvent(event)

标签: c++ boost concurrency producer-consumer


【解决方案1】:

似乎queueProcessor 函数将永远运行并且运行它的线程将永远不会退出。由该函数创建的任何线程都将完成它们的工作并退出,但是这个线程 - 在 publishthreadGroup 中创建的第一个线程 - 有一个无法停止的 while(true) 循环。因此,对join_all() 的调用将永远等待。您可以创建一些其他标志变量来触发该函数退出循环并返回吗?这应该可以解决问题!

【讨论】:

    猜你喜欢
    • 2013-04-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多