【问题标题】:Actor calculation model using boost::thread使用 boost::thread 的 Actor 计算模型
【发布时间】:2013-11-04 18:48:06
【问题描述】:

我正在尝试使用 boost::thread 在 C++ 上的线程上实现 Actor 计算模型。 但是程序在执行过程中抛出了奇怪的异常。异常不稳定,有时程序以正确的方式运行。

那里有我的代码:

actor.hpp

class Actor {

  public:
    typedef boost::function<int()> Job;

  private:
    std::queue<Job>             d_jobQueue;
    boost::mutex                d_jobQueueMutex;
    boost::condition_variable   d_hasJob;
    boost::atomic<bool>         d_keepWorkerRunning;
    boost::thread               d_worker;

    void workerThread();

  public:
    Actor();
    virtual ~Actor();

    void execJobAsync(const Job& job);

    int execJobSync(const Job& job);
};

actor.cpp

namespace {

int executeJobSync(std::string          *error,
                   boost::promise<int> *promise,
                   const Actor::Job     *job)
{
    int rc = (*job)();

    promise->set_value(rc);
    return 0;
}

}

void Actor::workerThread()
{
    while (d_keepWorkerRunning) try {
        Job job;
        {
            boost::unique_lock<boost::mutex> g(d_jobQueueMutex);

            while (d_jobQueue.empty()) {
                d_hasJob.wait(g);
            }

            job = d_jobQueue.front();
            d_jobQueue.pop();
        }

        job();
    }
    catch (...) {
        // Log error
    }
}

void Actor::execJobAsync(const Job& job)
{
    boost::mutex::scoped_lock g(d_jobQueueMutex);
    d_jobQueue.push(job);
    d_hasJob.notify_one();
}

int Actor::execJobSync(const Job& job)
{
    std::string error;
    boost::promise<int> promise;
    boost::unique_future<int> future = promise.get_future();

    {
        boost::mutex::scoped_lock g(d_jobQueueMutex);
        d_jobQueue.push(boost::bind(executeJobSync, &error, &promise, &job));
        d_hasJob.notify_one();
    }

    int rc = future.get();

    if (rc) {
        ErrorUtil::setLastError(rc, error.c_str());
    }

    return rc;
}

Actor::Actor()
: d_keepWorkerRunning(true)
, d_worker(&Actor::workerThread, this)
{
}

Actor::~Actor()
{
    d_keepWorkerRunning = false;
    {
        boost::mutex::scoped_lock g(d_jobQueueMutex);
        d_hasJob.notify_one();
    }
    d_worker.join();
}

实际上抛出的异常是 int rc = future.get(); 行中的 boost::thread_interrupted。但是形成 boost 文档我不能解释这个例外。文档说

抛出: - boost::thread_interrupted 如果与 *this 相关的结果在调用时尚未准备好,并且当前线程被中断。

但是我的工作线程不能处于中断状态。

当我使用 gdb 并设置“catch throw”时,我看到回溯看起来像

抛出线程中断

boost::detail::interruption_checker::check_for_interruption

boost::detail::interruption_checker::interruption_checker

boost::condition_variable::wait

boost::detail::future_object_base::wait_internal

boost::detail::future_object_base::wait

boost::detail::future_object::get

boost::unique_future::get

我查看了 boost 源,但不明白为什么 interrupt_checker 决定工作线程被中断。

所以有人 C++ 大师,请帮助我。我需要做什么才能获得正确的代码? 我正在使用:

提升 1_53

Linux 版本 2.6.18-194.32.1.el5 Red Hat 4.1.2-48

gcc 4.7

编辑

修好了!感谢 Evgeny Panasyuk 和 Lazin。问题出在 TLS 中 管理。 boost::thread 和 boost::thread_specific_ptr 正在使用 相同的 TLS 存储用于其目的。就我而言,当 他们都试图在创建时更改此存储(不幸的是,我 不明白为什么会发生这种情况)。所以 TLS 被破坏了。

我将代码中的 boost::thread_specific_ptr 替换为 __thread 指定变量。

Offtop:在调试过程中,我发现外部库中的内存损坏 并修复它=)

.

编辑 2 我得到了确切的问题......这是GCC中的一个错误=) _GLIBCXX_DEBUG 编译标志会破坏 ABI。 你可以看到关于 boost bugtracker 的讨论: https://svn.boost.org/trac/boost/ticket/7666

【问题讨论】:

  • "在我的情况下,当他们都试图在创建时更改此存储时出现问题" - 很高兴看到一些代码显示了这一点。可能是 Boost 中的错误 - 所以我们应该报告它。
  • “我将代码中的 boost::thread_specific_ptr 替换为 __thread 指定变量。” - 那么,另一个 TLS 插槽会损坏吗? =)

