1. 前提

1. 如果一个线程一直在申请锁,执行访问临界资源的代码,和释放锁,此时会导致其它线程饥饿问题,因为先申请到锁的线程,再次获得锁的概率很大,因为其它线程要被唤醒,所以先申请到锁的线程比它们快,更容易再次申请到锁,为了解决这个问题引入线程同步

2. 线程同步

在保证数据安全的前提下,让线程能够按照某种特定的顺序访问临界资源,从而有效避免饥饿问题,叫做同步

2.1 条件变量

在这里插入图片描述
放苹果的人和拿苹果的人都不知道盘子里有无苹果,当拿苹果的人去拿苹果先申请锁,然后进去拿苹果,不管拿没拿到都要出来,但不能二次申请锁,要在队列里等待,其它人也一样,只有当放苹果的人申请锁在盘子里放入苹果后,摇一下铃铛,然后按照顺序在队列里唤醒一个人,唤醒后去申请锁,然后拿到苹果后继续排队,此时就同步了,这里的人就是线程,其中铃铛+队列就是条件变量

当⼀个线程互斥地访问某个变量时,它可能发现在其它线程改变状态之前,它什么也做不了。

例如⼀个线程访问队列时,发现队列为空,它只能等待,直到其它线程将⼀个节点添加到队列
中,这种情况就需要用到条件变量
在这里插入图片描述
这是条件变量的接口,跟锁类似
第一个参数是要初始化或销毁的条件变量,第二个参数是属性设置为nullptr
在这里插入图片描述
第二个参数是锁
在这里插入图片描述
1. 第一个函数是唤醒在该条件变量上等待的所有线程
2. 第二个函数是唤醒一个在该条件变量上等待的一个线程

2.2 封装条件变量

Cond.hpp

#pragma once
#include <iostream>
#include <pthread.h>
#include "Mutex.hpp"
using namespace MutexModule;
namespace CondModule
{
    class Cond
    {
    public:
        Cond()
        {
            pthread_cond_init(&_cond, nullptr);
        }
        void Wait(Mutex &mutex)
        {
            pthread_cond_wait(&_cond, mutex.Get());
        }
        void Signal()
        {
            pthread_cond_signal(&_cond);
        }
        void Broadcast()
        {
            pthread_cond_broadcast(&_cond);
        }
        ~Cond()
        {
            pthread_cond_destroy(&_cond);
        }

    private:
        pthread_cond_t _cond;
    };
}

3. 生产者消费者模型

生产者消费者模型本质就是多线程通信
在这里插入图片描述

生产者要申请锁然后才能对队列进行push(生产),消费者也需要申请锁对队列进行pop(消费),因为队列是共享资源,这里的生产者和消费者都是线程,所以生产和消费的过程是需要加锁的,也就是串行的,就是消费者和其它生产者必须等我生产完才能申请锁进行生产和消费

为什么要有生产者和消费者模型?
生产者和消费者模型是为了完成线程通信的

生产者和消费者模型的好处?
1. 生产过程和消费过程解耦

生产者生产完数据之后不用等待消费者处理,直接扔给阻塞队列,消费者不找生产者要数据,而是直接从阻塞队列里取,阻塞队列就相当于⼀个缓冲区,平衡了生产者和消费者的处理能力,这个阻塞队列就是用来给生产者和消费者解耦的

2. 支持忙闲不均
就是生产者一直在向队列里生产,但消费者就是不消费队列里的数据

3. 提高效率

生产者是要从另一个模块获取任务的,获取任务较慢,但是消费者拿到之前生产者push到队列里的任务是很快的,但是处理任务比拿任务慢多了

生产和消费的过程是需要加锁的是串行的,那为什么说生产者和消费者模型高效呢?
生产者和消费者模型的高效并不是体现在入队列和出队列上,而是体现在未来生产者从另一个模块并发获取任务的同时,消费者也在同时并发的处理具体的任务,高效在这,也就是让获取任务和处理任务并行起来

321原则

在这里插入图片描述

4. BlockingQueue

在这里插入图片描述

在多线程编程中阻塞队列(Blocking Queue)是⼀种常用于实现生产者和消费者模型的数据结构,其与普通的队列区别在于,当队列为空时,从队列获取元素的操作将会被阻塞,直到队列中被放入了元素,当队列满时,往队列里存放元素的操作也会被阻塞,直到有元素被从队列中取出
(以上的操作都是基于不同的线程来说的,线程在对阻塞队列进程操作时会被阻塞)

