标签:大数据 通讯 阻塞 ESS main time 消费者 平衡 img
进程彼此之间互相隔离,要实现进程间通信(IPC),multiprocessing模块支持两种形式:队列和管道,这两种方式都是使用消息传递的
创建队列的类(底层就是以管道和锁定的方式实现):
Queue([maxsize]):创建共享的进程队列,Queue是多进程安全的队列,可以使用Queue实现多进程之间的数据传递。
参数介绍:
maxsize是队列中允许最大项数,省略则无大小限制。
但需要明确:
1、队列内存放的是消息而非大数据
2、队列占用的是内存空间,因而maxsize即便是无大小限制也受限于内存大小
主要方法:
q.put方法用以插入数据到队列中。
q.get方法可以从队列读取并且删除一个元素。
from multiprocessing import Process,Queue q=Queue(3) #put ,get ,put_nowait,get_nowait,full,empty q.put(1) q.put(2) q.put(3) print(q.full()) #满了 # q.put(4) #再放就阻塞住了 print(q.get()) print(q.get()) print(q.get()) print(q.empty()) #空了 # print(q.get()) #再取就阻塞住了 True 1 2 3 True
生产者指的是生产数据的任务,消费者指的是处理数据的任务,在并发编程中,
如果生产者处理速度很快,而消费者处理速度很慢,那么生产者就必须等待消费者处理完,
才能继续生产数据。同样的道理,如果消费者的处理能力大于生产者,那么消费者就必须等待生产者。
为了解决这个问题于是引入了生产者和消费者模式。
生产者消费者模式是通过一个容器来解决生产者和消费者的强耦合问题。生产者和消费者彼此之间不直接通讯,而通过阻塞队列来进行通讯,所以生产者生产完数据之后不用等待消费者处理,直接扔给阻塞队列,消费者不找生产者要数据,而是直接从阻塞队列里取,阻塞队列就相当于一个缓冲区,平衡了生产者和消费者的处理能力。
这个阻塞队列就是用来给生产者和消费者解耦的
import time def producer(): for i in range(3): res = f"包子 {i}" time.sleep(0.5) print(f"生产者生产了{res}") consumer(res) def consumer(res): time.sleep(1) print(f"消费者吃了{res}") if __name__ == ‘__main__‘: producer() 生产者生产了包子 0 消费者吃了包子 0 生产者生产了包子 1 消费者吃了包子 1 生产者生产了包子 2 消费者吃了包子 2
from multiprocessing import Process,Queue import time def producer(q): for i in range(3): res = f"包子 {i}" time.sleep(0.5) print(f"生产者生产了{res}") # 把生产的给队列保存 q.put(res) def consumer(q): while True:# 消费者一直接收 res = q.get() time.sleep(1) print(f"消费者吃了{res}") if __name__ == ‘__main__‘: q = Queue() p1 = Process(target=producer,args=(q,)) p2 = Process(target=consumer,args=(q,)) p1.start() p2.start() print(‘主‘) 主 生产者生产了包子 0 生产者生产了包子 1 生产者生产了包子 2 消费者吃了包子 0 消费者吃了包子 1 消费者吃了包子 2
此时的问题是主进程永远不会结束,原因是:生产者p在生产完后就结束了,但是消费者c在取空了q之后,则一直处于死循环中且卡在q.get()这一步。
解决方式无非是让生产者在生产完毕后,往队列中再发一个结束信号,这样消费者在接收到结束信号后就可以break出死循环
队列先进先出
from multiprocessing import Process,Queue import time def producer(q): for i in range(3): res = f"包子 {i}" time.sleep(0.5) print(f"生产者生产了{res}") # 把生产的给队列保存 q.put(res) def consumer(q): while True:# 消费者一直接收 res = q.get() if res == None: break time.sleep(1) print(f"消费者吃了{res}") if __name__ == ‘__main__‘: q = Queue() p1 = Process(target=producer,args=(q,)) p2 = Process(target=consumer,args=(q,)) p1.start() p2.start() p1.join()# 主进程等待p1子进程执行完毕--即生产者生产完毕 q.put(None) print(‘主‘) 生产者生产了包子 0 生产者生产了包子 1 生产者生产了包子 2 消费者吃了包子 0 主 消费者吃了包子 1 消费者吃了包子 2
但上述解决方式,在有多个生产者和多个消费者时,我们则需要用一个很low的方式去解决,有几个消费者就需要发送几次结束信号:相当low,例如
if __name__ == ‘__main__‘: q=Queue() #生产者们:即厨师们 p1=Process(target=producer,args=(q,‘egon1‘,‘包子‘)) p2=Process(target=producer,args=(q,‘egon2‘,‘骨头‘)) p3=Process(target=producer,args=(q,‘egon3‘,‘泔水‘)) #消费者们:即吃货们 c1=Process(target=consumer,args=(q,‘alex1‘)) c2=Process(target=consumer,args=(q,‘alex2‘)) #开始 p1.start() p2.start() p3.start() c1.start() c2.start() p1.join() p2.join() p3.join() q.put(None) q.put(None) q.put(None) print(‘主‘)
标签:大数据 通讯 阻塞 ESS main time 消费者 平衡 img
原文地址:https://www.cnblogs.com/foremostxl/p/9729652.html