【问题标题】:Handling multiple connections using QThreadPool使用 QThreadPool 处理多个连接
【发布时间】:2018-01-18 22:43:32
【问题描述】:

假设您需要与设备保持 256 个 tcp 连接,只是为了偶尔发送命令。我想并行执行此操作(它需要阻塞,直到它得到响应),我正在尝试使用 QThreadPool 来实现此目的,但我怀疑它是否可能。

我尝试使用 QRunnable,但我不确定套接字在线程之间的行为方式(套接字应该只在创建它们的线程中使用?)

我也担心这个解决方案的效率,如果有人可以提出一些替代方案,我会很高兴,不一定使用 QT。

下面我发布了一些sn-ps的代码。

class Task : public QRunnable {

    Task(){
        //creating TaskSubclass instance and socket in it
    }

private:
    TaskSubclass               *sub;

    void run() override {
        //some debug info and variable setting...
        sub->doSomething( args );
        return;
    }
};

class TaskSubclass {
    Socket         *sock;           // socket instance
    //...
    void doSomething( args )
    {
        //writing to socket here
    }
}

class MainProgram : public QObject{
    Q_OBJECT
private:
    QThreadPool *pool;
    Task *tasks;

public:
    MainProgram(){
        pool = new QThreadPool(this);
        //create tasks here
    }

    void run(){
        //decide which task to start
        pool->start(tasks[i]);
    }
};

【问题讨论】:

  • 请问您为什么要使用多个线程来解决问题? Qt 的事件系统能够通过一个线程轻松处理数百个 Tcp 连接,还是您的“套接字数据处理程序”逻辑阻塞了主线程?
  • 感谢您的评论,我刚刚编辑了帖子。它需要阻塞直到得到响应。
  • 你不使用信号槽机制有什么原因吗?您可以只使用 QTcpSocket 的 readyRead-Signal 并将其连接到插槽。
  • 如果我错了请纠正我,但在此解决方案中,插槽将按顺序执行,而不是并行执行。

标签: c++ multithreading qt sockets threadpool


【解决方案1】:

虽然答案已经被接受,但我想分享一下我的)

我从您的问题中了解到: 拥有 256 个当前活动连接,您不时向其中一个发送请求(如您命名的“命令”)他们并等待回应。同时,你想让这个进程成为多线程的,虽然你说“它需要阻塞直到它得到响应”,我假设你暗示阻塞一个处理 request-response 进程的线程,但不是主线程。

如果我确实正确理解了这个问题,那么我建议使用 Qt 来做到这一点:

#include <functional>

#include <QObject>          // need to add "QT += core" in .pro
#include <QTcpSocket>       // QT += network
#include <QtConcurrent>     // QT += concurrent 
#include <QFuture>           
#include <QFutureWatcher>

class CommandSender : public QObject
{
public:
    // Sends a command via connection and blocks 
    // until the response arrives or timeout occurs
    // then passes the response to a handler
    // when the handler is done - unblocks
    void SendCommand(
        QTcpSocket* connection,
        const Command& command,
        void(*responseHandler)(Response&&))
    {
        const int timeout = 1000;   // milliseconds, set it to -1 if you want no timeouts 

        // Sending a command (blocking)
        connection.write(command.ToByteArray());    // Look QByteArray for more details
        if (connection.waitForBytesWritten(timeout) {
            qDebug() << connection.errorString() << endl;
            emit error(connection);
            return;
        }

        // Waiting for a response (blocking)
        QDataStream in{ connection, QIODevice::ReadOnly };
        QString message;
        do {
            if (!connection.waitForReadyRead(timeout)) {
                qDebug() << connection.errorString() << endl;
                emit error(connection);
                return;
            }
            in.startTransaction();
            in >> message;
        } while (!in.commitTransaction());

        responseHandler(Response{ message }); // Translate message to a response and handle it
    }

    // Non-blocking version of SendCommand
    void SendCommandAsync(
        QTcpSocket* connection,
        const Command& command,
        void(*responseHandler) (Response&&))
    {
        QFutureWatcher<void>* watcher = new QFutureWatcher<void>{ this };
        connect(watcher, &QFutureWatcher<void>::finished, [connection, watcher] ()
        {
           emit done(connection);
           watcher->deleteLater();
        });

        // Does not block,
        // emits "done" when finished
        QFuture<void> future
            = QtConcurrent::run(this, &CommandSender::SendCommand, connection, command, responseHandler);
        watcher->setFuture(future);
    }

signals:
    void done(QTcpSocket* connection);
    void error(QTcpSocket* connection);
}

现在您可以使用从线程池中提取的单独线程向套接字发送命令:在后台QtConcurrent::run() 使用Qt 为您提供的QThreadPool 的全局实例。该线程阻塞,直到它得到响应,然后用 responseHandler 处理它。同时,管理所有命令和套接字的主线程保持畅通。只需捕获 done() 信号,该信号表明响应已成功接收并处理。

需要注意的一点:异步版本只有在线程池中有空闲线程时才发送请求,否则等待。当然,这是任何线程池的行为(这正是这种模式的重点),但不要忘记这一点。

另外,我在编写代码时没有使用 Qt,因此可能包含一些错误。

编辑:事实证明,这不是线程安全的,因为套接字在 Qt 中不可重入。

您可以做的是将互斥锁与套接字相关联,并在每次执行其功能时锁定它。这可以很容易地在 QTcpSocket 类周围创建一个包装器。如果我错了,请纠正我。

【讨论】:

