
正文
python_分布式进程中遇到的问题
提示:扫一扫查出行【扫一扫了解最新限行尾号】
复制提示
看文档学习分布式进程中遇到了一下问题,文档里面例题是python2.X,我用的python3.x,就出现了一下莫名奇妙的问题,最终版代码先呈上:
taskManager.py
# coding:utf-8
# taskManager.py for windows 服务器端 import queue
from multiprocessing.managers import BaseManager
from multiprocessing import freeze_support
# 任务个数
task_number = 10
# 定义收发队列
task_queue = queue.Queue(task_number);
result_queue = queue.Queue(task_number);
def get_task():
return task_queue
def get_result():
return result_queue
#创建类似的QueueManager:
class QueueManager(BaseManager):
pass
def win_run():
# Windows 下绑定调用接口不能使用lambda,所以只能先定义函数再绑定
QueueManager.register('get_task_queue',callable = get_task)
QueueManager.register('get_result_queue',callable=get_result)
#绑定端口并设置验证口令,windows下需要填写IP地址,Linux下不填默认为本地 ip地址为本地ip地址
manager = QueueManager(address = ('192.xxx.xx.xxx',8001),authkey =b'qiye')
# 启动
manager.start()
try:
# 通过网络获取任务队列和结果队列
task = manager.get_task_queue()
result = manager.get_result_queue()
#添加任务
for url in ["ImageUrl_"+str(i) for i in range(10)]:
print('put task %s ...'%url)
task.put(url)
print('try get result...')
for i in range(10):
print('result is %s' %result.get(timeout=10))
except:
print('Manager error')
finally:
# 一定要关闭,否则会报管道未关闭的错误
manager.shutdown() if __name__=='__main__':
# Windows 下多进程可能会有问题,添加这句可以缓解
freeze_support()
win_run()
taskWorker.py
# coding:utf-8
import time
from multiprocessing.managers import BaseManager
# 创建类似的 QueueManager:
class QueueManager(BaseManager):
pass
#第一步: 使用QueueManager注册用于获取Queue的方法名称
QueueManager.register('get_task_queue')
QueueManager.register('get_result_queue')
#第二步:连接到服务器:
server_addr = '192.xxx.xx.xxx'
print('Connect to server %s...'%server_addr)
#端口和验证口令注意保持与服务器进程保持一致:
m = QueueManager(address=(server_addr,8001),authkey=b'qiye')
#从网络连接
m.connect()
#第三部:获取Queue的对象:
task = m.get_task_queue()
result = m.get_result_queue()
#第四步:从task队列获取任务,并把结果写入result队列:
while(not task.empty()):
image_url = task.get(True,timeout = 5)
print('run task download %s ...'%image_url)
time.sleep(1)
result.put('%s--->success'%image_url)
#处理结束
print('worker exit.')
先运行 taskManager.py 服务器端代码,再快速运行 taskWorker.py 客户端代码 运行结果依次如下:









