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
相关产品推荐
相关产品推荐

