一、引入

下面的代码模拟多人在同一时间段抢固定数量的票。

  • 多个线程模拟多人抢票
  • 每次抢票时使用usleep()函数休眠,模拟抢票花费的时间
  • 使用全局变量tickets储存当前票数,该变量被所有线程共享

tip:函数int usleep(useconds_t usec)——休眠usec微秒

#include <iostream>
#include <unistd.h>
#include <pthread.h>

using namespace std;

int tickets = 1000; // 票数

// 线程函数
void* getTickets(void* args)
{
	int threadId = (int)(long)args;
	while(1)
	{
		if(tickets > 0)
		{
			usleep(1000);
			printf("%d号线程:第%d张票\n", threadId, tickets--);
		}
		else break;
	}
	return nullptr;
}
int main()
{
	pthread_t t1, t2, t3;

	// 多线程抢票
	pthread_create(&t1, nullptr, getTickets, (void*)1);
	pthread_create(&t2, nullptr, getTickets, (void*)2);
	pthread_create(&t3, nullptr, getTickets, (void*)3);

    pthread_join(t1, nullptr);
    pthread_join(t2, nullptr);
    pthread_join(t3, nullptr);
    return 0;
}

运行结果:抢到了负数的票,这是不符合实际的

原因:

当一个线程等待时,一般就会切换到其他线程执行

  1. 当tickets=1时,if 语句判断条件为真,向下执行,遇到usleep()函数后,进行休眠,线程发生切换,此时tickets并未减1
  2. 此时内存中ticket仍然为1,后一个线程执行函数,if 语句判断条件仍然为真,向下执行,然后进行休眠,线程发生切换
  3. 当一个线程结束休眠,继续执行,打印字符串,将tickets--后写入内存
  4. 之后的线程结束休眠,继续执行时,就会使用上个线程减1后的tickets,并且也会将tickets--
  5. 所以,后面线程使用的tickets就比1小

从实验结果看,因为在判断和更新数据之间,切换了线程,导致出现问题

但其实 tickets-- 这个单一操作也是不安全的,只是实验不好展示

对一个全局变量进行多线程更改本身也不安全(ticket--本身也不安全)

ticket-- 操作本身不是⼀个原子操作,它实际上包括三个步骤:

  1. load:将共享变量num从内存加载到寄存器中
  2. update:更新寄存器里面num的值,执行-1操作
  3. store:将num新值从寄存器中写回内存地址

这些步骤至少对应三条汇编语句,所以会导致一些问题

问题举例:

假设有线程A、线程B、全局变量num(初始为100)。假设使两个线程都执行num--:

1.首先,线程A被调度执行num--:首先从内存中读取num的值——100;将num值减1;还没有来得及将num变化后的结果99写回内存,就被调度器给切换了,同时带走了寄存器中线程A的上下文(包括寄存器中num的值99)

2.此时,线程B被调度执行num--:首先从内存中读取num的值——100(此时内存中num的值还没有发生变化),线程B没有被打断,执行完了3个完整的步骤,所以内存中num的值被线程B更新为99,然后线程B又被调度器给切换了。

3.然后,线程A再次被调度回来继续执行未完成的num--:恢复寄存器中线程A的上下文,从中断的地方继续执行,将寄存器中num变化后的结果99写回内存,那么最后num的值就是99。

本来使两个线程都执行num减1,num的预期结果是98,而现在是99,就已经出现了问题

解决方案——加锁对全局变量(共享资源)进行保护

