我们先来了解一下进程的通信

进程通信

为什么要有进程通信?

数据传输:⼀个进程需要将它的数据发送给另⼀个进程
资源共享:多个进程之间共享同样的资源。
通知事件:⼀个进程需要向另⼀个或⼀组进程发送消息,通知它(它们)发⽣了某种事件(如进
程终⽌时要通知⽗进程)。
进程控制:有些进程希望完全控制另⼀个进程的执⾏(如Debug进程),此时控制进程希望能够
拦截另⼀个进程的所有陷⼊和异常,并能够及时知道它的状态改变。
那么进程是怎么通信的?答案是用管道进行通信

什么是管道

管道是Unix中最古⽼的进程间通信的形式。
我们把从⼀个进程连接到另⼀个进程的⼀个数据流称为⼀个“管道”
我们可以来看看下面的图片
这里我们讲一个管道--匿名管道

匿名管道

匿名管道通常用于父子之间的进程通信

注意:匿名管道不是一个具体的文件他是存在于内核空间里的一段缓冲区

那么这个匿名管道是怎么做到父进程写入子进程读取的?

我们可以从fork层面来理解

在fork之前fd[0]读取端和fd[1]写入端都是指向父进程的

有了以上的认识,我们可以来实践一下

这里要介绍一个函数pipe;这个函数的作用就是如上面所说创建匿名管道的

他是包含在unstd.h头文件里面的

下面我们用匿名管道来实现一个功能子进程·1读取管道里面的内容,父进程进行写入

# include<iostream>
# include<unistd.h>
# include<sys/types.h>
# include<sys/wait.h>
# include<cstring>
void func()
{
    int fds[2];
    int ret = pipe(fds);
    if(ret == -1)
    {
        return ;
    }
    pid_t cid = fork();
    if(cid == 0)
    {
        char buffer[128] = {"\0"};
        close(fds[1]);
        ssize_t n = read(fds[0],buffer,sizeof(buffer));
        std::cout<<buffer<<std::endl;
        close(fds[0]);
        exit(0);
    }
    else
    {
        char buffer[] = {"i am father"};
        close(fds[0]);
        write(fds[1],buffer,strlen(buffer));
        close(fds[1]);
    }
}
int main()
{
    func();
    return 0;
}

运行结果

站在⽂件描述符⻆度-深度理解管道

我们不难发现我们对管道的读写操作,其实就是对文件的读写操作也调用的是read write这一类的文件操作接口,只不过这个文件是内存级别的,操作系统提供了这些内存级的文件,需要我们调用一些接口就可以打开这个文件,是读打开或者写打开,还有读写打开,打开文件会返回文件描述符而这个管道文件是以读写打开,把这两个文件描述符返回给上一层,这就完成了管道的创建,当子进程被创建时,他会继承父进程的代码还有父进程的文件描述符表,这时候子进程也指向父进程打开的文件,所以子进程能拿到父进程打开的管道文件

    这个管道文件是单向通信的,当子进程想要读取时就可以把写入的文件给关闭了,当父进程想要写入时就可以把读取的文件给关闭了,此时这就形成了一个单向通信的通道,这个就叫做管道,这个的工作模式是半双工的。

内核角度理解管道的本质

进程间通信的本质:让多个进程共享同一份内核内存资源,才有交换数据的基础。 匿名管道是纯内核内存结构,不落地磁盘。

管道其实也是文件,这也符合linux系统下一切皆文件的特点!

管道和普通文件都有 inode、struct file。

  1. 普通文件 inode、内容在磁盘;管道 inode、缓冲只在内核内存;
  2. 普通文件单 struct file 可双向随机读写;管道固定一对读 / 写 struct file,半双工流式通信;
  3. 普通文件持久化,管道随进程释放直接销毁。

利用匿名管道实现进程池

我们首先要知道一些特性

管道的容量: 默认是64kb
原子性: 原子性指的是一个操作要么完全执行,要么完全不执行,不会出现执行到一半的情况。

还有

管道通信的四种情况
写的慢,读的快 — 读端进程就需要阻塞(等写)
写的快,读的慢 — 写满缓冲区的时候,写就要阻塞等待(等读)
写端关,读端开 — read就会读到返回值为0,表示读到了文件结尾,读端也会自动关闭
写端开,读端关 — 写端再写入没有任何意义(OS不会做没有意义的事),OS会杀掉写端进程 ,发送异常信号13(SIGPIPE)

设计思路

我们先来看看我们需要实现一个什么样子的进程池

我们想要设计的进程池如上图所示,父进程创建几个管道和子进程,父进程通过写入管道告诉子进程要干什么,子进程通过读取管道来知道自己需要干什么

这就是设计思路

详细讲解

我们可以分成三个部分来设计,首先需要设计一个类为channel来描述一个管道,然后 一个类为channel_member来组织整个管道,然后一个process_pool的类来管理运行整个进程池,整体的思路总结下来就是先描述再组织

我们首先来看看单个channel

