Python自定义主从分布式架构实例分析
Python  /  管理员 发布于 7年前   233
本文实例讲述了Python自定义主从分布式架构。分享给大家供大家参考,具体如下:
环境:Win7 x64,Python 2.7,APScheduler 2.1.2。
原理图如下:
代码部分:
(1)、中心节点:
#encoding=utf-8#author: walker#date: 2014-12-03#function: 中心节点(主要功能是分配任务)import SocketServer, socket, QueueCenterIP = '127.0.0.1' #中心节点IPCenterListenPort = 9999 #中心节点监听端口CenterClient = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) #中心节点用于发送网络消息的socketTaskQueue = Queue.Queue() #任务队列#获取任务队列def GetTaskQueue(): for i in range(1, 11): TaskQueue.put(str(i))#CenterServer的回调函数,在接受到udp报文是触发class MyUDPHandler(SocketServer.BaseRequestHandler): def handle(self): data = self.request[0].strip() socket = self.request[1] print(data) if data.startswith('wait'): vec = data.split(':') if len(vec) != 3: print('Error: len(vec) != 3') else: nodeIP = vec[1] nodeListenPort = vec[2] nodeID = nodeIP + ':' + nodeListenPort if not TaskQueue.empty(): task = TaskQueue.get() print('send task ' + task + ' to ' + nodeID) CenterClient.sendto('task:' + task, (nodeIP, int(nodeListenPort))) else: print('TaskQueue is empty!')GetTaskQueue() #获取任务队列CenterServer = SocketServer.UDPServer((CenterIP, CenterListenPort), MyUDPHandler)print('Listen port ' + str(CenterListenPort) + ' ...')CenterServer.serve_forever()(2)、任务节点:
#encoding=utf-8#author: walker#date: 2014-12-03#function: 任务节点(请求/接收/执行任务)import time, socket, SocketServerfrom apscheduler.scheduler import SchedulerCenterIP = '127.0.0.1' #中心节点IPCenterListenPort = 9999 #中心节点监听端口NodeIP = socket.gethostbyname(socket.gethostname()) #任务节点自身IPNodeClient = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) #任务节点用于发送网络消息的socket#任务:发送网络信息def jobSendNetMsg(): msg = '' if NodeServer.TaskState == 'wait': msg = 'wait:' + NodeIP + ':' + str(NodeListenPort) elif NodeServer.TaskState == 'exec': msg = 'exec:' + NodeIP + ':' + str(NodeListenPort) print(msg) NodeClient.sendto(msg, (CenterIP, CenterListenPort))#添加并启动定时任务def InitTimer(): sched = Scheduler() sched.add_interval_job(jobSendNetMsg, seconds=1) sched.start()#执行任务def ExecTask(task): print('ExecTask ' + task + ' ...') time.sleep(2) print('ExecTask ' + task + ' over')#NodeServer的回调函数,在接受到udp报文是触发class MyUDPHandler(SocketServer.BaseRequestHandler): def handle(self): data = self.request[0].strip() socket = self.request[1] print('recv data: ' + data) if data.startswith('task'): vec = data.split(':') if len(vec) != 2: print('Error: len(vec) != 2') else: task = vec[1] self.server.TaskState = 'exec' ExecTask(task) self.server.TaskState = 'wait'InitTimer()NodeServer = SocketServer.UDPServer(('', 0), MyUDPHandler)NodeServer.TaskState = 'wait' #(exec/wait)NodeListenPort = NodeServer.server_address[1]print('NodeListenPort:' + str(NodeListenPort))NodeServer.serve_forever()更多关于Python相关内容感兴趣的读者可查看本站专题:《Python URL操作技巧总结》、《Python图片操作技巧总结》、《Python数据结构与算法教程》、《Python Socket编程技巧总结》、《Python函数使用技巧总结》、《Python字符串操作技巧汇总》、《Python入门与进阶经典教程》及《Python文件与目录操作技巧汇总》
希望本文所述对大家Python程序设计有所帮助。
122 在
学历:一种延缓就业设计,生活需求下的权衡之选中评论 工作几年后,报名考研了,到现在还没认真学习备考,迷茫中。作为一名北漂互联网打工人..123 在
Clash for Windows作者删库跑路了,github已404中评论 按理说只要你在国内,所有的流量进出都在监控范围内,不管你怎么隐藏也没用,想搞你分..原梓番博客 在
在Laravel框架中使用模型Model分表最简单的方法中评论 好久好久都没看友情链接申请了,今天刚看,已经添加。..博主 在
佛跳墙vpn软件不会用?上不了网?佛跳墙vpn常见问题以及解决办法中评论 @1111老铁这个不行了,可以看看近期评论的其他文章..1111 在
佛跳墙vpn软件不会用?上不了网?佛跳墙vpn常见问题以及解决办法中评论 网站不能打开,博主百忙中能否发个APP下载链接,佛跳墙或极光..
Copyright·© 2019 侯体宗版权所有·
粤ICP备20027696号