简单来说阻塞队列是一个容量具有上限的队列,不满足读写条件时,就要进行阻塞对应的线程

5. 实现多线程的生产者和消费者模型的阻塞队列

BlockQueue.hpp

// 阻塞队列的实现
#pragma once
#include <iostream>
using namespace std;
#include <string>
#include <queue>
#include <pthread.h>
const int defaultcap = 5;
template <class T>
class BlockQueue
{
private:
    bool IsFull() { return _q.size() >= _cap; }
    bool IsEmpty() { return _q.empty(); }

public:
    BlockQueue(int cap = defaultcap)
        : _csleep_num(0), _psleep_num(0), _cap(cap)
    {
        pthread_mutex_init(&_mutex, nullptr);
        pthread_cond_init(&_full_cond, nullptr);
        pthread_cond_init(&_empty_cond, nullptr);
    }
    void Equeue(const T &in) // 入队列
    {
        pthread_mutex_lock(&_mutex);
        while (IsFull())
        {
        
           // 应该让生产者线程进行等待
            // 重点1:pthread_cond_wait调用成功,挂起当前线程之前,要先自动释放锁!!
            // 重点2:当线程被唤醒的时候,默认就在临界区内唤醒!要从pthread_cond_wait
            // 成功返回,需要当前线程,重新申请_mutex锁!!!
            // 重点3:如果我被唤醒,但是申请锁失败了??我就会在锁上阻塞等待!!!
            _psleep_num++;
            std::cout << "生产者,进入休眠了: _psleep_num" <<  _psleep_num << std::endl;
            // 问题1: pthread_cond_wait是函数吗?有没有可能失败?pthread_cond_wait立即返回了
            // 问题2:pthread_cond_wait可能会因为,条件其实不满足,pthread_cond_wait 伪唤醒
            pthread_cond_wait(&_full_cond, &_mutex);
            _psleep_num--;
        }
        _q.push(in);
        if (_csleep_num > 0)
        {
            pthread_cond_signal(&_empty_cond);
            cout << "唤醒消费者..." << endl;
        }
        pthread_mutex_unlock(&_mutex);
    }
    T Pop()
    {
        pthread_mutex_lock(&_mutex);
        while (IsEmpty())
        {
        //这里和生产者那里问题一样
            _csleep_num++;
            cout << "消费者,进入休眠了: _csleep_num" << _csleep_num << endl;
            pthread_cond_wait(&_empty_cond, &_mutex);

            _csleep_num--;
        }
        T data = _q.front();
        _q.pop();
        if (_psleep_num > 0)
        {
        		//这里如果没有这个if,在这句话前面解锁可以,在后面解锁也可以
            pthread_cond_signal(&_full_cond);
            cout << "唤醒生产者..." << endl;
        }
        pthread_mutex_unlock(&_mutex);
        return data;
    }
    ~BlockQueue()
    {
        pthread_mutex_destroy(&_mutex);
        pthread_cond_destroy(&_full_cond);
        pthread_cond_destroy(&_empty_cond);
    }

private:
    queue<T> _q;                // 临界资源!!!
    int _cap;                   // 容量大小
    pthread_mutex_t _mutex;     // 锁
    pthread_cond_t _full_cond;  // 生产者的条件变量
    pthread_cond_t _empty_cond; // 消费者的条件变量
    int _csleep_num;            // 消费者休眠的个数
    int _psleep_num;            // 生产者休眠的个数
};

main.cc

#include "BlockQueue.hpp"
#include "Task.hpp"
#include <iostream>
#include <pthread.h>
#include <unistd.h>


void *consumer(void *args)
{
    BlockQueue<task_t> *bq = static_cast<BlockQueue<task_t> *>(args);

    while (true)
    {
        sleep(100);

        // 1. 消费任务
        task_t t = bq->Pop();

        // 2. 处理任务 -- 处理任务的时候,这个任务,已经被拿到线程的上下文中了,不属于队列了
        t();
    }
}

void *productor(void *args)
{
    BlockQueue<task_t> *bq = static_cast<BlockQueue<task_t> *>(args);
    while (true)
    {
        // 1. 获得任务
        std::cout << "生产了一个任务: " << std::endl;

        // 2. 生产任务
        bq->Equeue(Download);
    }
}