标签: c++ multithreading boost actor


【解决方案1】:

我发现了几个错误:


Actor::workerThread 函数在d_jobQueueMutex 上进行双重解锁。第一个解锁是手动d_jobQueueMutex.unlock();,第二个是boost::unique_lock&lt;boost::mutex&gt;的析构函数。

您应该阻止其中一种解锁,例如unique_lockmutex 之间的release 关联:

g.release(); // <------------ PATCH
d_jobQueueMutex.unlock();

或者添加额外的代码块 + 默认构造的Job


workerThread 可能永远不会离开以下循环:

while (d_jobQueue.empty()) {
    d_hasJob.wait(g);
}

想象以下情况:d_jobQueue 为空,Actor::~Actor() 被调用,它设置标志并通知工作线程:

d_keepWorkerRunning = false;
d_hasJob.notify_one();

workerThread 在 while 循环中醒来,发现队列为空并再次休眠。

通常的做法是发送特殊的最终作业来停止工作线程:

~Actor()
{
    execJobSync([this]()->int
    {
        d_keepWorkerRunning = false;
        return 0;
    });
    d_worker.join();
}

在这种情况下,d_keepWorkerRunning 不需要是原子的。


LIVE DEMO on Coliru


编辑

我已将事件队列代码添加到您的示例中。

EventQueueImplActor 中都有并发队列,但类型不同。可以将公共部分提取到适用于任何类型的单独实体concurrent_queue&lt;T&gt;。在一个地方调试和测试队列比捕获分散在不同类中的错误要容易得多。

所以,你可以试试这个concurrent_queue&lt;T&gt;(on Coliru)

【讨论】:

  • 谢谢你,叶夫根尼。我删除了双重解锁(请参阅更新的代码)。但这并没有帮助=(仍然有同样奇怪的崩溃。我注意到当我在另一个类中添加另一个条件变量时出现了这个问题(类似于阻塞事件队列)。
  • @inkooboo 我已将现场演示添加到答案底部 - 您可以尝试使用它(也许添加其他部分)以重现崩溃。
  • 我已将事件队列代码添加到您的示例中。 coliru.stacked-crooked.com/a/1429e27abf63efa1 不幸的是我无法重现崩溃。但是您可以查看代码并可能有一些想法。谢谢。
  • 感谢您的帮助。仍然无法解决我的问题。在 SO 上支持您的其他答案。
  • @inkooboo 此外,您可以通过安全检查(无双重解锁或锁定)制作DebugMutex。顺便说一句,您是否仍然收到与原始问题相同的错误消息?
【解决方案2】:

这只是一个猜测。我认为有些代码实际上可以调用boost::tread::interrupt()。您可以为此函数设置断点并查看对此负责的代码。你可以在execJobSync测试中断:

int Actor::execJobSync(const Job& job)
{
    if (boost::this_thread::interruption_requested())
        std::cout << "Interruption requested!" << std::endl;
    std::string error;
    boost::promise<int> promise;
    boost::unique_future<int> future = promise.get_future();

在这种情况下最可疑的代码是引用线程对象的代码。

无论如何,让你的 boost::thread 代码中断是个好习惯。在某些范围内也可以disable interruption

如果不是这种情况 - 您需要检查与线程本地存储一起使用的代码,因为线程中断标志存储在 TLS 中。也许你的一些代码会重写它。您可以检查此类代码片段之前和之后的中断。

另一种可能是您的内存已损坏。如果没有代码调用 boost::thread::interrupt() 并且您不使用 TLS。这是最困难的情况,尝试使用一些动态分析器 - valgrind 或 clang 内存清理器。

题外话: 您可能需要使用一些并发队列。 std::queue 会因为高内存争用而变得非常慢,并且最终会导致缓存性能不佳。良好的并发队列允许您的代码并行地使元素入队和出队。

另外,actor 不是应该执行任意代码的东西。 Actor 队列必须接收简单的消息,而不是函数!您正在编写作业队列 :) 您需要查看一些演员系统,例如 Akkalibcpa

【讨论】:

  • 谢谢。经过一天的调试,我可以说它看起来像内存损坏。我将 break 设置为 boost::thread::interrupt() 并看到没有人调用它。同时,它也是一个可以设置“interrupt_requested”标志的地方。目前我正在尝试定位发生内存损坏的位置。
  • 您在使用 Visual Studio 吗?据我记得,它有一些可以检查堆一致性的宏定义。
  • 不,这是 linux 解决方案。我使用了 valgrind 并从外部库中检测到一些“无效写入”。幸运的是我有它的来源,所以可以解决这个问题。
猜你喜欢
  • 2012-11-13
  • 2011-09-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-07-11
  • 2013-10-05
  • 1970-01-01
相关资源
最近更新 更多