爱和冰阔落头像
关注
【Linux】从匿名管道到进程池:任务派发、fd 继承 Bug 与完整实现封面图

【Linux】从匿名管道到进程池:任务派发、fd 继承 Bug 与完整实现

封面

🔥个人主页:爱和冰阔乐
📚专栏传送门:《数据结构与算法》C++
🐶学习方向:C++方向学习爱好者
⭐人生格言:得知坦然 ,失之淡然

在这里插入图片描述


🏠博主简介
在这里插入图片描述

前言

上篇已经把匿名管道的基本规则理清了:父进程先调用 pipe(),再通过 fork() 让子进程继承对同一条管道的访问关系。

接下来把一条管道扩展成多条:

父进程
 ├── 管道 1 ── 子进程 1
 ├── 管道 2 ── 子进程 2
 ├── 管道 3 ── 子进程 3
 └── ...

父进程通过每条管道的写端发送任务码,子进程阻塞读取并执行任务,这就是一个最小的进程池模型。

代码本身不算难,真正容易出错的是 fd 继承

当父进程在循环中不断 fork() 时,后创建的子进程会把父进程当前已经持有的历史写端一起继承下来。如果这些写端没有关闭,父进程即使执行了 Stop(),某些子进程仍然可能永远读不到 EOF,最后把 waitpid() 卡住。

这篇就围绕这个坑,把进程池的设计、代码和退出流程完整走一遍。


一、先把进程池模型画出来

前面的匿名管道只有一个子进程。

接下来把场景扩展一下:

父进程
 ├── 管道 1 ── 子进程 1
 ├── 管道 2 ── 子进程 2
 ├── 管道 3 ── 子进程 3
 ├── 管道 4 ── 子进程 4
 └── 管道 5 ── 子进程 5

父进程通过每条管道的写端给对应子进程发送任务码,子进程阻塞等待任务。

1.1 为什么要提前创建子进程

传统方式是:

任务到达
   ↓
fork() 创建子进程
   ↓
执行任务
   ↓
子进程退出

如果任务很多,就需要频繁创建和销毁进程。

进程池的思路是:

程序启动
   ↓
提前创建 N 个子进程
   ↓
任务到达时挑选一个子进程
   ↓
子进程执行完成后继续等待下一次任务

这和 vector 预留容量的池化思想类似:把创建资源的成本提前支付,后面反复复用已经存在的资源

1.2 父进程到底给子进程发什么

父进程不需要把一整个函数通过管道发送给子进程。

父子进程使用的是同一份程序代码,所以双方可以提前约定:

任务码 0 → 打印日志
任务码 1 → 下载任务
任务码 2 → 上传任务

父进程只需要往管道中写一个 int,子进程读取后再从任务表中找到对应函数。管道传的是任务码,不是函数本身。

每次写入 sizeof(int),远小于通常的 PIPE_BUF,因此多个写者场景下也可以获得小块写入的原子性保证。不过本例中每条管道只有父进程一个写者,模型更加简单。

二、代码怎么拆

草稿中使用了 .hpp.cxx,这里顺手说明一下:

  • .hpp 通常表示 C++ 头文件;
  • .cxx.cc.cpp 都可以作为 C++ 源文件后缀;
  • .hpp 并不代表“必须把头文件和源文件混在一起”,它只是一个常见命名约定;
  • 模板或较小的类实现经常直接放在 .hpp 中。

下面把程序拆成:

Task.hpp
Channel.hpp
ProcessPool.hpp
main.cc

2.1 Task.hpp:管理任务

#pragma once

#include <iostream>
#include <vector>
#include <functional>
#include <cstdlib>
#include <ctime>
#include <unistd.h>

using task_t = std::function<void()>;

void PrintLog()
{
    std::cout
        << "子进程[" << getpid() << "]执行打印日志任务"
        << std::endl;
}

void Download()
{
    std::cout
        << "子进程[" << getpid() << "]执行下载任务"
        << std::endl;
}

void Upload()
{
    std::cout
        << "子进程[" << getpid() << "]执行上传任务"
        << std::endl;
}

class TaskManager
{
public:
    TaskManager()
    {
        std::srand(static_cast<unsigned int>(std::time(nullptr)));

        _tasks.push_back(PrintLog);
        _tasks.push_back(Download);
        _tasks.push_back(Upload);
    }

    int SelectTaskCode() const
    {
        return std::rand() % static_cast<int>(_tasks.size());
    }

    void Execute(int code)
    {
        if (code >= 0 && code < static_cast<int>(_tasks.size()))
        {
            _tasks[code]();
        }
    }

private:
    std::vector<task_t> _tasks;
};