class channel
{
public:
    channel(int fd,int subid)
    :_fd(fd)
    ,_subid(subid)
    {
        _name="channel-"+std::to_string(fd)+" "+std::to_string(subid);
    }
    int ret_fd()
    {
        return _fd;
    }
    int ret_subid()
    {
        return _subid;
    }
    void print_name()
    {
        std::cout<<_name<<std::endl;
    }
    bool channel_close()
    {
        int ret = close(_fd);
        if(ret == 0)
        {
            sleep(1);
            std::cout<<"已关闭子进程id为-"<<_subid<<std::endl;
            return true;
        }
        std::cout<<"关闭失败-"<<_fd<<std::endl;
        return false;
    }
    void channel_wait()
    {
        pid_t n = waitpid(_subid,nullptr,0);
        (void)n;
        std::cout<<"父进程等待子进程成功id为: "<<_subid<<std::endl;
        sleep(1);
    }
    const std::string name()
    {
        return _name;
    }
private:
    int _fd;
    int _subid;
    std::string _name;
};

这里我们对单个管道进行描述,其中定义了名字,文件描述符,进程id等,还有对单个管道的关闭与等待

class channel_member
{   
public:
    channel_member()
    :_next(0)
    {}

    void insert(int fd,int subid)
    {
        channel c(fd,subid);
        _channels.push_back(c);
    }

    channel& select()
    {
        channel& ret = _channels[_next];
        _next++;
        _next%=_channels.size();
        return ret;
    }

    void print_channels()
    {
        for(auto& a: _channels)
        {
            a.print_name();
        }
    }

    void close_all_channel()
    {
        for(auto& a : _channels)
        {
            a.channel_close();
        }
    }
    void close_channel()
    {
         for(auto& a : _channels)
        {
            a.channel_close();
        }
    }
    void wait_channel()
    {
       for(auto& a : _channels)
        {
            a.channel_wait();
        }
    }
    void end()
    {
        for(auto& a : _channels)
        {
            a.channel_close();
            a.channel_wait();
        }
    }

private:
    int _next;
    std::vector<channel> _channels;
};

这里就是组织整个管道,对于整个管道的管理这里用了一个vector来进行组织管理这个select函数是用来进行轮询的,保证每个管道都有数据进行流通,每个子进程都会有任务执行,防止进程饥饿

接下来看看process_pool

class processpool
{
public:
    processpool(int num = 1)
    :_processnum(num)
    {}
    void execute( char* buffer)
    {
         task t1;
         int code = buffer[0] - '0';
         t1.send_code(code);
    }
    void work(int rfd)
    {
        char buffer[10] = {"\0"};
        int ret = read(rfd,buffer,sizeof(buffer));
        if(ret>0)
        {
            std::cout<<"子进程获得一个任务码"<<buffer<<std::endl;
            execute(buffer);
        }
        else if(ret == 0)
        {
            return ;
        }
        else
        {
            std::cout<<"读取错误!"<<std::endl;
            return ;
        }
    }
    void send_code(int wfd)
    {
        std::srand(std::time(nullptr));
        int num = std::rand() % 3;
        char buffer[10];
        sprintf(buffer, "%d", num); 
        write(wfd,buffer,strlen(buffer));
    }
    void get_member()
    {
        _cm.print_channels();
    }
    void close_all()
    {
        _cm.close_all_channel();
    }
    bool start()
    {
        
       for(int i  = 0;i<_processnum;i++)
       {
         int fds[2];
         int ret = pipe(fds);
        if(ret<0)
        {
            std::cout<<"管道创建错误"<<std::endl;
            return false;
        }
        pid_t rid = fork();
        if(rid == 0)
        {
            close_all();
            sleep(5);
            close(fds[1]);
            work(fds[0]);
            close(fds[0]);
            exit(0);
        }
        else if(rid < 0)
        {
            std::cout<<"子进程创建失败!"<<std::endl;
            return false;
        }
        else 
        {
            close(fds[0]);
            send_code(fds[1]);
            _cm.insert(fds[1],rid);
        }
       }
       return true;
    }
    void close_processpool()
    {
       _cm.end();
    }

private:
    channel_member _cm;
    int _processnum;
};

这里重点来说一下这里,为什么子进程在创建之后要调用close_all();

首先我们要知道,父进程创建子进程,子进程会继承父进程的代码还有对应的文件描述符

对于这种情况,我们需要进行处理,我们需要保证父进程的对于每个管道的写入端都必须是唯一的,只有父进程进行控制,这样才行。

那么怎么做?我们可以在子进程创建的第一时间去关闭std::vector<channel> _channels;里面的管道这里面只会记录父进程的写端。这样就能保证子进程不会再去指向父进程的读端,为什么可以直接操作这个数组,这样直接关闭,不是会把现在正在使用的管道进行关闭码?答案是不是的,这里要知道进程之间具有独立性!当我们要对内容进行修改时,会触发写时拷贝,子进程关闭的是父进程之前的内容!和现在的没关系,可以正常放心关闭。

下面是模拟发送任务的一个类,用来供给进程池使用。

typedef void(*task_t)();
void ping()
{
    std::cout << "我是一个网络测试的任务" << std::endl;
    
}

void Download()
{
    std::cout << "我是一个下载的任务" << std::endl;
}

void Upload()
{
    std::cout << "我是一个上传的任务" << std::endl;
}
class task
{
public:
    task()
    {
       Register (ping);
       Register (Download);
       Register (Upload);
    }
    void Register(task_t t)
    {
        task_arr.push_back(t);
    }

    void send_code(int code)
    {
        task_arr[code]();
    }
    
private:
    std::vector<task_t> task_arr;
};

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