int main()
{
   
    // 申请阻塞队列
    BlockQueue<task_t> *bq = new BlockQueue<task_t>();

    // 构建生产和消费者
    pthread_t c[2], p[3];

    pthread_create(c, nullptr, consumer, bq);
    pthread_create(c + 1, nullptr, consumer, bq);
    pthread_create(p, nullptr, productor, bq);
    pthread_create(p + 1, nullptr, productor, bq);
    pthread_create(p + 2, nullptr, productor, bq);

    pthread_join(c[0], nullptr);
    pthread_join(c[1], nullptr);
    pthread_join(p[0], nullptr);
    pthread_join(p[1], nullptr);
    pthread_join(p[2], nullptr);

    return 0;
}

Makefile

CP:main.cc
	g++ -o $@ $^ -l pthread -std=c++11
.PHONY:clean
clean:
	rm -f CP

Task.hpp

#pragma once
#include <iostream>
#include <unistd.h>
#include <functional>
using task_t = std::function<void()>;
void Download()
{
    std::cout << "我是一个下载任务..." << std::endl;
    sleep(3); // 假设处理任务比较耗时
}

1. pthread_cond_wait调用成功,挂起当前线程之前,要先自动释放锁

2. 当线程被唤醒的时候,默认就在临界区(访问共享资源的代码)内唤醒,要从pthread_cond_wait 成功返回, 需要当前线程重新申请_mutex锁

3. 如果线程被唤醒后,线程申请锁失败了,线程就会在锁上阻塞等待

4. pthread_cond_wait是函数,有可能会失败,失败后pthread_cond_wait立即返回了,不会阻塞,如果是if向下执行push代码就会出现错误,因为队列大小满了不能继续push,所以循环是while,如果出现错误继续循环判断,再次阻塞该线程

5. pthread_cond_wait可能会因为条件其实不满足,或者如果把线程全部唤醒了空间还有一个,然后有一个线程push了,其它线程不能继续push,因为没空间了,这叫做pthread_cond_wait 伪唤醒

6. pthread_cond_signal(&_full_cond);
在这句话前面解锁可以,在后面解锁也可以,在这句话前面解锁唤醒对应条件变量下等待的线程,对应的线程一定会申请锁失败,因为我还没释放锁,但它会在锁上等待

在后面解锁唤醒对应条件变量下等待的线程也可以,因为对应的线程可能申请倒锁,也可能不能,不能就在锁上等待,能就执行临界区

这里的线程可以是生产者也可以是消费者,在代码的 while (IsFull()) 和 while(IsEmpty()) 这里会有这些问题和细节

6. POSIX信号量

跟之前的信号量的概念一样,只不过这是一个标准
本质是一个计数器,是对特定资源的预定机制

多线程使用资源有两种场景
1. 将目标资源整体使用,mutex+ -> 两元信号量
2. 将目标资源按照不同块,分别使用 -> 信号量

所有线程必须先看到sem,计数器sem–,sem++,信号量本质也是临界资源
p-- :原子
v++ : 原子

1. 理解信号量

在这里插入图片描述
如果n是1就是二元信号量

2. 信号量的接口

1. 初始化信号量

#include <semaphore.h>
int sem_init(sem_t *sem, int pshared, unsigned int value);

参数:
 pshared:0表⽰线程间共享,⾮零表⽰进程间共享
 value:信号量初始值

2. 销毁信号量

int sem_destroy(sem_t *sem);

3. 等待信号量

功能:等待信号量,会将信号量的值减1
int sem_wait(sem_t *sem); //P()

4. 发布信号量

功能:发布信号量,表⽰资源使⽤完毕,可以归还资源了。将信号量值加1int sem_post(sem_t *sem);//V()

3. 实现实现多线程的生产者和消费者模型的环形队列

RingQueue.hpp

#pragma once
#include <iostream>
#include <vector>
#include "Sem.hpp"
#include "Mutex.hpp"
using namespace SemModule;
using namespace MutexModule;
static const int gcap = 5;
template <class T>
class RingQueue
{

public:
    RingQueue(int cap = gcap)
        : _cap(cap), _p_step(0), _c_step(0), _blank(cap), _data(0), _rq(cap)
    {
    }
    void Equeue(const T &in)
    {
        _blank.P();
        {
            //这里是先申请信号量再申请锁,因为有一个线程如果申请锁成功了,访问临界资源的同时,其它线程可以同时在并发的申请信号量,高效
            LockGuard lockguard(_pmutex);
            _rq[_p_step++] = in;
            _p_step %= _cap;
        }
        _data.V();
    }
    void Pop(T *out)
    {
        _data.P();
        {
            LockGuard lockguard(_cmutex);
            *out = _rq[_c_step++];
            _c_step %= _cap;
        }
        _blank.V();
    }
    // ~RingQueue()
    // {
    // }
private:
    std::vector<T> _rq; // 组织环形队列
    int _cap;           // 环形队列容量
    int _p_step;        // 生产者下标
    int _c_step;        // 消费者下标
    Sem _blank;         // 空盘子
    Sem _data;          // 苹果
    Mutex _pmutex;      // 生产者和生产者之间的锁
    Mutex _cmutex;      // 消费者和消费者之间的锁
};