任务表只在程序中注册一次。

父进程负责随机选择任务码,子进程负责根据任务码执行函数。

2.2 Channel.hpp:描述一条通信信道

父进程要管理的不只是一个 fd,还需要知道这条管道对应哪个子进程。

#pragma once

#include <iostream>
#include <string>
#include <unistd.h>
#include <sys/types.h>
#include <sys/wait.h>

class Channel
{
public:
    Channel(int wfd, pid_t child_id)
        : _wfd(wfd),
          _child_id(child_id),
          _name(
              "channel-" +
              std::to_string(wfd) +
              "-" +
              std::to_string(child_id)
          )
    {}

    const std::string& Name() const
    {
        return _name;
    }

    int WriteFd() const
    {
        return _wfd;
    }

    pid_t ChildId() const
    {
        return _child_id;
    }

    bool Send(int task_code) const
    {
        ssize_t n = write(_wfd, &task_code, sizeof(task_code));
        return n == static_cast<ssize_t>(sizeof(task_code));
    }

    void Close()
    {
        if (_wfd >= 0)
        {
            close(_wfd);
            _wfd = -1;
        }
    }

    void Wait() const
    {
        pid_t ret = waitpid(_child_id, nullptr, 0);

        if (ret > 0)
        {
            std::cout
                << "回收子进程:" << _child_id
                << std::endl;
        }
    }

private:
    int _wfd;
    pid_t _child_id;
    std::string _name;
};

这就是先描述

多个 Channel 再放进 vector,就是再组织

三、ProcessPool.hpp:创建、派发和回收

#pragma once

#include <iostream>
#include <vector>
#include <unistd.h>
#include <sys/types.h>
#include "Task.hpp"
#include "Channel.hpp"

const int DEFAULT_PROCESS_NUM = 5;

class ProcessPool
{
public:
    explicit ProcessPool(int process_num = DEFAULT_PROCESS_NUM)
        : _process_num(process_num),
          _next(0)
    {}

    bool Start()
    {
        for (int i = 0; i < _process_num; ++i)
        {
            int pipefd[2] = {0};

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

            pid_t child_id = fork();

            if (child_id < 0)
            {
                perror("fork");
                close(pipefd[0]);
                close(pipefd[1]);
                return false;
            }
            else if (child_id == 0)
            {
                /*
                 * 这个细节非常重要:
                 * 当前子进程会继承父进程此前保存的所有写端。
                 * 如果不关闭这些历史写端,父进程将来关闭写端后,
                 * 某些子进程仍可能因为持有额外写端而一直读不到 EOF。
                 */
                for (auto& channel : _channels)
                {
                    close(channel.WriteFd());
                }

                close(pipefd[1]);

                Worker(pipefd[0]);

                close(pipefd[0]);
                _exit(0);
            }

            close(pipefd[0]);

            _channels.emplace_back(pipefd[1], child_id);
        }

        return true;
    }

    void RunOnce()
    {
        if (_channels.empty())
        {
            return;
        }

        int task_code = _task_manager.SelectTaskCode();
        Channel& channel = SelectChannel();

        std::cout
            << "选择子进程:" << channel.Name()
            << ",发送任务码:" << task_code
            << std::endl;

        if (!channel.Send(task_code))
        {
            perror("write task code");
        }
    }

    void Stop()
    {
        /*
         * 先关闭父进程持有的所有写端。
         * 子进程把管道中剩余任务读完以后,read() 返回 0 并退出。
         */
        for (auto& channel : _channels)
        {
            std::cout
                << "关闭写端:" << channel.Name()
                << std::endl;

            channel.Close();
        }

        /*
         * 再统一回收所有子进程。
         */
        for (const auto& channel : _channels)
        {
            channel.Wait();
        }
    }

private:
    void Worker(int rfd)
    {
        while (true)
        {
            int task_code = 0;

            ssize_t n = read(rfd, &task_code, sizeof(task_code));

            if (n == static_cast<ssize_t>(sizeof(task_code)))
            {
                std::cout
                    << "子进程[" << getpid()
                    << "]收到任务码:" << task_code
                    << std::endl;

                _task_manager.Execute(task_code);
            }
            else if (n == 0)
            {
                std::cout
                    << "子进程[" << getpid()
                    << "]读取到 EOF,准备退出"
                    << std::endl;

                break;
            }
            else if (n < 0)
            {
                perror("read task code");
                break;
            }
            else
            {
                std::cerr << "读取到不完整任务码,子进程退出" << std::endl;
                break;
            }
        }
    }

