# encoding: utf-8
__author__ = 'zhanghe'
import random
import time
import Queue
from multiprocessing.managers import BaseManager
# :
task_queue = Queue.Queue()
# :
result_queue = Queue.Queue()
# BaseManagerQueueManager:
class QueueManager(BaseManager):
pass
# Queue, callableQueue:
QueueManager.register('get_task_queue', callable=lambda: task_queue)
QueueManager.register('get_result_queue', callable=lambda: result_queue)
# 5000, 'abc':
manager = QueueManager(address=('', 5000), authkey='abc')
# Queue:
manager.start()
# Queue:
task = manager.get_task_queue()
result = manager.get_result_queue()
# :
for i in range(10):
n = random.randint(0, 10000)
print('Put task %d...' % n)
task.put(n)
# result:
print('Try get results...')
# for i in range(10):
while True:
try:
r = result.get(timeout=2)
print('Result: %s' % r)
except Queue.Empty:
# 2S
print('result queue is empty.')
# :
manager.shutdown()