Makefile

sem:main.cc
	g++ -o $@ $^ -l pthread -g -std=c++11
.PHONY:clean
clean:
	rm -f sem

main.cc

#include <iostream>
#include <pthread.h>
#include <unistd.h>
#include "RingQueue.hpp"

struct threaddata
{
    RingQueue<int> *rq;
    std::string name;
};

void *consumer(void *args)
{
    threaddata *td = static_cast<threaddata*>(args);

    while (true)
    {
        sleep(3);
        // 1. 消费任务
        int t = 0;
        td->rq->Pop(&t);

        // 2. 处理任务 -- 处理任务的时候,这个任务,已经被拿到线程的上下文中了,不属于队列了
        std::cout << td->name << " 消费者拿到了一个数据:  " << t << std::endl;
        // t();
    }
}

int data = 1;

void *productor(void *args)
{
    threaddata *td = static_cast<threaddata*>(args);
    
    while (true)
    {
        //sleep(1);
        // sleep(2);
        // 1. 获得任务
        // std::cout << "生产了一个任务: " << x << "+" << y << "=?" << std::endl;
        std::cout << td->name << " 生产了一个任务: " << data << std::endl;

        // 2. 生产任务
        td->rq->Equeue(data);

        data++;
    }
}

int main()
{
    // 扩展认识: 阻塞队列: 可以放任务吗?
    // 申请阻塞队列
    RingQueue<int> *rq = new RingQueue<int>();

    // 构建生产和消费者
    // 如果我们改成多生产多消费呢??
    // 单单: cc, pp -> 互斥关系不需要维护,互斥与同步
    // 多多:cc, pp -> 之间的互斥关系!
    pthread_t c[2], p[3];

    threaddata *td = new threaddata();
    td->name = "cthread-1";
    td->rq = rq;
    pthread_create(c, nullptr, consumer, td);

    threaddata *td2 = new threaddata();
    td2->name = "cthread-2";
    td2->rq = rq;
    pthread_create(c + 1, nullptr, consumer, td2);

    threaddata *td3 = new threaddata();
    td3->name = "pthread-3";
    td3->rq = rq;
    pthread_create(p, nullptr, productor, td3);

    threaddata *td4 = new threaddata();
    td4->name = "pthread-4";
    td4->rq = rq;
    pthread_create(p + 1, nullptr, productor, td4);

    threaddata *td5 = new threaddata();
    td5->name = "pthread-5";
    td5->rq = rq;
    pthread_create(p + 2, nullptr, productor, td5);

    pthread_join(c[0], nullptr);
    pthread_join(c[1], nullptr);
    pthread_join(p[0], nullptr);
    pthread_join(p[1], nullptr);
    pthread_join(p[2], nullptr);

    return 0;
}

Mutex.hpp

#pragma once
#include <iostream>
#include <pthread.h>
namespace MutexModule
{
    class Mutex
    {
    public:
        Mutex()
        {
            pthread_mutex_init(&_mutex, nullptr);
        }
        void Lock()
        {
            pthread_mutex_lock(&_mutex);
        }
        void Unlock()
        {
            pthread_mutex_unlock(&_mutex);
        }
        ~Mutex()
        {
            pthread_mutex_destroy(&_mutex);
        }

    private:
        pthread_mutex_t _mutex;
    };

    class LockGuard
    {
    public:
        LockGuard(Mutex &mutex)
            : _mutex(mutex)
        {
            _mutex.Lock();
        }
        ~LockGuard()
        {
            _mutex.Unlock();
        }

    private:
        Mutex &_mutex;
    };
}

4. 封装信号量

Sem.hpp

#include <iostream>
#include <semaphore.h>
#include <pthread.h>
namespace SemModule
{
    const int defaultvalue = 1;
    class Sem
    {
    public:
        Sem(unsigned int sem_value = defaultvalue)
        {
            sem_init(&_sem, 0, sem_value);
        }
        void P()
        {
            sem_wait(&_sem);
        }
        void V()
        {
            sem_post(&_sem);
        }
        ~Sem()
        {
            sem_destroy(&_sem);
        }

    private:
        sem_t _sem;
    };
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