    Channel& SelectChannel()
    {
        Channel& channel = _channels[_next];

        ++_next;
        _next %= _channels.size();

        return channel;
    }

private:
    int _process_num;
    std::size_t _next;
    std::vector<Channel> _channels;
    TaskManager _task_manager;
};

四、最容易漏掉的 Bug:历史写端继承

这段代码里最关键的不是轮询算法,而是关闭历史写端

假设父进程已经创建了子进程 1 和管道 1,接着再 fork() 创建子进程 2,那么子进程 2 会继承父进程当前持有的管道 1 写端。

如果不把这个历史写端关掉,后面父进程即使关闭了管道 1 写端,子进程 1 也可能因为系统中仍有其他写端引用,而一直读不到 EOF。

这就是草稿中“藏得比较深的系统级错误”真正对应的地方。

核心踩坑:关闭无关 fd,是多进程管道程序能否正确退出的关键。

4.1 Bug 是怎么出现的

创建第 0 个子进程时,父进程手里只有管道 0 的写端。

创建第 1 个子进程之前,父进程已经把管道 0 的写端保存进 _channels。这时再次 fork(),子进程 1 会把父进程当前持有的 fd 一起继承,因此它不只拿到自己的管道,还拿到了管道 0 的写端。

继续创建下去:

子进程 0:持有自己的读端
子进程 1:额外继承管道 0 写端
子进程 2:额外继承管道 0、1 写端
子进程 3:额外继承管道 0、1、2 写端

如果不清理这些历史写端,父进程关闭自己的写端后,内核仍然能看到其他子进程持有写端引用。只关父进程自己的 fd 并不等于系统里已经没有写端。

结果就是:

子进程 read()
    ↓
管道暂时没有数据
    ↓
但写端引用计数不为 0
    ↓
不能返回 EOF
    ↓
子进程一直阻塞
    ↓
父进程 waitpid() 一直等

这个现象看起来像“进程池退出死锁”,本质上其实是fd 没关干净

4.2 修复方式

每个新创建的子进程进入分支后,先把父进程此前保存的历史写端全部关闭:

for (auto& channel : _channels)
{
    close(channel.WriteFd());
}

然后再关闭当前管道的写端,只保留自己的读端:

close(pipefd[1]);
Worker(pipefd[0]);

这样父进程最终关闭所有写端后,每个子进程都能在数据读完时收到 EOF,随后正常退出。

五、main.cc:把整个流程跑起来

#include "ProcessPool.hpp"

int main()
{
    ProcessPool pool(5);

    if (!pool.Start())
    {
        std::cerr << "进程池启动失败" << std::endl;
        return 1;
    }

    int count = 10;

    while (count--)
    {
        pool.RunOnce();
        sleep(1);
    }

    pool.Stop();

    return 0;
}

主进程完成三件事:

Start()   → 创建管道和子进程
RunOnce() → 轮询选择子进程并发送任务码
Stop()    → 关闭写端并回收子进程

子进程则一直阻塞在自己的读端:

收到任务码 → 执行任务 → 继续等待
读取到 EOF → 退出

这样,一个基于匿名管道的简单进程池就完成了。创建、派发、退出三条链路都已经闭合。


总结

这个进程池模型并不复杂:

Start()   → 创建管道和子进程
RunOnce() → 轮询选择 Channel,发送任务码
Stop()    → 关闭所有写端,等待子进程退出

真正拉开“能跑”和“写对了”差距的,是 fd 管理。

父进程在循环里不断创建管道和子进程时,后创建的子进程会继承父进程已经打开的历史写端。只关闭当前管道的写端不够,历史写端也必须一起关掉,否则 EOF 条件永远无法成立。

这也是这篇文章最重要的结论:

排查重点:多进程程序出现退出卡死时,不要只盯着 waitpid(),先检查每个进程到底继承并持有哪些 fd。

进程池只是匿名管道的一次扩展应用,但它把管道最容易忽略的几个问题都暴露出来了:继承、引用计数、EOF 和退出顺序。把这些关系弄明白,后面再写更复杂的进程调度代码会稳很多。


资源分享:
【Linux】进程为什么不能直接通信?匿名管道原理、fd 继承与读写规则
【Linux】没有公网 IP 也能 SSH 回内网:反向 SSH、ProxyJump 与 systemd 保活实战
【Linux】系统彻底进不去怎么办?initramfs、recovery 与 chroot 救援指南

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/2402_87731470/article/details/163474266

文章来源crawl

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--