Python producer/consumer with process and pool
multiprocessing, producer-consumer, python
Solution
Hi I had found this answer,
from multiprocessing import Manager, Process, Pool,Queue
from Queue import Empty
def writer(queue):
for i in range(10):
queue.put(i)
print 'put %i size now %i'%(i, queue.qsize())
def reader(id, queue):
for i in range(10):
try:
cnt = queue.get(1,1)
print '%i got %i size now %i'%(id, cnt, queue.qsize())
except Empty:
pass
class Managerss:
def __init__(self):
self.queue= Queue()
self.NUMBER_OF_PROCESSES = 3
def start(self):
self.p = Process(target=writer, args=(self.queue,))
self.p.start()
self.workers = [Process(target=reader, args=(i, self.queue,))
for i in xrange(self.NUMBER_OF_PROCESSES)]
for w in self.workers:
w.start()
def join(self):
self.p.join()
for w in self.workers:
w.join()
if __name__ == '__main__':
m= Managerss()
m.start()
m.join()
Hope it helps
Problem
I'm trying to write simple code for producer consumer with process. Producer is a process. For consumer I'm getting processes from Pool. ``` from multiprocessing import Manager, Process, Pool from time import sleep def writer(queue): for i in range(10): queue.put(i) print 'put 1 size now ',queue.qsize() sleep(1) def reader(queue): print 'in reader' for i in range(10): queue.get(1) print 'got 1 size now ', queue.qsize() if __name__ == '__main__': q = Manager().Queue() p = Process(target=writer, args=(q,)) p.start() pool = Pool() c = pool.apply_async(reader,q) ``` But I'm getting error, ``` Process Process-2: Traceback (most recent call last): File "/usr/lib/python2.7/multiprocessing/process.py", line 258, in _bootstrap self.run() File "/usr/lib/python2.7/multiprocessing/process.py", line 114, in run self._target(*self._args, **self._kwargs) File "pc.py", line 5, in writer queue.put(i) File "<string>", line 2, in put File "/usr/lib/python2.7/multiprocessing/managers.py", line 758, in _callmethod conn.send((self._id, methodname, args, kwds)) IOError: [Errno 32] Broken pipe ``` Can anyone point me, where am I going wrong.