File tree Expand file tree Collapse file tree
Expand file tree Collapse file tree Original file line number Diff line number Diff line change 1+ # --coding: utf-8--
2+ # 2017/8/7
3+
4+ '分布式进程,服务器端'
5+
6+ import random , time , Queue
7+ from multiprocessing .managers import BaseManager
8+
9+ # 发送任务队列
10+ task_queue = Queue .Queue ()
11+ # 接收结果的队列
12+ result_queue = Queue .Queue ()
13+
14+ class QueueManager (BaseManager ):
15+ pass
16+
17+ # 把两个Queue都注册到网络上,callable参数关联了Queue对象
18+ QueueManager .register ('get_task_queue' , callable = lambda : task_queue )
19+ QueueManager .register ('get_result_queue' , callable = lambda : result_queue )
20+ # 绑定端口5000,设置验证码‘abc’
21+ manager = QueueManager (address = ('' , 5000 ), authkey = 'abc' )
22+ manager .start () # 启动queue
23+ # 获得通过网络访问的Queue对象
24+ task = manager .get_task_queue ()
25+ result = manager .get_result_queue ()
26+ # 放入任务
27+ for i in range (10 ):
28+ n = random .randint (0 , 10000 )
29+ print 'Put task %d...' % n
30+ task .put (n )
31+ # 从result队列读取结果
32+ print 'Try get results...'
33+ for i in range (10 ):
34+ r = result .get (timeout = 10 )
35+ print 'Result: %s' % r
36+ # 关闭
37+ manager .shutdown ()
Original file line number Diff line number Diff line change 1+ # --coding: utf-8--
2+ # 2017/8/7
3+
4+ '分布式进程,客户端'
5+ import time , sys , Queue
6+ from multiprocessing .managers import BaseManager
7+
8+ # 创建类似的QueueManager:
9+ class QueueManager (BaseManager ):
10+ pass
11+
12+ # 由于这个QueueManager只从网络上获取Queue,所以注册时只提供名字:
13+ QueueManager .register ('get_task_queue' )
14+ QueueManager .register ('get_result_queue' )
15+
16+ # 连接到服务器,也就是运行taskmanager.py的机器:
17+ server_addr = '127.0.0.1'
18+ print ('Connect to server %s...' % server_addr )
19+ # 端口和验证码注意保持与taskmanager.py设置的完全一致:
20+ m = QueueManager (address = (server_addr , 5000 ), authkey = 'abc' )
21+ # 从网络连接:
22+ m .connect ()
23+ # 获取Queue的对象:
24+ task = m .get_task_queue ()
25+ result = m .get_result_queue ()
26+ # 从task队列取任务,并把结果写入result队列:
27+ for i in range (10 ):
28+ try :
29+ n = task .get (timeout = 1 )
30+ print ('run task %d * %d...' % (n , n ))
31+ r = '%d * %d = %d' % (n , n , n * n )
32+ time .sleep (1 )
33+ result .put (r )
34+ except Queue .Empty :
35+ print ('task queue is empty.' )
36+ # 处理结束:
37+ print ('worker exit.' )
You can’t perform that action at this time.
0 commit comments