You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

C++ ThreadPool多线程运行挂起但单线程正常问题求助

问题分析与解决方案

问题根源

当线程数大于1时程序挂起、子进程残留的核心原因是非阻塞写导致的无限循环:

  • 多线程场景下,多个子进程cat无法及时读取管道数据,导致管道写缓冲区被占满。
  • 非阻塞模式下write会返回EAGAIN错误,代码直接break内部循环后进入下一次for循环,重复尝试写入,陷入无限循环。
  • 线程永远无法执行到close(pipefd[1]),子进程cat因未收到EOF持续运行,线程池无法完成join,最终程序挂起。

而线程数为1时,单个cat进程能及时处理管道数据,管道不会满,write顺利完成,流程正常结束。

解决方案

方案一:移除非阻塞写设置(最简单有效)

去掉写端的非阻塞标志,让write在管道满时阻塞,直到子进程读取数据后再继续写入,保证数据全部写入后正常结束循环,关闭写端触发子进程退出。

修改workerFunction中设置非阻塞的代码:

// 注释或删除以下两行非阻塞设置代码
// int flags = fcntl(pipefd[1], F_GETFL, 0);
// fcntl(pipefd[1], F_SETFL, flags | O_NONBLOCK);

方案二:正确处理非阻塞写的EAGAIN错误

如果必须保留非阻塞写,需在write返回EAGAIN时等待写端可写,而非直接break。可以用select等待写事件:

while (written < data.size() && running) {
    ssize_t bytes = write(pipefd[1], data.c_str() + written, data.size() - written);
    if (bytes == -1) {
        if (errno == EAGAIN || errno == EWOULDBLOCK) {
            // 等待写端可写,设置100ms超时
            fd_set write_fds;
            FD_ZERO(&write_fds);
            FD_SET(pipefd[1], &write_fds);
            struct timeval tv = {0, 100000};
            int ret = select(pipefd[1] + 1, nullptr, &write_fds, nullptr, &tv);
            if (ret == -1 && errno != EINTR) {
                perror("select");
                break;
            }
            continue; // 重新尝试写入
        } else {
            perror("write");
            break;
        }
    }
    written += bytes;
}

额外优化:处理信号中断的系统调用

当收到SIGINT时,running设为false,但线程可能卡在阻塞的write或waitpid中,需让系统调用中断后检查running状态:

  • 修改write的错误处理:
if (bytes == -1) {
    if (errno == EINTR) {
        // 被信号中断,检查running状态
        if (!running) break;
        continue;
    }
    perror("write");
    break;
}
  • 修改waitpid的调用,处理被信号中断的情况:
int status;
while (waitpid(pid, &status, 0) == -1) {
    if (errno == EINTR) continue;
    perror("waitpid");
    break;
}

修改后的完整代码(方案一)

#include <atomic>
#include <csignal>
#include <functional>
#include <iostream>
#include <mutex>
#include <string>
#include <thread>
#include <vector>
#include <unistd.h>
#include <sys/wait.h>
#include <fcntl.h>
#include <errno.h>
#include <condition_variable>
#include <queue>

using namespace std;

atomic<bool> running(true);
mutex coutMutex;

void handle_signal(int signal) {
    if (signal == SIGINT) {
        running = false;
    }
}

class ThreadPool {
    vector<thread> threads;
    queue<function<void()>> taskQueue;
    mutex queueMutex;
    condition_variable cv;
    atomic<bool> stop;

public:
    ThreadPool(size_t threadCount) : stop(false) {
        for (size_t i = 0; i < threadCount; ++i) {
            threads.emplace_back([this] {
                while (true) {
                    function<void()> task;
                    {
                        unique_lock<mutex> lock(queueMutex);
                        cv.wait(lock, [this] { return stop || !taskQueue.empty(); });
                        if (stop && taskQueue.empty()) return;
                        task = std::move(taskQueue.front());
                        taskQueue.pop();
                    }
                    task();
                }
            });
        }
    }

    ~ThreadPool() {
        {
            unique_lock<mutex> lock(queueMutex);
            stop = true;
        }
        cv.notify_all();
        for (thread& t : threads) {
            t.join();
        }
    }

    void enqueueTask(function<void()> task) {
        {
            unique_lock<mutex> lock(queueMutex);
            taskQueue.push(std::move(task));
        }
        cv.notify_one();
    }
};

void workerFunction(size_t start, size_t end, const vector<const char*>& args) {
    int pipefd[2];

    if (pipe(pipefd) < 0) {
        perror("pipe");
        return;
    }

    // 移除非阻塞写设置
    // int flags = fcntl(pipefd[1], F_GETFL, 0);
    // fcntl(pipefd[1], F_SETFL, flags | O_NONBLOCK);

    pid_t pid = fork();
    if (pid < 0) {
        perror("fork");
        return;
    }

    if (pid == 0) {
        // 子进程
        close(pipefd[1]);  // 关闭未使用的写端
        dup2(pipefd[0], STDIN_FILENO);
        close(pipefd[0]);

        execvp("cat", const_cast<char* const*>(args.data()));
        perror("execvp");
        exit(EXIT_FAILURE);
    } else {
        // 父进程
        close(pipefd[0]);  // 关闭未使用的读端

        string data = "Example data ";
        
        for (size_t i = start; i < end && running; ++i) {
            size_t written = 0;
            while (written < data.size() && running) {
                ssize_t bytes = write(pipefd[1], data.c_str() + written, data.size() - written);
                if (bytes == -1) {
                    if (errno == EINTR) {
                        // 被信号中断,检查running状态
                        if (!running) break;
                        continue;
                    }
                    perror("write");
                    break;
                }
                written += bytes;
            }
            if (!running) break;
        }
        {
            lock_guard<mutex> lock(coutMutex);
            cout << "Reached end!" << endl;
        }
        
        close(pipefd[1]); // 确保写端关闭
        int status;
        // 处理waitpid被信号中断的情况
        while (waitpid(pid, &status, 0) == -1) {
            if (errno == EINTR) continue;
            perror("waitpid");
            break;
        }

        if (WIFEXITED(status)) {
            lock_guard<mutex> lock(coutMutex);
            cout << "Worker finished. Exit code: " << WEXITSTATUS(status) << endl;
        } else if (WIFSIGNALED(status)) {
            lock_guard<mutex> lock(coutMutex);
            cerr << "Worker terminated by signal: " << WTERMSIG(status) << endl;
        }
    }
}


int main(int argc, char* argv[]) {
    signal(SIGINT, handle_signal);

    vector<const char*> args = {"cat", nullptr};

    size_t totalCombinations = 100;
    size_t threadCount = 4; // 现在可以设置大于1的线程数
    ThreadPool threadPool(threadCount);

    size_t chunkSize = totalCombinations / threadCount;
    size_t remainder = totalCombinations % threadCount;

    size_t start = 0;
    for (size_t i = 0; i < threadCount; ++i) {
        size_t end = start + chunkSize + (i < remainder ? 1 : 0);
        threadPool.enqueueTask([=, &args]() {
            workerFunction(start, end, args);
        });
        start = end;
    }

    return 0;
}

验证方法

编译运行修改后的代码:

g++ -std=c++17 -o thread_pipe thread_pipe.cpp -pthread
./thread_pipe

设置任意大于1的线程数(如4),程序会正常执行完毕,所有子进程cat会退出,无残留,程序不挂起。

内容的提问来源于stack exchange,提问作者GhostDog98

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 06:07:32