不确定我的想法是否 100% 正确,但为什么不将工作线程数组传递给超级线程,保留一个表示当前活动线程索引的索引,仅将信号连接到那个,调度需要时发出信号,等待完成,断开信号,推进索引并重复?如果将序列化信号分派给线程是您真正想要的,这可能会起作用。
编辑
好吧,我真的拼命制作了一个基于 Qt 的示例,该示例实现了所需的工作流程并将其放在 github 上。
#pragma once
#include <QThread>
#include <QApplication>
#include <QMetaType>
#include <QTimer>
#include <vector>
#include <memory>
#include <cstdio>
#include <algorithm>
struct Work
{
int m_work;
};
struct Result
{
int m_result;
int m_workerIndex;
};
Q_DECLARE_METATYPE(Work);
Q_DECLARE_METATYPE(Result);
class Worker : public QThread
{
Q_OBJECT
public:
Worker(int workerIndex) : m_workerIndex(workerIndex)
{
moveToThread(this);
connect(this, SIGNAL(WorkReceived(Work)), SLOT(DoWork(Work)));
printf("[%d]: worker %d initialized\n", reinterpret_cast<int>(currentThreadId()), workerIndex);
}
void DispatchWork(Work work)
{
emit WorkReceived(work);
}
public slots:
void DoWork(Work work)
{
printf("[%d]: worker %d received work %d\n", reinterpret_cast<int>(currentThreadId()), m_workerIndex, work.m_work);
msleep(100);
Result result = { work.m_work * 2, m_workerIndex };
emit WorkDone(result);
}
signals:
void WorkReceived(Work work);
void WorkDone(Result result);
private:
int m_workerIndex;
};
class Master : public QObject
{
Q_OBJECT
public:
Master(int workerCount) : m_activeWorker(0), m_workerCount(workerCount)
{
printf("[%d]: creating master thread\n", reinterpret_cast<int>(QThread::currentThreadId()));
}
~Master()
{
std::for_each(m_workers.begin(), m_workers.end(), [](std::unique_ptr<Worker>& worker)
{
worker->quit();
worker->wait();
});
}
public slots:
void Initialize()
{
printf("[%d]: initializing master thread\n", reinterpret_cast<int>(QThread::currentThreadId()));
for (int workerIndex = 0; workerIndex < m_workerCount; ++workerIndex)
{
auto worker = new Worker(workerIndex);
m_workers.push_back(std::move(std::unique_ptr<Worker>(worker)));
connect(worker, SIGNAL(WorkDone(Result)), SLOT(WorkDone(Result)));
worker->start();
}
m_timer = new QTimer();
m_timer->setInterval(500);
connect(m_timer, SIGNAL(timeout()), SLOT(GenerateWork()));
m_timer->start();
}
void GenerateWork()
{
Work work = { m_activeWorker };
printf("[%d]: dispatching work %d to worker %d\n", reinterpret_cast<int>(QThread::currentThreadId()), work.m_work, m_activeWorker);
m_workers[m_activeWorker]->DispatchWork(work);
m_activeWorker = ++m_activeWorker % m_workers.size();
}
void WorkDone(Result result)
{
printf("[%d]: received result %d from worker %d\n", reinterpret_cast<int>(QThread::currentThreadId()), result.m_result, result.m_workerIndex);
}
void Terminate()
{
m_timer->stop();
delete m_timer;
}
private:
int m_workerCount;
std::vector<std::unique_ptr<Worker>> m_workers;
int m_activeWorker;
QTimer* m_timer;
};
QtThreadExample.cpp:
#include "QtThreadExample.hpp"
#include <QTimer>
int main(int argc, char** argv)
{
qRegisterMetaType<Work>("Work");
qRegisterMetaType<Result>("Result");
QApplication application(argc, argv);
QThread masterThread;
Master master(5);
master.moveToThread(&masterThread);
master.connect(&masterThread, SIGNAL(started()), SLOT(Initialize()));
master.connect(&masterThread, SIGNAL(terminated()), SLOT(Terminate()));
masterThread.start();
// Set a timer to terminate the program after 10 seconds
QTimer::singleShot(10 * 1000, &application, SLOT(quit()));
application.exec();
masterThread.quit();
masterThread.wait();
printf("[%d]: master thread has finished\n", reinterpret_cast<int>(QThread::currentThreadId()));
return 0;
}
一般来说,解决方案实际上是不从主线程本身发出信号,而是从工作线程发出信号 - 这样您就可以为每个线程获得一个唯一的信号,并且可以在事件循环中异步处理工作,并且然后发出一个线程完成的信号。该示例可以并且应该根据您的需要进行重构,但总的来说,它使用索引和信号线程的想法在 Qt 中演示了生产者/消费者模式。我正在使用通用计时器在主线程(Master::m_timer)中生成工作 - 我猜在你的情况下,你将使用来自套接字、文件或其他东西的信号来生成工作。然后我在一个活动的工作线程上调用一个方法,该方法向工作线程的事件循环发出一个信号以开始执行工作,然后发出一个关于完成的信号。这是一般描述,请查看示例,尝试一下,如果您有任何后续问题,请告诉我。
如果您使用 Qt 对象,我想这会相当不错,但在传统意义上的消费者/生产者模式中,信号/插槽的东西实际上让生活变得更加艰难,一个标准的 C++11 管道,带有 @987654328 @ 和一个调用 condition_variable::notify_one() 的主线程和简单地等待条件变量的工作线程会更容易,但是 Qt 对所有 I/O 东西都有很好的包装器。因此,请尝试一下并做出决定。
下面是示例的示例输出,我猜线程日志表明达到了所需的效果:
还有一点需要注意,因为QApplication 本身运行一个事件循环,如果你没有 GUI,你实际上可以让你的所有 I/O 对象和主类都存在于主线程中并从那里发出信号,从而消除需要一个单独的主线程。当然,如果你有一个 GUI,你不会想用这些东西给它增加负担。