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.

Original source