上面的过程展示了多个线程在交替执行时产生了数据不安全问题——数据不一致。所以,我们定义的全局变量在没有保护的时候,往往是不安全的。

    所以需要一把锁,保证在任一时刻,只有一个线程可以访问共享的资源或代码段,从而避免数据的不一致和错误。

    Linux上提供的这把锁叫互斥量。

    二、互斥量mutex

    1.线程间互斥相关的背景概念

    • 临界资源:多线程执行流进行安全访问的共享资源
    • 临界区:每个线程内部访问临界资源的代码(只占代码的很小一部分)
    • 互斥:任何时刻,互斥保证有且只有一个执行流进入临界区,访问临界资源,通常对临界资源起保护作用——多个线程串行访问临界资源
    • 原子性:不会被任何调度机制打断的操作,该操作只有两态,要么不做,要么做完(只用一条汇编完成的操作就是原子性的)

    2.互斥量的概念

    互斥锁的基本概念

    互斥锁(mutex)是一种用于实现多线程之间的同步机制的工具,它可以保证在任一时刻,只有一个线程可以访问共享的资源或代码段。互斥锁可以避免多线程程序中出现数据竞争(data race)或者死锁(deadlock)等问题,提高程序的正确性和稳定性。

    互斥锁的基本用法

    1. 创建一个互斥锁对象,然后在需要访问临界区域的代码前,调用互斥锁的lock()函数,以获取锁的所有权。
    2. 在访问完临界区域后,调用互斥锁的unlock()函数,以释放锁的所有权。

    当线程执行完任务释放锁以后,锁才会被传递给其他等待获取锁的线程,它们会重复以上操作以安全地完成任务。

    互斥锁的类型

    互斥锁可以分为全局锁和局部锁,它们的用法有所不同,但都要进行初始化、加锁和解锁操作

    全局锁是指在全局变量区定义的互斥锁,它可以被程序中的任何线程使用。

    • 全局锁的优点是简单易用,不需要传递参数,也不需要动态分配内存
    • 全局锁的缺点是可能造成资源浪费,因为不同的线程可能需要访问不同的共享资源,但是只能使用同一个互斥锁,这会导致不必要的等待和阻塞。另外,全局锁也不利于模块化编程,因为它破坏了数据的封装性。

    局部锁是指在局部变量区或堆区定义的互斥锁,它只能被定义它的函数或结构体中的线程使用。

    • 局部锁的优点是可以根据需要创建多个互斥锁,每个互斥锁只保护一个共享资源,这样可以提高并发性和效率。另外,局部锁也有利于模块化编程,因为它保持了数据的封装性。
    • 局部锁的缺点是 需要传递参数 或者 动态分配内存 ,这会增加编程的复杂度和开销。

    加锁和解锁

    加锁和解锁是一种实现临界区互斥性的方法,使线程串行执行临界区代码。

    • 加锁是指在进入临界区之前,线程获取一个锁对象。如果锁对象已经被其他线程占用,就必须等待或者阻塞,直到锁对象被释放。
    • 解锁是指在退出临界区之后,线程释放锁对象。从而让其他等待的线程有机会获取锁对象,并进入临界区。

    互斥锁的理解

    加锁保护全局资源时,所有线程对应区域都应该加上锁。谁持有锁,谁就进入临界区。

    当一个线程申请锁成功,访问临界资源时,其他申请锁的线程会阻塞等待。如果该线程被切走(比如时间片到了),其余线程仍然不能申请锁成功,不能访问临界资源,也就不能向后执行。

    因此,对于其余线程而言,有意义的锁状态只有 申请锁前(lock前) 和 释放锁后(unlock后)

    所以,在其余线程看来,持有锁的线程 执行临界区代码 是原子性的操作

    3.互斥量接口

    互斥锁的类型是pthread库中的一个结构体类型pthread_mutex_t,这个结构体包含了一些内部变量,用来表示互斥锁的状态和属性。我们暂时不需要关心这些变量的具体含义,只需要知道它是用来实现线程间的互斥操作的。

    (因为要使用pthread库,所以在编译时要加上 -lpthread 

    初始化互斥锁

    要使用pthread_mutex_t类型的变量,首先要对它进行初始化。

    初始化有两种方式:

    • 静态分配:在编译时就给互斥锁赋值为一个常量,表示它是一个默认属性的互斥锁(我们暂时不需要关心默认属性是什么),这种方式只能用于初始化静态或全局变量
    pthread_mutex_t mutex = PTHREAD_MUTEX_INITIALIZER // 它是一个宏
    

    好处是简单方便,不需要调用函数进行初始化和销毁缺点是只能使用默认属性,不能指定其他属性,例如是否递归、是否健壮等。

    • 动态分配:在运行时调用函数来初始化变量,这种方式可以用于初始化局部、静态和全局变量
    pthread_mutex_t mutex;
    pthread_mutex_init(&mutex, NULL);
    pthread_mutex_destroy(&mutex);
    

    好处是可以指定其他属性;缺点使用前必须初始化,使用后必须销毁

    pthread_mutex_init函数用于初始化互斥锁

     int pthread_mutex_init(pthread_mutex_t *mutex,const pthread_mutexattr_t *mutexattr);

    参数:

    • mutex:指向 pthread_mutex_t 类型变量的指针
    • mutexattr:指向 pthread_mutexattr_t 类型变量的指针,用于设置互斥锁的属性。如果使用默认属性,可以将第二个参数设置为 NULL。

    返回值:

    • 成功时返回 0,失败时返回错误码

    关于互斥锁的属性,暂时不用关心,一般设置为nullptr/NULL。

      销毁互斥锁

      在使用完互斥锁后,应调用函数来释放资源

      • 静态分配(使用 PTHREAD_ MUTEX_ INITIALIZER 初始化)不需要调用函数销毁互斥锁,动态分配才需要
      • 不要销毁一个已经加锁的互斥量
      • 已经销毁的互斥量,要确保后面不会有线程再尝试加锁

      pthread_mutex_destroy函数用于销毁互斥锁

      int pthread_mutex_destroy(pthread_mutex_t *mutex);

      参数:

      • mutex:指向 pthread_mutex_t 类型变量的指针

      返回值:

      • 成功时返回 0,失败时返回错误码

      加锁和解锁

      pthread_mutex_lock函数用于加锁互斥锁

      int pthread_mutex_lock(pthread_mutex_t *mutex);

      参数:

      • mutex:指向 pthread_mutex_t 类型变量的指针

      返回值:

      • 成功时返回 0,失败时返回错误码

      如果互斥量处于未锁状态,该函数会将互斥量锁定,同时返回成功

      如果互斥锁已经被锁定,调用该函数的线程将阻塞,直到互斥锁被解锁

      pthread_mutex_trylock函数用于尝试加锁互斥锁

      int pthread_mutex_trylock(pthread_mutex_t *mutex);

      参数:

      • mutex:指向 pthread_mutex_t 类型变量的指针

      返回值:

      • 成功时返回 0,失败时返回错误码

      如果互斥锁已经被锁定,该函数会立即返回而不会阻塞。

      在使用完共享资源后,应调用该函数来解锁互斥锁,以便其他线程可以访问共享资源。

      pthread_mutex_unlock函数用于解锁互斥锁

      int pthread_mutex_unlock(pthread_mutex_t *mutex);

      参数:

      • mutex:指向 pthread_mutex_t 类型变量的指针

      返回值:

      • 成功时返回 0,失败时返回错误码

      4.使用互斥锁

      改进上面的抢票系统:

      使用全局锁

      #include <iostream>
      #include <unistd.h>
      #include <pthread.h>
      
      using namespace std;
      
      pthread_mutex_t mutex =PTHREAD_MUTEX_INITIALIZER; // 定义一个全局锁
      
      int tickets = 1000; // 票数
      
      // 线程函数
      void* getTickets(void* args)
      {
      	int threadId = (int)(long)args;
      	while(1)
      	{
      		pthread_mutex_lock(&mutex); // 加锁
      		if(tickets > 0)
      		{
      			usleep(1000);
      			printf("%d号线程:第%d张票\n", threadId, tickets--);
      			pthread_mutex_unlock(&mutex); // 解锁
      		}
      		else
      		{
      			pthread_mutex_unlock(&mutex); // 解锁
      			break;
      		}
      	}
      	return nullptr;
      }
      int main()
      {
      	pthread_t t1, t2, t3;
      
      	// 多线程抢票
      	pthread_create(&t1, nullptr, getTickets, (void*)1);
      	pthread_create(&t2, nullptr, getTickets, (void*)2);
      	pthread_create(&t3, nullptr, getTickets, (void*)3);
      
          pthread_join(t1, nullptr);
          pthread_join(t2, nullptr);
          pthread_join(t3, nullptr);
          return 0;
      }
      

      运行结果:最后全局变量tickets的值不会是0或-1

      使用局部锁

      由于锁是局部变量,存放在主线程的栈上,需要传递给其余线程使用

      所以,创建了结构体ThreadData,向新线程传递 主线程中创建的局部锁 和 线程索引

      #include <iostream>
      #include <unistd.h>
      #include <pthread.h>
      
      using namespace std;
      
      // pthread_mutex_t mutex =PTHREAD_MUTEX_INITIALIZER; // 定义一个全局锁
      
      class ThreadData // 为了线程
      {
      public:
      	// 构造函数
      	ThreadData(int index, pthread_mutex_t *pmtx)
      		: _index(index), _pmtx(pmtx)
      	{
      	}
      
      public:
      	int _index;				// 线程索引
      	pthread_mutex_t *_pmtx; // 锁的地址
      };
      
      int tickets = 1000; // 票数
      
      // 线程函数
      void *getTickets(void *args)
      {
      	ThreadData *pData = (ThreadData *)args;
      	int threadId = pData->_index;
      
      	while (1)
      	{
      		pthread_mutex_lock(pData->_pmtx); // 加锁
      		if (tickets > 0)
      		{
      			usleep(1000);
      			printf("%d号线程:第%d张票\n", threadId, tickets--);
      			pthread_mutex_unlock(pData->_pmtx); // 解锁
      		}
      		else
      		{
      			pthread_mutex_unlock(pData->_pmtx); // 解锁
      			break;
      		}
      	}
      	return nullptr;
      }
      int main()
      {
      	pthread_mutex_t mtx;			   // 定义局部锁
      	pthread_mutex_init(&mtx, nullptr); // 初始化局部锁
      
      	ThreadData td[3] = {
      		ThreadData(1, &mtx),
      		ThreadData(2, &mtx),
      		ThreadData(3, &mtx)};
      
      	pthread_t t1, t2, t3;
      	// 多线程抢票
      	pthread_create(&t1, nullptr, getTickets, (void *)&td[0]);
      	pthread_create(&t2, nullptr, getTickets, (void *)&td[1]);
      	pthread_create(&t3, nullptr, getTickets, (void *)&td[2]);
      
      	pthread_join(t1, nullptr);
      	pthread_join(t2, nullptr);
      	pthread_join(t3, nullptr);
      
      	pthread_mutex_destroy(&mtx); // 销毁锁
      
      	return 0;
      }
      

      运行结果:

      运行现象

      • 对比加锁之前,加锁之后运行速度明显变慢。这是因为同一时间只允许一个线程访问临界资源,会带来性能损耗
      • 出现同一个线程连续抢票的情况。这是因为锁只是规定互斥访问,执行流竞争到锁就能够执行临界区代码(进行抢票),线程执行也就没有顺序

      性能损耗

      在加锁之前,同一时间可以有多个线程对临界资源进行访问。但是加锁之后,同一时间只允许一个线程访问临界资源(串行执行),会带来一定程度上的性能损耗。

      互斥锁虽然能保护共享资源的安全,但同时也会带来一些性能上的开销,主要有以下几个方面:

      • 互斥锁的创建和销毁需要调用操作系统的API,这会消耗一定的时间和内存资源。
      • 互斥锁的加锁和解锁需要进行原子操作(atomic operation),这会增加CPU的指令数和内存访问次数。
      • 互斥锁的等待和唤醒需要进行上下文切换(context switch),这会导致CPU缓存(cache)的失效和线程调度(scheduling)的延迟。
      • 互斥锁的竞争会造成线程的阻塞(blocking)或者忙等待(busy waiting),这会降低线程的利用率和并发度。

      因此,互斥锁在一定程度上会降低多线程程序的效率,尤其是在互斥锁保护的代码段或资源:

      • 非常频繁地被访问,导致锁的竞争很激烈。
      • 非常耗时地被执行,导致锁的持有时间很长。
      • 非常简单地被处理,导致锁的开销占比很高。

      如何减少互斥锁对效率的影响呢?

      一般来说,有以下几个建议:

      • 尽量减少互斥锁的数量和范围,只保护必要的共享数据或临界区(critical section),避免过度同步(oversynchronization)。
      • 尽量缩短互斥锁的持有时间,尽快释放锁,避免在持有锁的情况下进行I/O操作或其他耗时操作。
      • 尽量使用更高效的同步机制,如读写锁(read-write lock)、自旋锁(spin lock)、条件变量(condition variable)等,根据不同场景选择合适的工具。

      总之,互斥锁是一种有利有弊的同步机制,它可以保证多线程程序的正确性和稳定性,但也会降低程序的效率。因此,在使用互斥锁时,需要权衡利弊,合理设计和优化代码,以达到最佳的性能表现。

      5.互斥量实现原理

      要访问临界资源,每一个线程都要申请锁,也就是说这个锁必须被所有线程共享,因此锁本身也是一种共享资源锁保证了临界资源的安全,那么谁来保证锁本身的安全?

      锁本身的安全是由操作系统的原子操作保证的。原子操作可以保证在多线程环境下在任何时刻只有一个线程能够访问锁,申请锁(lock)和释放锁(unlock)的操作是原子的,这样就保证了锁本身的安全。

      lock和unlock的实现原理:

      伪代码:

      lock具体实现

      1. 线程向寄存器al中放入0
      2. 交换寄存器al中的值和内存中mutex的值
      3. 判断寄存器al中的值:若为1,则申请锁成功,lock函数返回;若为0,则申请失败,线程挂起等待,等待另一线程解锁

      内存中全局变量mutex的初始值为1,代表锁。

      • 当mutex为1时,代表锁没有被申请
      • 当mutex为0时,代表锁已经被申请了

      当线程的上下文中,寄存器al的值为1,则代表持有锁。

      • 一个线程申请锁成功后,该线程寄存器al中的值为1,内存中mutex变为0
      • 若此时其余线程申请锁,交换寄存器al中的值和内存中mutex的值后,寄存器al中的值仍为0,申请锁失败

      为什么申请锁(lock)的操作是原子的?

      如果线程在步骤1执行完过后被切换,寄存器中的0会作为线程的上下文被带走,对后续步骤没有影响

      步骤2(swap或exchange指令)作用是把寄存器和内存单元的数据相交换,由于只有一条汇编指令,保证了原子性。

      unlock具体实现

      通过swap或exchange指令作用是把寄存器和内存单元的数据相交换,寄存器al中的值变为0,内存中mutex重新变为1,同时唤醒申请锁时被阻塞的进程。这就相当于将锁归还回原处。

      由于只有一条汇编指令,保证了原子性。

      6.死锁

      死锁的概念

      死锁是指两个或多个线程在执行过程中,由于竞争资源而造成的一种互相等待的现象,若无外力作用,它们都将无法继续执行。死锁通常发生在多个线程同时请求多个资源时,由于资源分配的不当,导致线程之间相互等待,无法继续执行。

      例如:线程A和线程B各自拥有锁a和锁b,但是它们有了锁还要申请对方的锁,因为它们申请的锁已经被占用,最后会导致代码无法推进。

      一把锁也可能造成死锁

      比如当一个线程连续申请同一把锁两次时,也会造成死锁,使线程阻塞

      pthread_mutex_lock(&mutex); 
      pthread_mutex_lock(&mutex); 

      死锁的必要条件

      • 互斥条件:一个资源每次只能被一个执行流使用
      • 请求与保持条件:一个执行流因请求资源而阻塞时,对已获得的资源保持不放
      • 不剥夺条件:一个执行流已获得的资源,在未使用完之前,不能强行被剥夺

      这里我们理解该条件可以结合上锁的函数:

      在这里插入图片描述
      其中lock就是阻塞式申请锁,申请不到就去阻塞等待
      trylock时如果锁已经被占用,就会去释放对方的锁再让自己去申请锁(这样就破坏了不剥夺条件,避免死锁)

      • 循环等待条件:若干执行流之间形成一种头尾相接的循环等待资源的关系

      避免死锁的方法

      • 破坏这四个必要条件中的一个或多个:
      1. 使用信号量或互斥锁来实现对资源的互斥访问,避免多个进程或线程同时竞争同一个资源。
      2. 使用银行家算法或者预分配算法来分配资源,避免进程或线程在占有资源的同时请求新的资源,导致资源不足。
      3. 使用优先级机制或者超时机制来实现对资源的抢占,为不同类型的锁分配不同的优先级,按照优先级顺序获取锁。避免低优先级的进程或线程长时间占用资源,阻塞高优先级的进程或线程。
      4. 使用拓扑排序或者有序分配法来分配资源,避免进程或线程之间形成循环等待的链条,资源一次性分配。或者访问完临界资源以后,就马上干净地释放锁。
      • 设置锁超时:为每个锁设置一个超时时间,如果在超时时间内无法获取锁,则放弃获取并释放已经获取的锁,避免锁未释放的情况。
      • 使用死锁检测算法:定期运行死锁检测算法,检测系统中是否存在死锁。如果检测到死锁,则采取相应措施进行解除。
         

      三、可重入和线程安全

      1.概念

      • 线程安全:就是多个线程在访问共享资源时,能够正确地执行,不会相互干扰或破坏彼此的执行结果。⼀般而言,多个线程并发同⼀段只有局部变量的代码时,不会出现不同的结果。但是对全局变量或者静态变量进行操作,并且没有锁保护的情况下,容易出现该问题。
      • 重入:同一个函数被不同的执行流调用,当前一个流程还没有执行完,就有其他的执行流再次进入,我们称之为重入。⼀个函数在重入的情况下,运行结果不会出现任何不同或者任何问题,则该函数被称为可重入函数,否则,是不可重入函数

      这两个概念本身没有太大关系,只是线程安全问题很多是由不可重入函数导致的

      2.常见的线程不安全的情况

      • 对全局变量或静态变量进行操作:如果多个线程并发访问同一个全局变量或静态变量,并且对它进行了修改操作,那么可能会出现线程安全问题。
      • 使用非线程安全的函数:一些函数(如strtok和gmtime)在多线程环境中使用时可能会出现线程安全问题。这些函数通常都有线程安全的替代版本(如strtok_r和gmtime_r),应该尽量使用这些替代版本。
      • 没有正确使用锁:如果多个线程需要并发访问同一块数据,那么应该使用锁来保护这块数据。如果没有正确使用锁,或者锁的粒度不够细,那么可能会出现线程安全问题。

      3.常见线程安全的情况

      • 每个线程对全局变量或者静态变量只有读取的权限,而没有写入的权限,一般来说这些线程是安全的
      • 类或者接口对于线程来说都是原子操作
      • 多个线程之间的切换不会导致该接口的执行结果存在二义性

      4.常见不可重入的情况

      • 调用了malloc/free函数,因为malloc函数是用全局链表来管理堆的
      • 调用了标准I/O库函数,标准I/O库的很多实现都以不可重入的方式使用全局数据结构
      • 可重入函数体内使用了静态的数据结构

      5.常见可重入情况

      • 不使用全局变量或静态变量
      • 不使用用malloc或者new开辟出的空间
      • 不调用不可重入函数
      • 不返回静态或全局数据,所有数据都有函数的调用者提供
      • 使用本地数据,或者通过制作全局数据的本地拷贝来保护全局数据

      6.可重入与线程安全联系

      • 函数是可重入的,那就是线程安全的
      • 函数是不可重入的,那就不能由多个线程使用,有可能引发线程安全问题
      • 如果一个函数中有全局变量,那么这个函数既不是线程安全也不是可重入的

      7.可重入与线程安全区别

      • 可重入函数是线程安全函数的一种
      • 线程安全不一定是可重入的,而可重入函数则一定是线程安全的
      • 假设有一个不可重入的函数,如果不加锁,它就不是线程安全的;如果加上锁,它是线程安全的,但仍是不可重入函数
      • 如果将对临界资源的访问加上锁,则这个函数是线程安全的,但如果这个重入函数,若锁还未释放,会产生死锁,因此是不可重入的

      四、线程同步

      1.引入

      举抢票的例子来说

      我们可以看到,2号线程连续抢了多张票,这是因为2号线程的优先级比较高。这样访问临界资源是不会出错的,但是它是不合理的,会造成其他线程的饥饿问题

      加锁能够保证在同一时间只有一个线程执行临界区代码访问临界资源,但它不能保证让每一个线程都能访问临界资源。

      所以,我们才需要一个同步机制,让多个线程按照一定顺序访问临界资源。

      2.相关概念

      线程同步

      同步和异步通常用来描述两个或多个事件之间的关系:

      • 同步是指两个或多个事件按照一定的顺序发生,一个事件的发生依赖于另一个事件的完成
      • 异步则是指两个或多个事件之间没有固定的先后顺序,它们可以独立发生

      对于线程而言,线程同步指的是协调多个线程按照某种特定的顺序执行,以确保它们能够正确地访问共享资源。这通常需要使用一些同步机制,如互斥锁、信号量和条件变量等,来控制线程之间的执行顺序。

      竞态条件

      竞态条件是指在多线程程序中,多个线程同时访问和修改共享资源,可能会导致程序出现不确定的行为,甚至产生错误的结果

      3.生产者消费者模型

      学习生产者消费者模型是为了理解线程同步更广泛的应用场景

      引入

      超市、厂商和顾客是一个很好的例子

      • 厂商可以被看作是生产者,它生产商品并将其运送到超市。
      • 超市可以被看作是缓冲区,它存储厂商生产的商品。
      • 顾客可以被看作是消费者,它从超市购买商品。

      当超市的库存充足时,厂商不需要再运送更多的商品;当超市的库存不足时,厂商需要生产更多的商品并将其运送到超市。

      同样,当超市有足够的商品时,顾客可以购买它们;当超市缺货时,顾客需要等待厂商运送更多的商品。

      生产者消费者模式的基本思想:生产者负责生成数据并将其放入缓冲区,而消费者则从缓冲区中取出数据并进行处理。

      • 当缓冲区为空时,消费者需要等待生产者生成新的数据
      • 当缓冲区已满时,生产者需要等待消费者取出数据

      缓冲区——BlockingQueue

      生产者消费者模型的缓冲区可以由堵塞队列(BlockingQueue)实现的。

      Blocking Queue常用于生产者消费者模型中,作为生产者和消费者之间的缓冲区,生产者向Blocking Queue中添加数据,消费者从Blocking Queue中取出数据,遵循队列的先进先出(FIFO)的原则:

      • 当Blocking Queue已满时,生产者线程将会被阻塞;
      • 当Blocking Queue为空时,消费者线程将会被阻塞。
         

      生产者消费者模式

      • 1个交易场所:一段缓冲区,生产者向缓冲区生产数据,而消费者从缓冲区取走数据。
      • 2种角色:生产者线程和消费者线程。
      • 3种关系:生产者之间:互斥  ;消费者之间:互斥  ;生产者和消费者:互斥与同步

      互斥:同一时间只有一个线程访问临界资源

      同步:多个生产者线程/消费者线程按照一定顺序进行访问临界资源

      生产者和消费者彼此之间不直接通讯,而通过阻塞队列来进行通讯,所以生产者生产完数据之后不用等待消费者处理,直接扔给阻塞队列,消费者不找生产者要数据,而是直接从阻塞队列里取,阻塞队列就相当于一个缓冲区,平衡了生产者和消费者的处理能力。这个阻塞队列就是用来给生产者和消费者解耦的。

      好处:

      • 生产者线程和消费者线程解耦
      • ⽀持并发
      • ⽀持生产和消费一段时间的忙闲不均
      • 提高效率

      如果没有缓冲区,生产者和消费者是强耦合

      在没有缓冲区的情况下,生产者必须直接将数据传递给消费者,而消费者也必须直接从生产者那里获取数据。这样一来,生产者必须等待消费者准备好接收数据,而消费者也必须等待生产者生成新的数据,是强耦合关系

      相反,如果有了缓冲区,那么生产者和消费者之间就可以通过缓冲区来解耦。生产者只需要将数据放入缓冲区,而不需要关心消费者何时获取数据。同样,消费者也只需要从缓冲区中取出数据,而不需要关心生产者何时生成新的数据。这样一来,生产者和消费者之间就可以独立地运行,它们之间的耦合度也会降低,使程序运行更加高效。 

      4.应用条件变量的思路

      互斥锁和条件变量相互配合,实现线程同步。

      一个条件变量会维护一个等待队列,让多个线程 依次 申请互斥锁访并问临界资源

      线程同步可以应用在抢票的例子中

      问题出现:

      优先级较高的线程可能会不停抢票,造成其他线程的饥饿问题

      问题解决:

      当一个线程获取锁后,不直接去抢票,先将该线程放入条件变量等待(同时自动释放锁),直到被条件变量唤醒。条件变量唤醒该线程后,该线程自动再去获取锁,进行抢票,释放锁。

      条件变量内部是以队列的形式管理等待的线程,后加入的线程排在队尾,队头的线程最先被唤醒

       理解图(只用于理解,不是内部的实际结构)

      这样一来,就实现了多个线程依次访问临界资源

      在生产者消费者模式中,线程同步有更广泛的应用

      问题出现:

      生产者线程获取锁后,判断缓冲区是否已满,若当前缓冲区已满,生产者线程不能向缓冲区写入数据,只能释放锁。

      但是如果生产者线程的优先级较高,就会不停 加锁-判断-解锁 ,一直访问临界资源,执行无意义的操作,从而导致消费者线程无法访问临界资源,降低效率

      当缓冲区为空时,消费者线程也可能有类似的问题

      问题解决:

      使用条件变量和互斥锁配合。

      当生产者线程获取锁后,若缓冲区已满,判断条件不满足,就将该线程放入条件变量等待(同时自动释放锁),直到被条件变量唤醒。当消费者线程 获取锁-从缓冲区取出数据-释放锁 后,通知条件变量,生产者线程被唤醒,获取锁-向缓冲区加入数据-释放锁。

      当缓冲区为空时,消费者线程也是类似的解决方式

      5.条件变量的接口

      条件变量的类型是pthread库中的一个结构体类型pthread_cond_t

      pthread_cond族函数是Linux下的一组用于线程同步的函数。它们包括:

      • pthread_cond_init:初始化条件变量。
      • pthread_cond_wait:阻塞等待条件变量满足。
      • pthread_cond_signal:唤醒一个等待条件变量的线程。
      • pthread_cond_broadcast:唤醒所有等待条件变量的线程。
      • pthread_cond_timedwait:阻塞等待条件变量满足,直到指定时间
      • pthread_cond_destroy:销毁条件变量。

      这些函数的返回值都是一样的:当函数执行成功时,它们都返回0。任何其他返回值都表示错误。

      pthread_cond_init

      初始化条件变量

      int pthread_cond_init(pthread_cond_t *restrict cond, const pthread_condattr_t *restrict attr);

      参数:

      • cond:需要初始化的条件变量。
      • attr:初始化条件变量的属性,一般设置为 NULL/nullptr 表示默认属性

      和定义互斥锁类似,调用 pthread_cond_init 函数初始化条件变量叫做动态分配,除此之外,还可以静态分配(一般在全局使用):

      pthread_cond_t cond = PTHREAD_COND_INITIALIZER; // 它是一个宏

      注意:静态分配的条件变量不需要手动调用函数销毁。

      pthread_cond_destroy

      销毁条件变量

      int pthread_cond_destroy(pthread_cond_t *cond);

      参数:

      • cond:需要销毁的条件变量。

      pthread_cond_wait

      释放互斥锁,并将线程挂起

      int pthread_cond_wait(pthread_cond_t *restrict cond, pthread_mutex_t *restrict mutex);

      参数:

      • cond:需要等待的条件变量。
      • mutex:当前线程所处临界区对应的互斥锁。

      为什么需要mutex(互斥锁参数)?

      当前线程是持有锁的,当他陷入休眠后,如果不释放锁,其他线程无法申请锁,也无法访问临界资源,就会出现问题。所以需要传入mutex参数,将锁自动释放。当线程被唤醒时,因为要访问临界资源,也会自动重新申请锁,从当前位置继续向下执行。

      • 该函数调用时,会以原子性的方式将锁释放,将线程自身挂起
      • 当线程被唤醒时,自动重新获取互斥锁:若获取锁成功,该函数返回;若获取锁失败,处于竞争锁状态,直到获取锁成功

      pthread_cond_broadcast 和 pthread_cond_signal

      唤醒等待条件变量的所有/首个线程

      int pthread_cond_broadcast(pthread_cond_t cond);
      int pthread_cond_signal(pthread_cond_t cond);

      参数:

      • cond:唤醒在 cond 条件变量下等待的线程。

      区别:

      • pthread_cond_signal 函数用于唤醒等待队列中首个线程。
      • pthread_cond_broadcast 函数用于唤醒等待队列中的全部线程。

      6.应用条件变量(线程同步)

      在抢票系统中的应用

      说明:

      1. 抢票线程在获取锁后,加入条件变量管理的队列等待,同时自动释放锁
      2. 主线程每隔1秒发送条件变量信号,唤醒一个抢票线程
      3. 抢票线程重新获取锁,执行抢票,然后释放锁

      这样一来,多个线程就能按照一定顺序访问临界资源,执行抢票

      #include <iostream>
      #include <unistd.h>
      #include <pthread.h>
      
      using namespace std;
      
      pthread_mutex_t mutex = PTHREAD_MUTEX_INITIALIZER; // 定义一个全局锁
      pthread_cond_t cond = PTHREAD_COND_INITIALIZER;    // 定义一个全局条件变量
      
      int tickets = 1000; // 票数
      
      // 线程函数
      void *getTickets(void *args)
      {
          while (true)
          {
              pthread_mutex_lock(&mutex);       // 加锁
              pthread_cond_wait(&cond, &mutex); // 等待条件变量
              // 执行抢票操作(这里先省略判断票数)
              printf("%s:第%d张票\n", (char *)args, tickets--);
              pthread_mutex_unlock(&mutex); // 解锁
          }
          return nullptr;
      }
      
      int main()
      {
          // 多线程抢票
          pthread_t t1, t2, t3;
      
          pthread_create(&t1, nullptr, getTickets, (void *)"thread 1");
          pthread_create(&t2, nullptr, getTickets, (void *)"thread 2");
          pthread_create(&t3, nullptr, getTickets, (void *)"thread 3");
      
          while (true)
          {
              pthread_cond_signal(&cond); // 发出条件变量信号,唤醒一个等待线程
              sleep(1);                   // 模拟其他操作,避免过于频繁唤醒
          }
      
          pthread_join(t1, nullptr);
          pthread_join(t2, nullptr);
      
          return 0;
      }
      

      执行结果:

      在生产者消费者模型中的应用

      实现基于BlockingQueue的生产者消费者模型

      生产者线程向阻塞队列放入int类型的数据,消费者线程通过阻塞队列获取int类型的数据

      BlockQueue类的框架(BlockQueue.hpp)

      我们使用互斥锁条件变量保证线程安全。

      生产数据(push)消费数据(pop)的过程:

      • 生产数据:获取锁->判断缓冲区是否满了->若满了就阻塞等待条件变量->被唤醒->将数据加入缓冲区->释放锁
      • 消费数据:获取锁->判断缓冲区是否为空->若为空就阻塞等待条件变量->被唤醒->将数据从缓冲区中取出->释放锁

      向条件变量发信号,唤醒线程的过程:

      • 生产者唤醒消费者:生产者将数据加入缓冲区->缓冲区不为空(货物不为空)->生产者向消费者条件变量发信号->唤醒消费者
      • (当生产者生产完,把商品放到超市中,就能唤醒消费者继续取出)
      • 消费者唤醒生产者:消费者将数据从缓冲区中取出->缓冲区没有满(货物没有满)->消费者向生产者条件变量发信号->唤醒生产者
      • (当消费者取出商品后,消费者就能唤醒生产者继续生产)

      这样就实现了线程同步,保证同一时间只有一个线程能够访问容器,避免两者同时在缓冲区中操作数据,从而避免竞态条件。

      说明:

      • 一个计数器maxSize记录着阻塞队列的最大容量
      • 为了保存各种数据类型,使用了模板
      #pragma once
      
      #include <queue>
      #include <pthread.h>
      #include <iostream>
      using namespace std;
      
      template <class T>
      class BlockQueue
      {
      public:
        BlockQueue(int maxSize)
        {
          _maxSize = maxSize;
          pthread_mutex_init(&_mutex, nullptr);
          pthread_cond_init(&_condConsumer, nullptr);
          pthread_cond_init(&_condProducer, nullptr);
        }
        ~BlockQueue()
        {
          pthread_mutex_destroy(&_mutex);
          pthread_cond_destroy(&_condConsumer);
          pthread_cond_destroy(&_condProducer);
        }
      
        void push(const T &in) // 输入型参数
        {
          pthread_mutex_lock(&_mutex); // 加锁
          while (isFull())             // 细节1:使用while而不是if
          {
            pthread_cond_wait(&_condProducer, &_mutex); // 生产条件不满足,阻塞生产者
          }
          _queue.push(in);                     // 生产数据
          pthread_cond_signal(&_condConsumer); // 通知消费者可以消费了
          pthread_mutex_unlock(&_mutex);       // 解锁
          //pthread_cond_signal(&_condConsumer);//细节2:可以放在解锁之前或之后
        }
      
        void pop(T *out) // 输出型参数
        {
          pthread_mutex_lock(&_mutex); // 加锁
          while (isEmpty())            // 细节1:使用while而不是if
          {
            pthread_cond_wait(&_condConsumer, &_mutex); // 消费条件不满足,阻塞消费者
          }
          *out = _queue.front(); // 消费数据
          _queue.pop();
          pthread_cond_signal(&_condProducer); // 通知生产者可以生产了
          pthread_mutex_unlock(&_mutex);       // 解锁
          //pthread_cond_signal(&_condProducer);//细节2:可以放在解锁之前或之后
        }
      
      private:
        bool isFull()
        {
          return _queue.size() >= _maxSize;
        }
        bool isEmpty()
        {
          return _queue.empty();
        }
      
      private:
        queue<T> _queue;
        int _maxSize; // 队列的最大容量
      
        pthread_mutex_t _mutex;
        pthread_cond_t _condConsumer; // 消费者条件变量
        pthread_cond_t _condProducer; // 生产者条件变量
      };
      • 细节1:在push和pop进行条件判断时,使用while而不是if

      这样做的目的是:防止虚假唤醒。举个例子,如果使用的是pthread_cond_broadcast(),会同时唤醒所有休眠的线程,开始竞争互斥锁,获取锁后直接向下执行。可能会造成缓冲区为空,仍在取出数据/缓冲区已满,仍在放入数据

      • 细节2:唤醒其他线程(pthread_cond_signal),可以放在解锁之前或之后,没有影响
      生产消费过程(main.cc)

      生产和消费int类型的数据

      #include "BlockQueue.hpp"
      #include <unistd.h>
      
      void *produce(void *arg) // 生产者线程函数
      {
          BlockQueue<int> *bq = static_cast<BlockQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              cout << "生产者生产了一个数据: " << i << endl;
              bq->push(i);
              sleep(1); // 模拟生产时间,使生产者比消费者慢一点
          }
          return nullptr;
      } 
      
      void *consume(void *arg) // 消费者线程函数
      {
          BlockQueue<int> *bq = static_cast<BlockQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              int data;
              bq->pop(&data);
              cout << "消费者消费了一个数据: " << data << endl;
              //sleep(1); // 模拟消费时间,使消费者比生产者慢一点
             
          }
          return nullptr;
      } 
      
      int main()
      {
          BlockQueue<int> *bq = new BlockQueue<int>(5);
          pthread_t consumer, producer;
          pthread_create(&producer, nullptr, produce, bq); // 创建生产者线程
          pthread_create(&consumer, nullptr, consume, bq); // 创建消费者线程
      
          pthread_join(producer, nullptr);
          pthread_join(consumer, nullptr);
      
          delete bq;
          return 0;
      }
      运行结果
      • 当生产者比消费者慢时

      消费一个数据,阻塞队列为空,等待数据生产,再进行消费:生产一个,消费一个

      • 当消费者比生产者慢时

      最开始生产者会将阻塞队列填满,等待数据消费,之后等待数据消费,再进行生产:消费一个,生产一个

      消费顺序和生产顺序也是一一对应的。

      在生产者消费者模型中的扩展应用

      刚刚的例子中,我们生产和消费的是int类型的数据,这只是用于测试,是没有实际意义的。现在我们向阻塞队列放入的是需要处理的任务,消费者线程会从阻塞队列中取出任务,并完成任务。

      BlockQueue类的框架(BlockQueue.hpp)

      与上面相同

      任务类(Task.hpp)

      任务类Task中包含了 处理任务的方式(回调函数),任务的数据(传入函数的参数)

      #pragma once
      #include <iostream>
      #include <functional>
      
      class Task
      {
          using func_t = std::function<int(int, int)>; // 函数对象类型
          // 等同于typedef std::function<int(int,int)> func_t;
      public:
          Task() // 默认构造函数
              : _x(0), _y(0), _func(nullptr)
          {
          }
          Task(int x, int y, func_t func) // 带参构造函数
              : _x(x), _y(y), _func(func)
          {
          }
          int operator()() // 函数调用运算符重载
          {
              return _func(_x, _y);
          }
      
      private:
          int _x;       // 参数1
          int _y;       // 参数2
          func_t _func; // 函数对象
      };
      生产消费过程(main.cc)

      生产和消费Task类型的数据

      #include "BlockQueue.hpp"
      #include "Task.hpp"
      #include <unistd.h>
      
      void *produce(void *arg) // 生产者线程函数
      {
          BlockQueue<Task> *bq = static_cast<BlockQueue<Task> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              Task task(i, i + 1, [](int a, int b)
                        { return a + b; });// 使用lambda表达式创建任务处理方式(回调函数)
              bq->push(task);// 生产任务
              cout << "生产者生产了一个需要处理的任务: " << i << "+" << i + 1 << endl;
              sleep(1); // 模拟生产时间,使生产者比消费者慢一点
          }
          return nullptr;
      }
      
      void *consume(void *arg) // 消费者线程函数
      {
          BlockQueue<Task> *bq = static_cast<BlockQueue<Task> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              Task task;
              bq->pop(&task);// 消费任务
              cout << "消费者获取了一个任务并处理: " << i << "+" << i + 1 << "=" << task() << endl;
              //sleep(1); // 模拟消费时间,使消费者比生产者慢一点
          }
          return nullptr;
      }
      
      int main()
      {
          BlockQueue<Task> *bq = new BlockQueue<Task>(5);
          pthread_t consumer, producer;
          pthread_create(&producer, nullptr, produce, bq); // 创建生产者线程
          pthread_create(&consumer, nullptr, consume, bq); // 创建消费者线程
      
          pthread_join(producer, nullptr);
          pthread_join(consumer, nullptr);
      
          delete bq;
          return 0;
      }
      运行结果
      • 当生产者比消费者慢时

      消费一个数据,阻塞队列为空,等待数据生产,再进行消费:生产一个,消费一个

      • 当消费者比生产者慢时

      最开始生产者会将阻塞队列填满,等待数据消费,之后等待数据消费,再进行生产:消费一个,生产一个

      7.生产者消费者模型的理解

      在之前的例子中,都只有一个生产者线程和一个消费者线程,那可不可以直接改为多个生产者线程和多个消费者线程呢

      这个问题的关键就是各个线程间能否互斥访问临界资源(阻塞队列)

      可以。因为多个生产者线程会在条件变量中等待,依次申请锁,向阻塞队列写入;多个消费者线程也会在条件变量中等待,依次申请锁,从阻塞队列读取。而管理阻塞队列的只有一把锁,也就是说,任意时刻访问(写入或读取)阻塞队列的只有一个线程,实现了互斥访问

      在生产者消费者模型中,同一时间只有一个线程访问阻塞队列(串行访问),那么它的高效体现在哪里呢

      • 对于生产者线程,不只是向阻塞队列中放入任务。放入任务之前,生产者线程还需要构建任务(从网络中、数据库中获取……)
      • 构建任务可以由多个生产者线程同时进行
      • 对于消费者线程,不只是从阻塞队列中获取任务。获取任务之后,消费者线程还需要处理任务
      • 处理任务可以由多个消费者线程同时进行
      • 生产者线程放入任务时,消费者线程可以处理任务;消费者线程获取任务时,生产者线程可以构建任务

      所以,生产者消费者模型的高效性并不体现在访问阻塞队列上,而是体现在放入任务之前和获取任务之后,多个线程并发执行。

      五、POSIX信号量

      1.引入

      在之前的代码中,要判断生产或消费条件是否满足,就必须访问临界资源进行检测就必须先获取锁

      先加锁=》再检测=》再操作=》再解锁

      有没有一种方式,在访问临界资源之前,就获取临界资源的使用情况,从而判断条件是否就绪?

      POSIX信号量就能在访问临界资源之前,得知临界资源的状态,决定下一个步骤。

      2.POSIX信号量的概念

      什么是信号量

      进行加锁后,临界资源只能作为一个整体被使用。而在实际使用中,我们需要允许不同线程同时访问一份临界资源的不同区域,从而提升效率。

      信号量(Semaphore) 是用于 进程/线程 间的 同步机制,可以控制多个进程对共享资源的访问。它提供了一种机制,将一整块公共资源划分成了一个个小资源同一时间,每个小块资源允许一个线程进行访问。

      信号量的本质

      信号量本质是一个计数器,描述临界资源可使用的小资源数目。

      信号量开始有一个可分配数值(可使用的小资源数目):

      • 申请信号量成功,则计数器 -1
      • 信号量被释放,则计数器 +1
      • 如果信号量 <= 0,则申请信号量线程需要进行等待

      (类似于count++,count--)

      信号量的使用过程

      所有线程在访问临界资源前,必须先申请信号量

      线程需要使用临界资源时:(类似于对公共资源进行预定)

      1. 申请信号量,预定资源。(信号量申请成功,就一定会有一小块资源)
      2. 申请信号量成功,等待资源的分配。
      3. 找到对应访问资源进行访问。
      4. 访问完成后,释放信号量,归还资源,以供其它进程使用。

      信号量的作用

      由于线程在访问临界资源前,必须先申请信号量,我们就能够在访问临界资源之前,提前知道临界资源的使用情况:

      • 申请信号量成功,条件就绪,可以使用一小块临界资源
      • 申请信号量失败,条件未就绪,线程等待

      通过信号量申请的成功与否,间接判断使用临界资源的条件是否就绪

      信号量的PV操作

      多个线程在访问临界资源前,都会申请信号量,所以信号量本身也是公共资源,必须保证自身操作的安全性。所以,在 申请信号量 和 释放信号量 的时候,都必须要保证 申请资源(++) 和 归还资源(- -) 操作是原子的

      而对信号量++和- - 的操作我们就叫做PV操作

      信号量的工作基于两个基本操作:P操作(wait操作)和 V操作(signal操作)

      • P操作(申请资源--):当一个进程或线程需要访问共享资源时,它会执行 P 操作,会将信号量的值减 1 。减1后,若信号量的值大于等于 0 时,表示当前有可用的资源,进程或线程可以继续访问;若信号量的值小于 0 时,表示没有可用的资源,进程或线程会被阻塞,直到有其他进程或线程归还资源。
      • V操作(归还资源++):当一个进程或线程使用完共享资源后,它会执行 V 操作,会将信号量的值加 1 。加1后,信号量的值小于等于 0 ,表示有其他进程或线程正在等待该资源,此时会唤醒一个等待的进程或线程。

      3.POSIX信号量的接口

      初始化信号量

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

      参数:

      • pshared:0表示线程间共享,非零表示进程间共享
      • value:信号量初始值

      销毁信号量

      int sem_destroy(sem_t *sem);

      等待信号量(P操作)

      功能:等待信号量,会将信号量的值减1

      int sem_wait(sem_t *sem); //P()
      • 若申请信号量成功,继续向下运行
      • 若申请信号量失败,执行流阻塞在当前位置

      发布信号量(V操作)

      功能:发布信号量,表示资源使用完毕,可以归还资源了。将信号量值加1。

      int sem_post(sem_t *sem);//V()

      4.基于环形队列的生产者消费者模型概念

      学习环形队列是为了理解POSIX信号量更广泛的应用场景

      环形队列的基本概念

      环形队列首尾相连的先进后出的数据结构,可以采用数组模拟⽤模运算来模拟环状特性

      环形结构起始状态和结束状态都是⼀样的,不好判断为空或者为满。

      判空或判满的方法:

      • 加计数器或者标记位
      • 预留⼀个空的位置,空的时候,front和tail指向同一位置;满的时候,tail和front相差一个空位置

      基于环形队列的生产者消费者模型

      环形队列作为缓冲区,实现生产者消费者模型:

      • 生产者向环形队列中放入任务
      • 消费者从环形队列中取出任务

      只有当环形队列为空或满的时候,生产者和消费者指向同一个区域

      理解:

      • 阻塞队列作为缓冲区时,是将缓冲区作为一个整体使用
      • 环形队列作为缓冲区时,是将缓冲区分为多块使用,每一块空间就是一小块临界资源

      生产者消费者模型与信号量的联系

      我们需要维护的核心规则:

      • 消费者不能超过生产者
      • 生产者不能把消费者套一个圈以上
      • 当生产者和消费者指向同一个位置:若为空,生产者先放入任务;若为满,消费者先取出任务

      所以,为了维护核心规则:

      • 对于生产者线程,关注队列中的剩余空间——空间资源定义一个信号量
      • 对于消费者线程,关注放入队列中的数据——数据资源定义一个信号量

      那么如何实现呢?

      实现思路

      举例说明:

      创建信号量space_sem和data_sem,分别管理环形队列中的空间资源和数据资源

      • 当生产者放入数据后,空间减少(space_sem--),数据增多(data_sem++)
      • 当消费者取出数据后,空间增多(space_sem++),数据减少(data_sem--)

      假设环形队列的大小为10(临界资源被分为10块),生产和消费的过程如下:

      • 当环形队列为空时,space_sem=10,data_sem=0,生产者申请空间资源信号量成功,消费者申请数据资源信号量失败,阻塞等待,所以只有生产者能向下运行
      • 当环形队列为满时,space_sem=0,data_sem=10,生产者申请空间资源信号量失败,消费者申请数据资源信号量成功,阻塞等待,所以只有消费者能向下运行
      • 其余情况,生产者和消费者都能申请信号量成功,可以并发访问不同位置

      两种生产者消费者模型的比较

      两种生产者消费者模型:

      • 阻塞队列实现缓冲区——基于阻塞队列的生产者消费者模型

      同一时间,只有一个线程(消费者或生产者)访问临界资源

      这是因为阻塞队列是一整块临界资源,同一时间只有一个线程能够访问

      • 环形队列实现缓冲区——基于环形队列的生产者消费者模型

      同一时间,可以有一个生产者和一个消费者共同访问临界资源

      这是因为环形队列被划分为不同区域,同一时间可以有多个线程访问临界资源的不同区域

      5.模型的代码实现

      实现基于RingQueue的生产者消费者模型

      生产者线程向环形队列中放入数据,消费者线程从环形队列中获取数据

      单生产者单消费者的实现

      RingQueue类的框架(RingQueue.hpp)

      我们使用POSIX信号量保证线程安全。

      生产数据(push)消费数据(pop)的过程:

      • 生产数据:申请空间资源信号量(P操作)->放入数据->释放数据资源信号量(V操作)
      • 放入数据后,可用空间变少,有效资源变多
      • 消费数据:申请数据资源信号量(P操作)->取出数据->释放空间资源信号量(V操作)
      • 放入数据后,可用空间变多,有效资源变少

      如果申请信号量失败,线程阻塞等待

      POSIX信号量变化,唤醒线程的过程:

      • 生产者唤醒消费者:生产者将数据加入缓冲区->缓冲区不为空(货物不为空)->数据资源信号量增加(data_sem++)->唤醒消费者
      • (当生产者生产完,把商品放到超市中,就能唤醒消费者继续取出)
      • 消费者唤醒生产者:消费者将数据从缓冲区中取出->缓冲区没有满(货物没有满)->空间资源信号量增加(space_sem++)->唤醒生产者
      • (当消费者取出商品后,消费者就能唤醒生产者继续生产)

      说明:

      • 生产者和消费者在队列中的位置是用两个不同的下标实现的,队列为空或满时,两个下标相同。(生产者下标——productorStep,消费者下标——consumerStep)
      • vector模拟实现环形队列
      • _capacity表示环形队列的容量
      #pragma once
      #include <iostream>
      #include <vector>
      #include <semaphore.h>
      using namespace std;
      
      template <typename T>
      class RingQueue
      {
      public:
          RingQueue(const size_t &capacity)
              : _capacity(capacity), _queue(capacity)
          {
              sem_init(&_spaceSem, 0, capacity);// 初始化空间资源信号量,初始值为队列容量
              sem_init(&_dataSem, 0, 0);// 初始化数据资源信号量,初始值为0
              productorStep=0;
              consumerStep=0;
          }
      
          ~RingQueue()
          {
              sem_destroy(&_spaceSem);
              sem_destroy(&_dataSem);
          }
      
          void push(const T &in)
          {
              sem_wait(&_spaceSem);// 等待空间资源信号量(P操作)
              _queue[productorStep++] = in;// 放入数据
              productorStep %= _capacity;
              sem_post(&_dataSem);// 释放数据资源信号量(V操作)
          }
          
      
          void pop(T *out)
          {
              sem_wait(&_dataSem);// 等待数据资源信号量(P操作)
              *out = _queue[consumerStep++];// 取出数据
              consumerStep %= _capacity;
              sem_post(&_spaceSem);// 释放空间资源信号量(V操作)
          }
      private:
          vector<T> _queue;// 环形队列容器
          size_t _capacity;// 环形队列容量
          sem_t _spaceSem;// 空间资源信号量
          sem_t _dataSem;// 数据资源信号量
          int productorStep = 0;// 生产者下标
          int consumerStep = 0;// 消费者下标
      };
      生产消费过程(main.cc)

      生产和消费int类型的数据

      #include "RingQueue.hpp"
      #include <pthread.h>
      #include <unistd.h>
      
      void *produce(void *arg) // 生产者线程函数
      {
          RingQueue<int> *rq = static_cast<RingQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              rq->push(i);
              cout << "生产者线程生产了数据: " << i << endl;
              sleep(1); // 模拟生产时间,使生产者比消费者慢一点
          }
         
          return nullptr;
      }
      
      void *consume(void *arg) // 消费者线程函数
      {
          RingQueue<int> *rq = static_cast<RingQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              int data;
              rq->pop(&data);
              cout << "消费者线程消费了数据: " << data << endl;
              //sleep(1); // 模拟消费时间,使消费者比生产者慢一点
          }
        
          return nullptr;
      }
      
      int main()
      {
          RingQueue<int> *rq = new RingQueue<int>(5); // 创建容量为5的环形队列
          pthread_t consumer, producer;
          pthread_create(&producer, nullptr, produce, rq); // 创建生产者线程
          pthread_create(&consumer, nullptr, consume, rq); // 创建消费者线程
      
          pthread_join(producer, nullptr);
          pthread_join(consumer, nullptr);
      
          delete rq;
          return 0;
      }
      运行结果
      • 当生产者比消费者慢时

      消费一个数据,环形队列为空,等待数据生产,再进行消费:生产一个,消费一个

      • 当消费者比生产者慢时

      最开始生产者会将环形队列填满,等待数据消费,之后等待数据消费,再进行生产:消费一个,生产一个

      消费顺序和生产顺序也是一一对应的。

      • 当生产者和消费者速度一致时

      只要环形队列不为空或满,生产者和消费者会并发执行(大多数情况)

      (不方便用实验证明)

      多生产者多消费者的实现

      在之前的代码中,只有一个生产者进行生产,一个消费者进行消费,那如果有多个生产者进行生产,多个消费者进行消费呢?

      我们知道在生产者消费者模型中,生产者之间是互斥的,消费者之间也是互斥的。也就是说,同一时间,最多只有一个生产者和一个消费者访问环形队列就是问题的关键

      所以,我们可以通过加锁来实现。

      RingQueue类的框架(RingQueue.hpp)

      增加的内容:

      • 增加管理生产者的互斥锁(_pmutex)和管理消费者的互斥锁(_cmutex)
      • 生产者放入数据之前要申请锁(_pmutex),消费者取出数据之前也要申请锁(_cmutex)
      #pragma once
      #include <iostream>
      #include <vector>
      #include <semaphore.h>
      #include <pthread.h>
      using namespace std;
      
      template <typename T>
      class RingQueue
      {
      public:
          RingQueue(const size_t &capacity)
              : _capacity(capacity), _queue(capacity)
          {
              sem_init(&_spaceSem, 0, capacity);// 初始化空间资源信号量,初始值为队列容量
              sem_init(&_dataSem, 0, 0);// 初始化数据资源信号量,初始值为0
              productorStep=0;
              consumerStep=0;
              pthread_mutex_init(&_pmutex, nullptr);
              pthread_mutex_init(&_cmutex, nullptr);
          }
      
          ~RingQueue()
          {
              sem_destroy(&_spaceSem);
              sem_destroy(&_dataSem);
              pthread_mutex_destroy(&_pmutex);
              pthread_mutex_destroy(&_cmutex);
          }
      
          void push(const T &in)
          {
              //!!!优化细节
              sem_wait(&_spaceSem);// 等待空间资源信号量(P操作)
              pthread_mutex_lock(&_pmutex);// 生产者加锁
              _queue[productorStep++] = in;// 放入数据
              productorStep %= _capacity;
              pthread_mutex_unlock(&_pmutex);// 生产者解锁
              sem_post(&_dataSem);// 释放数据资源信号量(V操作)
          }
          
      
          void pop(T *out)
          {
              //!!!优化细节
              sem_wait(&_dataSem);// 等待数据资源信号量(P操作)
              pthread_mutex_lock(&_cmutex);// 消费者加锁
              *out = _queue[consumerStep++];// 取出数据
              consumerStep %= _capacity;
              pthread_mutex_unlock(&_cmutex);// 消费者解锁
              sem_post(&_spaceSem);// 释放空间资源信号量(V操作)
          }
      private:
          vector<T> _queue;// 环形队列容器
          size_t _capacity;// 环形队列容量
          sem_t _spaceSem;// 空间资源信号量
          sem_t _dataSem;// 数据资源信号量
          int productorStep = 0;// 生产者下标
          int consumerStep = 0;// 消费者下标
          pthread_mutex_t _pmutex;// 生产者互斥锁
          pthread_mutex_t _cmutex;// 消费者互斥锁
      };

      优化细节:先申请信号量,再加锁

      原因:

      1. 信号量的申请过程是原子的,不需要锁保护
      2. 当一个线程持有锁时,其余线程也能够同时申请信号量。
      3. 先加锁,再申请信号量,当一个线程持有锁时,其余线程只能阻塞等待,不能申请信号量
      生产消费过程(main.cc)

      说明:

      • 创建和等待多个生产者线程和消费者线程
      #include "RingQueue.hpp"
      #include <pthread.h>
      #include <unistd.h>
      
      void *produce(void *arg) // 生产者线程函数
      {
          RingQueue<int> *rq = static_cast<RingQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              rq->push(i);
              cout << "生产者线程:" << pthread_self() << " 生产了数据: " << i << endl;
              //sleep(1); // 模拟生产时间,使生产者比消费者慢一点
          }
         
          return nullptr;
      }
      
      void *consume(void *arg) // 消费者线程函数
      {
          RingQueue<int> *rq = static_cast<RingQueue<int> *>(arg);
          for (int i = 0; i < 20; ++i)
          {
              int data;
              rq->pop(&data);
              cout << "消费者线程:" << pthread_self() << " 消费了数据: " << data << endl;
              sleep(1); // 模拟消费时间,使消费者比生产者慢一点
          }
        
          return nullptr;
      }
      
      int main()
      {
          RingQueue<int> *rq = new RingQueue<int>(5); // 创建容量为5的环形队列
          pthread_t consumer[4], producer[4];
          for (int i = 0; i < 4; ++i) {
              pthread_create(&producer[i], nullptr, produce, rq); // 创建生产者线程
              pthread_create(&consumer[i], nullptr, consume, rq); // 创建消费者线程
          }
      
          for (int i = 0; i < 4; ++i) {
              pthread_join(producer[i], nullptr);
              pthread_join(consumer[i], nullptr);
          }
      
          delete rq;
          return 0;
      }
      运行结果

      多个生产者生产数据,多个消费者消费数据,但同一时间只有一个生产者和一个消费者访问环形队列

      高效性

      生产者消费者模型的高效性并不体现在访问环形队列上,而是体现在放入任务之前和获取任务之后,多个线程并发执行。

      六、线程池

      线程池概念

      线程池(Thread Pool)是线程的一种使用模式,用于管理线程的创建和生命周期,以及提供一个用于并行执行任务的线程队列。

      • 当没有任务时,线程池中就有已经创建的线程,处于休眠状态
      • 当需要处理任务时,可以直接唤醒线程池中的线程,任务完成后线程休眠,从而省去了创建线程和销毁线程的开销

      线程池优点

      在前面的情况中,我们都是遇到任务然后创建线程再执行。但是线程的频繁创建就类似于内存的频繁申请,会给操作系统带来更大的压力,进而影响整体的性能

      一次申请好一定数量的线程,然后将线程的管理操作交给线程池,就避免了在短时间内不断创建与销毁线程的代价,线程池不但能够保证内核的充分利用,还能防止过分调度,并根据实际业务情况进行修改。 

      应用场景

      • 需要大量的线程来完成任务,且完成任务的时间比较短 

      例如:WEB服务器完成网页请求这样的任务,使用线程池技术是非常合适的。因为单个任务小,而任务数量巨大,你可以想象一个热门网站的点击次数。 但对于长时间的任务,比如一个Telnet连接请求,线程池的优点就不明显了。因为Telnet会话时间比线程的创建时间大多了。

      • 对性能要求苛刻的应用

      比如要求服务器迅速响应客户请求。

      • 接受突发性的大量请求,但不至于使服务器因此产生大量线程的应用

      突发性大量客户请求,在没有线程池情况下,将产生大量线程,虽然理论上大部分操作系统线程数目最大值不是问题,短时间内产生大量线程可能使内存到达极限,出现错误。
       

      简单实现

      示意图

      相关说明

              我们创建好了线程池之后,首次我们先是对其进行初始化操作;然后不断的向任务队列塞数据,由线程池中的线程去获取任务并执行相关操作;

      1. 任务队列(即临界资源)是会被多个执行流同时访问,因此我们需要引入互斥锁对任务队列进行保护。
      2. 线程池中的线程想要获取到任务队列中的任务,那么就必须要确保任务队列中有任务,所以我们还需引入条件变量来进行判断,如果队列中没有任务,线程池中的线程将会被挂起,直到任务队列中有任务后才被唤醒;
      3. 在thread_pool.hpp中,多线程去执行对应的方法的时候,采用的是静态成员函数,这样做的目的是解决类中存在隐藏的this指针问题,因为多线程在调用对应的函数时,该函数只有一个形参,不加static的话,那么形参个数就有两个,是不可以的;所以我们可以将this指针作为参数传递过去,就可以访问类内的成员函数了;
         

      thread_pool.hpp

      #pragma once        
      #include <iostream>    
      #include <string>    
      #include <queue>    
      #include <unistd.h>    
      #include <pthread.h> 
         
      using namespace std;
          
      namespace ns_threadpool    
      {    
          const int g_num = 3; //默认线程池大小   
          template <class T>    
          class ThreadPool    
          {    
          private:    
              int num_; //固定大小的线程池   
              vector<pthread_t> threads_; //存放线程的容器
              queue<T> task_queue_; //任务队列,使用STL的queue实现   
              pthread_mutex_t mtx_; //定义一把锁  
              pthread_cond_t cond_; //定义一个条件变量
         
          public:    
              void Lock() { pthread_mutex_lock(&mtx_);} //加锁操作     
          
              void Unlock() { pthread_mutex_unlock(&mtx_);} //解锁操作   
         
              bool IsEmpety() { return task_queue_.empty();} //判断任务队列是否为空  
          
              void Wait() { pthread_cond_wait(&cond_, &mtx_);} //让线程在条件变量下等待   
                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           
              void WakeUp() { pthread_cond_signal(&cond_);} //唤醒在条件变量下等待的线程   
         
          public:    
              ThreadPool(int num = g_num):num_(num)    
              {    
                  pthread_mutex_init(&mtx_, nullptr);    
                  pthread_cond_init(&cond_, nullptr);    
              }    
          
               ~ThreadPool()    
              {    
                  pthread_mutex_destroy(&mtx_);    
                  pthread_cond_destroy(&cond_);    
              }    
      
              //在类中要让线程执行类内成员方法,是不可行的    
              //必须让线程执行静态方法    
              static void* Rountine(void* args)    
              {    
                  pthread_detach(pthread_self());//线程分离,不需要回收   
                  ThreadPool<T>* tp = (ThreadPool<T>*)args;    
                  while(true)    
                  {    
                      tp->Lock();    
                      while(tp->IsEmpety())    
                      {    
                          tp->Wait();    
                      }    
                      T t;    
                      tp->PopTask(&t);    
                      tp->Unlock();    
                      t.Run(); //执行任务对象的Run方法,注意:先解锁再执行任务   
          
                  }    
              }    
          
              void InitThreadPool()    
              {    
                  pthread_t tid;    
                  for(int i = 0; i < num_; i++)    
                  {    
                      pthread_create(&tid, nullptr, Rountine, (void*)this/*传this指针给静态方法*/); 
                      threads_.push_back(tid);//保存线程ID   
                  }    
                  for(size_t i = 0; i < threads_.size(); i++)    
                  {    
                      cout << "线程池中第" << i << "个线程ID为:" << threads_[i] << endl;    
                  }
              }    
          
              void PushTask(const T& in)//向任务队列添加任务    
              {    
                  Lock();    
                  task_queue_.push(in);    
                  Unlock();    
                  WakeUp();//唤醒一个等待的线程    
              }    
          
              void PopTask(T* out)//从任务队列获取任务    
              {    
                  *out = task_queue_.front();    
                  task_queue_.pop();    
              }    
          };    
      }  

      Task.hpp

      #pragma once                                                                                                                                                                                                                                                                                                                                                                
      #include <iostream>    
      #include <pthread.h>    
      using namespace std;    
          
      namespace ns_task    
      {    
          class Task    
          {    
          private:    
              int x_;    
              int y_;    
              char op_;//用来表示:+ 、- 、* 、/ 、%    
          public:    
              Task(){}    
              Task(int x, int y, char op):x_(x), y_(y), op_(op){}    
          
              string show()    
              {    
                  string message = to_string(x_);    
                  message += op_;    
                  message += to_string(y_);    
                  message += "=?";    
                  return message;    
              }    
              int Run()    
              {    
                  int res = 0;    
                  switch(op_)    
                  {    
                      case '+':    
                        res = x_ + y_;    
                        break;    
                      case '-':    
                        res = x_ - y_;    
                        break;    
                      case '*':    
                        res = x_ * y_;    
                        break;    
                      case '/':    
                        res = x_ / y_;    
                        break;    
                      case '%':    
                        res = x_ % y_;    
                        break;    
                      default:    
                        cout << "bug" << endl;    
                        break;    
                  }    
                  printf("当前任务正在被:线程%lu处理,处理结果为:%d %c %d = %d\n",pthread_self(), x_, op_, y_, res);     
                  return res;    
              }    
          
              int operator()()    
              {    
                  return Run();    
              }    
          
              ~Task(){}    
          };    
      } 

      main.cpp

      #include "thread_pool.hpp"    
      #include "Task.hpp"    
      #include <ctime>    
      #include <cstdlib> 
         
      using namespace ns_threadpool;    
      using namespace ns_task; 
         
      int main()                                                                                                                                                                            
      {    
          ThreadPool<Task>* tp = new ThreadPool<Task>();//创建线程池    
          tp->InitThreadPool();  //进行初始化  
          srand((long long)time(nullptr));//生产随机数    
          while(true) //不断向任务队列塞数据   
          {    
              Task t(rand() % 20 + 1, rand() % 10 + 1, "+-*/%"[rand() % 5]);    
              tp->PushTask(t);    
              sleep(1);    
          }    
          delete tp; //释放线程池资源
          return 0;    
      }

      运行结果

      七、单例模式

      概念

      单例(Singleton)模式,是一种常用的软件设计模式。在它的核心结构中只包含一个被称为单例的特殊类。通过单例模式可以保证系统中,应用该模式的类一个类只有一个实例。即一个类只有一个对象实例 ;

      使用场景

      1. 语义上只需要一个
      2. 该对象内部存在大量的空间,保存了大量的数据,如果允许该对象存在多份,或者允许发生各种拷贝,内存中存在冗余数据;

      单例模式通常有两种形式

      1. 饿汉式:吃完饭, 立刻洗碗, 这种就是饿汉方式. 因为下一顿吃的时候可以立刻拿着碗就能吃饭。
      2. 懒汉式:吃完饭, 先把碗放下, 然后下一顿饭用到这个碗了再洗碗, 就是懒汉方式。

      懒汉方式最核心的思想是 "延时加载",能够优化服务器的启动速度。

      举个简单的例子,使用malloc后,如果立刻在物理内存中开辟了空间,就是饿汉模式;如果在使用这块空间时,才真正开辟物理空间,就是懒汉模式(和写时拷贝很像)

      1.饿汉实现方式

      该模式在类被加载时就会实例化一个对象

      template <typename T> 
      class Singleton 
      { 
      private:
          static Singleton<T> data;//饿汉模式,在加载的时候对象就已经存在了 
      public: 
          static Singleton<T>* GetInstance() 
          { 
              return &data; 
          } 
      };

      该模式能简单快速的创建一个单例对象,而且是线程安全的(只在类加载时才会初始化,以后都不会)。

      但缺点就是,不管你是否使用,都会直接创建一个对象,消耗一定的性能(当然很小很小,几乎可以忽略不计,所以这种模式在很多场合十分常用

      2. 懒汉实现方式

      该模式只在你使用对象时,才会生成单例对象(比如调用GetInstance方法) 

      template <typename T> 
      class Singleton 
      { 
      private:
          static Singleton<T>* inst; //懒汉式单例,只有在调用GetInstance时才会实例化一个单例对象
      public: 
          static Singleton<T>* GetInstance() 
          { 
              if (inst == NULL) 
              { 
                  inst = new Singleton<T>(); 
              }
          return inst; 
          } 
      };

      看上去,这段代码没什么明显问题,但它不是线程安全的

      假设当前有多个线程同时调用GetInstance()方法,由于当前还没有对象生成,那么就会由多个线程创建多个对象,不符合单例模式。

      懒汉模式(线程安全版本)

      // 懒汉模式, 线程安全 
      template <typename T> 
      class Singleton 
      {
      private: 
          static Singleton<T>* inst; 
          static std::mutex lock; 
      public: 
          static T* GetInstance() 
          { 
              if (inst == NULL) // 双重判定空指针, 降低锁冲突的概率, 提高性能 
              {                 
                  lock.lock();  // 使用互斥锁, 保证多线程情况下也只调用一次 new
                  if (inst == NULL) 
                  { 
                      inst = new T(); 
                  }
                  lock.unlock(); 
              }
          return inst;
          } 
      };

      这种形式是在懒汉方式的基础上增加的,当多个线程调用GetInstance方法时,此时类中没有对象,那么多个线程就会来到锁的位置,竞争锁。必然只能有一个线程竞争锁成功,此时再次判断有没有对象被创建(就是inst指针),如果没有就会new一个对象,如果有就会解锁,并返回已有的对象;

      总的来说,这样的形式使得多个线程调用GetInstance方法时,无论成功与否,都会有返回值;

      单例模式实现线程池(懒汉方式)

      thread_pool.hpp

      #pragma once        
      #include <iostream>    
      #include <string>    
      #include <queue>    
      #include <unistd.h>    
      #include <pthread.h> 
         
      using namespace std;
          
      namespace ns_threadpool    
      {    
          const int g_num = 3;   
          template <class T>    
          class ThreadPool    
          {    
          private:    
              int num_;   
              vector<pthread_t> threads_; 
              queue<T> task_queue_;   
              pthread_mutex_t mtx_; 
              pthread_cond_t cond_; 
              static ThreadPool<T>* ins;//类内的静态指针
         
              private:
              //构造函数私有化
              ThreadPool(int num = g_num):num_(num)    
              {    
                  pthread_mutex_init(&mtx_, nullptr);    
                  pthread_cond_init(&cond_, nullptr);    
              }    
              ThreadPool(const ThreadPool<T>&) = delete; //禁止拷贝构造函数
              ThreadPool<T>& operator=(const ThreadPool<T>&) = delete; //禁止赋值操作符重载
          public:    
              static ThreadPool<T>* GetInstance()//单例模式获取线程池对象    
              {    
                  static pthread_mutex_t lock = PTHREAD_MUTEX_INITIALIZER;
                 if(ins == nullptr)//先判断,不用直接申请锁,提升效率
                 {
                      pthread_mutex_lock(&lock);
                      if(ins == nullptr)    
                      {    
                          ins = new ThreadPool<T>(); 
                          ins->InitThreadPool();
                          cout << "线程池单例对象创建成功!" << endl;   
                      }    
                      pthread_mutex_unlock(&lock);
                 }
                  return ins;    
              }
              void Lock() { pthread_mutex_lock(&mtx_);}     
          
              void Unlock() { pthread_mutex_unlock(&mtx_);}   
         
              bool IsEmpety() { return task_queue_.empty();} 
          
              void Wait() { pthread_cond_wait(&cond_, &mtx_);} 
                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           
              void WakeUp() { pthread_cond_signal(&cond_);} 
         
          public:    
             
          
               ~ThreadPool()    
              {    
                  pthread_mutex_destroy(&mtx_);    
                  pthread_cond_destroy(&cond_);    
              }    
      
              //在类中要让线程执行类内成员方法,是不可行的    
              //必须让线程执行静态方法    
              static void* Rountine(void* args)    
              {    
                  pthread_detach(pthread_self());
                  ThreadPool<T>* tp = (ThreadPool<T>*)args;    
                  while(true)    
                  {    
                      tp->Lock();    
                      while(tp->IsEmpety())    
                      {    
                          tp->Wait();    
                      }    
                      T t;    
                      tp->PopTask(&t);    
                      tp->Unlock();    
                      t.Run(); 
          
                  }    
              }    
          
              void InitThreadPool()    
              {    
                  pthread_t tid;    
                  for(int i = 0; i < num_; i++)    
                  {    
                      pthread_create(&tid, nullptr, Rountine, (void*)this/*传this指针给静态方法*/); 
                      threads_.push_back(tid);
                  }    
                  for(size_t i = 0; i < threads_.size(); i++)    
                  {    
                      cout << "线程池中第" << i << "个线程ID为:" << threads_[i] << endl;    
                  }
              }    
          
              void PushTask(const T& in)    
              {    
                  Lock();    
                  task_queue_.push(in);    
                  Unlock();    
                  WakeUp();
              }    
          
              void PopTask(T* out)   
              {    
                  *out = task_queue_.front();    
                  task_queue_.pop();    
              }    
          };    
          template <class T>    
          ThreadPool<T>* ThreadPool<T>::ins = nullptr;
      }  

      main.cpp

      #include "thread_pool.hpp"    
      #include "Task.hpp"    
      #include <ctime>    
      #include <cstdlib> 
         
      using namespace ns_threadpool;    
      using namespace ns_task; 
         
      int main()                                                                                                                                                                            
      {    
          ThreadPool<Task>* tp = ThreadPool<Task>::GetInstance();//获取线程池单例对象 
          srand((long long)time(nullptr));
          while(true) 
          {    
              Task t(rand() % 20 + 1, rand() % 10 + 1, "+-*/%"[rand() % 5]);    
              tp->PushTask(t);    
              cout<<"实例化线程池对象的地址为:" << tp << endl;//打印实例化线程池对象地址
              sleep(1);    
          }    
          return 0;    
      }

      Task.hpp同上

      运行结果

      运行后发现对象的地址是一样的,表明单例调用成功了,只存在一份

      八、STL,智能指针和线程安全

      STL中的容器是否是线程安全的?

      不是。

      原因是,STL 的设计初衷是将性能挖掘到极致,而一旦涉及到加锁保证线程安全,会对性能造成巨大的影响。而且对于不同的容器,加锁方式的不同,性能可能也不同(例如hash表的锁表和锁桶)。

      因此 STL 默认不是线程安全。 如果需要在多线程环境下使用, 往往需要调用者自行保证线程安全。

      智能指针是否是线程安全的?

      1. 对于 unique_ptr,由于只是在当前代码块范围内生效, 因此不涉及线程安全问题。
      2. 对于 shared_ptr,多个对象需要共用一个引用计数变量,所以会存在线程安全问题。但是标准库实现的时候考虑到了这个问题,基于原子操作(CAS)的方式保证 shared_ptr 能够高效,原子的操作引用计数。

      九、其他常见的各种锁

      悲观锁

      在每次取数据时,总是担心数据会被其他线程修改,所以会在取数据前先加锁(读锁,写锁,行锁等),当其他线程想要访问数据时,被阻塞挂起。

      悲观锁适用于写多读少的情况下,即:需要频繁的写数据时候,可以考虑使用悲观锁。

      乐观锁

      每次取数据时候,总是乐观的认为数据不会被其他线程修改,因此不上锁。但是在更新数据前,会判断其他数据在更新前有没有对数据进行修改。

      乐观锁适用于读多写少的情况下,即:读数据多的时候,可以考虑使用乐观锁。

      乐观锁主要采用两种方式:版本号机制和CAS操作。

      CAS操作(compare and swap):当需要更新数据时,判断当前内存值和之前取得的值是否相等。如果相等则用新值更新。若不等则失败,失败则重试,一般是一个自旋的过程,即不断重试。

      自旋锁

      当一个线程在获取锁的时候,如果锁已经被其它线程获取,那么该线程将循环等待,然后不断的判断锁是否能够被成功获取,直到获取到锁才会退出循环。  

      自旋锁适用于 线程在占用临界资源的时间较少 的情况

      其他的锁:

      自旋锁:

      Linux 提供的自旋锁接口:

      #include <pthread.h>

      函数名 描述 参数 返回值
      pthread_spin_init 初始化自旋锁 pthread_spinlock_t *lock, int pshared 成功返回0,失败返回错误代码
      pthread_spin_lock 获取自旋锁 pthread_spinlock_t *lock 成功返回0,失败返回错误代码
      pthread_spin_trylock 尝试获取自旋锁 pthread_spinlock_t *lock 成功返回0,锁已被持有时返回 EBUSY
      pthread_spin_unlock 释放自旋锁 pthread_spinlock_t *lock 成功返回0,失败返回错误代码
      pthread_spin_destroy 销毁自旋锁,释放资源 pthread_spinlock_t *lock 成功返回0,失败返回错误代码

      读写锁   

      有些时候,公共数据修改少,读取多。

      通常而言,在读的过程中,往往伴随着查找的操作,中间耗时很长。给这种代码段加锁,会极大地
      降低我们程序的效率。

      读写锁就可以专门处理这种读多写少的情况

      读者写者模型

      读者写者关系:

      • 写者之间:互斥
      • 读者和写着之间:同步和互斥
      • 读者之间:没有关系

      读者写者模型 和 生产者消费者模型 的区别:消费者会拿走数据,读者不会

      读锁和写锁的关系

      注意:写独占,读共享,读锁优先级高

      读写锁接口

      #include <pthread.h>
      
      // 初始化读写锁
      int pthread_rwlock_init(pthread_rwlock_t *restrict rwlock, const pthread_rwlockattr_t *restrict attr);
      
      // 加读锁(阻塞式)
      int pthread_rwlock_rdlock(pthread_rwlock_t *rwlock);
      // 尝试加读锁(非阻塞,成功返回0,失败返回EBUSY)
      int pthread_rwlock_tryrdlock(pthread_rwlock_t *rwlock);
      
      // 加写锁(阻塞式)
      int pthread_rwlock_wrlock(pthread_rwlock_t *rwlock);
      // 尝试加写锁(非阻塞)
      int pthread_rwlock_trywrlock(pthread_rwlock_t *rwlock);
      
      // 解锁
      int pthread_rwlock_unlock(pthread_rwlock_t *rwlock);
      
      // 销毁锁
      int pthread_rwlock_destroy(pthread_rwlock_t *rwlock);
      
      Logo

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

      更多推荐