  • 谢谢,看起来不错,但是如果我们在 QtConcureent.run() 中传递指向套接字的指针,线程安全应该不会有问题吗?我听说套接字只能在创建它们的线程中使用,并且在不同的线程之间传递它们是不安全的(我不知道应该在哪里创建套接字,因为主线程会不安全?)。如果我错了或者你有一些解释说这是安全的,请纠正我。
  • 你是对的,它不是线程安全的。原因是 - Qt 中的套接字不可重入(我忘了,抱歉)。所以我遇到了我尝试从主线程管理套接字连接/断开连接并读/写另一个工作线程的情况。我的错。 但是它得出了一个合乎逻辑的结论,即 所有 套接字上的操作无论如何都必须在同一个线程中执行(至少在 Qt 中)。
【解决方案2】:

正如 OMD_AT 已经指出的那样,最好的解决方案是使用 Select() 并让内核为您完成工作:-)

这里有一个异步方法和同步多线程方法的示例。

在这个例子中,我们创建了 10 个到 google 网络服务的连接并向服务器发出一个简单的 get 请求,我们测量每种方法中的所有连接需要多长时间才能接收到来自 google 服务器的响应。

请注意,您应该使用更快的网络服务器来进行真正的测试,例如 localhost,因为网络延迟对结果有很大影响。

#include <QCoreApplication>
#include <QTcpSocket>
#include <QtConcurrent/QtConcurrentRun>
#include <QElapsedTimer>
#include <QAtomicInt>

class Task : public QRunnable
{
    public:
        Task() : QRunnable() {}
        static QAtomicInt counter;
        static QElapsedTimer timer;
        virtual void run() override
        {
            QTcpSocket* socket = new QTcpSocket();
            socket->connectToHost("www.google.com", 80);
            socket->write("GET / HTTP/1.1\r\nHost: www.google.com\r\n\r\n");
            socket->waitForReadyRead();
            if(!--counter) {
                 qDebug("Multiple Threads elapsed: %lld nanoseconds", timer.nsecsElapsed());
            }
        }
};

QAtomicInt Task::counter;
QElapsedTimer Task::timer;

int main(int argc, char *argv[])
{
    QCoreApplication app(argc, argv);

    // init
    int connections = 10;
    Task::counter = connections;
    QElapsedTimer timer;

    /// Async via One Thread (Select)

    // handle the data
    auto dataHandler = [&timer,&connections](QByteArray data) {
        Q_UNUSED(data);
        if(!--connections) qDebug("  Single Threads elapsed: %lld nanoseconds", timer.nsecsElapsed());
    };

    // create 10 connection to google.com and send an http get request
    timer.start();
    for(int i = 0; i < connections; i++) {
        QTcpSocket* socket = new QTcpSocket();
        socket->connectToHost("www.google.com", 80);
        socket->write("GET / HTTP/1.1\r\nHost: www.google.com\r\n\r\n");
        QObject::connect(socket, &QTcpSocket::readyRead, [dataHandler,socket]() {
            dataHandler(socket->readAll());
        });
    }


   /// Async via Multiple Threads

    Task::timer.start();
    for(int i = 0; i < connections; i++) {
        QThreadPool::globalInstance()->start(new Task());
    }

    return app.exec();
}

打印:

Multiple Threads elapsed: 62324598 nanoseconds
  Single Threads elapsed: 63613967 nanoseconds

【讨论】:

  • 正如我在上面的 cmets 中提到的,我真的很想利用计算机中的许多内核,这些内核将专门用于这个软件。
  • 在主线程中创建套接字并在不同的线程中对其进行一些操作是否安全?另外,为了提高效率,您是否能够比较此解决方案和 select() 方法的性能? (区别大吗)
  • 是的,你是对的,它不是线程保存,我已经让它线程保存... Qt 在后台使用 Select 主事件循环只是休眠,直到来自内核的事件发送到应用程序.如果你想绕过 Select,你可以使用像 waitForReadyRead() 或 waitForBytesWritten() 这样的同步数据访问。是的,在这种情况下,如果您使用多个连接,则需要线程,因为线程正在阻塞...
  • 好的,现在它是线程安全的,但它不像普通的顺序方法吗?我没有看到在这里使用 QtConcurrent 的意义,因为关于写入和读取套接字的所有事情(信号和 QtConcurrent 调用)都是在一个线程中完成的,并且只能同时处理数据。在我的情况下,没有理由只使用 QtConcurrent 来处理数据。如果我错了,请纠正我。
  • 是的,您不需要 QtConcurrent,但您可以使用它来处理您在不同线程中收到的数据,我只是想向您展示...我添加了一个阻塞线程我的例子的方法这应该是你想要的。
【解决方案3】:

对于这个问题,我最喜欢的解决方案是使用select() 多路复用您的套接字。这样您就不需要创建额外的线程,这是一种“非常 POSIX”的方式。

例如看这个教程:

http://www.binarytides.com/multiple-socket-connections-fdset-select-linux/

或相关问题:

Using select(..) on client

【讨论】:

  • 我希望通过使用 QTcpSockets 或 boost 套接字等包装器将其保持在更高级别的抽象上,因此我不需要在原始内存上进行操作。
  • 在考虑我可以从那些包装中获取原始描述符只是为了使用 select() 但我想让它平行。它仍然是不错的选择,谢谢。
猜你喜欢
  • 2015-06-01
  • 2012-12-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-06-18
  • 2019-02-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多