1:事件机制共享队列:
利用消息机制在两个队列中,通过传递消息,实现可以控制的生产者消费者问题
要求:readthread读时,writethread不能写;writethread写时,readthread不能读。
基本方法 时间类(Event)
·set:设置事件。将标志位设为True。
wait:等待事件。会将当前线程阻塞,直到标志位变为True。
clear:清除事件。将标志位设为False。
set() clear() 函数的交替执行 也就是消息传递的本质
1模版: 2 3基本code 4# 事件消息机制 5import queue 6import threading 7import random 8from threading import Event 9from threading import Thread 10class WriteThread(Thread): 11 def __init__(self,q,wt,rt): 12 super().__init__(); 13 self.queue=q; 14 self.rt=rt; 15 self.wt=wt; 16 def run(self): 17 self.rt.set() 18 19 self.wt.wait(); 20 self.wt.clear(); 21 22class ReadThread(Thread): 23 def __init__(self,q,wt,rt): 24 super().__init__(); 25 self.queue=q; 26 self.rt=rt; 27 self.wt=wt; 28 def run(self): 29 while True: 30 self.rt.wait(); 31 self.wt.wait(); 32 self.wt.clear()
参考代码:
1# -*- coding: utf-8 -*- 2""" 3Created on Tue Sep 10 20:10:10 2019 4 5@author: DGW-PC 6""" 7# 事件消息机制 8import queue 9import threading 10import random 11from threading import Event 12from threading import Thread 13 14class WriteThread(Thread): 15 def __init__(self,q,wt,rt): 16 super().__init__(); 17 self.queue=q; 18 self.rt=rt; 19 self.wt=wt; 20 def run(self): 21 data=[random.randint(1,100) for _ in range(0,10)]; 22 self.queue.put(data); 23 print("WriteThread写队列:",data); 24 self.rt.set(); # 发送读事件 25 print("WriteThread通知读"); 26 print("WriteThread等待写"); 27 self.wt.wait(); 28 print("WriteThread收到写事件"); 29 self.wt.clear(); 30 6 31class ReadThread(Thread): 32 def __init__(self,q,wt,rt): 33 super().__init__(); 34 self.queue=q; 35 self.rt=rt; 36 self.wt=wt; 37 def run(self): 38 while True: 39 self.rt.wait();# 等待写事件 带来 40 print("ReadThread 收到读事件"); 41 print("ReadThread 开始读{0}".format(self.queue.get())); 42 print("ReadThread 发送写事件"); 43 self.wt.set(); 44 self.rt.clear(); 45q=queue.Queue(); 46rt=Event(); 47wt=Event(); 48writethread=WriteThread(q,wt,rt); # 实例化对象的 49readthread=ReadThread(q,wt,rt); # 实例化对象的 50 51writethread.start(); 52readthread.start();
2:条件锁同步生产者消费者
作用: 在保护互斥资源的基础上,增加了条件判断的机制
即为使用wait() 函数 判断不满足当前条件的基础上,让当前线程的阻塞。
其他线程如果生成了满足了条件的资源 使用notify() notifyALl() 函数将刮起线程唤醒。
使用了 threading 的Condition 类
acquire() : 锁住当前资源
relarse() :释放当前锁住的资源
wait:挂起当前线程, 等待唤起 。
• notify:唤起被 wait 函数挂起的线程 。
• notif计All:唤起所有线程,防止线程永远处于沉默状态 。
模版:
1基本code 2from threading import Thread 3from threading import Condition 4import random 5import time 6lock=Condition(); # 声明条件锁 7flag=0; 8def cnsumer(): 9 lock.acquire(); 10 while flag==0: 11 lock.wait(); 12 13 业务代码--- 14lock.relarse(); 15 16def product(): 17 lock.acquire(); 18 19 释放锁之前对控制变量进行操作,数据的操作控制 可以作为全局变量来锁定 20 lock.notifyALl(); 21 lock.relarse();
参考代码code:
1# -*- coding: utf-8 -*- 2""" 3Created on Wed Sep 11 21:40:41 2019 4 5@author: DGW-PC 6""" 7# 条件锁生产者消费者 8from threading import Thread 9from threading import Condition 10import random 11import time 12 13flag=0; # 声明控制标志 14goods=0; # 事物表示 15lock=Condition(); 16def consumer(x): 17 global flag; 18 global goods; 19 lock.acquire(); # 取得锁 20 while flag==0: # 便于多次进行消费 21 print("consumer %d进入等待" % x); 22 lock.wait(); 23 print("consumer {0}:消费了{1}".format(x,goods));# format 次序从0开始 24 flag-=1; 25 lock.release(); #释放锁 26 27def product(x): 28 global flag; 29 global goods; 30 time.sleep(3); 31 lock.acquire(); 32 goods=random.randint(1,1000); 33 print("product {0} 产生了{1}".format(x,goods)); 34 flag+=1; 35 lock.notifyAll(); 36 lock.release(); 37 38threads=[]; 39 40for i in range(0,2): 41 t1=Thread(target=consumer,args=(i,)); 42 t2=Thread(target=product,args=(i,)); 43 t1.start(); 44 t2.start(); 45 threads.append(t1); 46 threads.append(t2); 47 48for x in threads: 49 x.join();