이 글은 주로 파이썬의 커스텀 마스터-슬레이브 분산 아키텍처를 소개하고, 마스터-슬레이브 분산 아키텍처의 구조와 원리, 구체적인 코드 구현 기법을 예시 형태로 분석해 도움이 필요한 친구들이 참고할 수 있습니다
이 문서의 예에서는 Python의 사용자 정의 마스터-슬레이브 분산 아키텍처를 설명합니다. 참고용으로 모든 사람과 공유하세요. 세부 사항은 다음과 같습니다:
환경: Win7 x64, Python 2.7, APScheduler 2.1.2.
구성도는 다음과 같습니다.
코드 부분:
( 1), 센터 노드:
#encoding=utf-8 #author: walker #date: 2014-12-03 #function: 中心节点(主要功能是分配任务) import SocketServer, socket, Queue CenterIP = '127.0.0.1' #中心节点IP CenterListenPort = 9999 #中心节点监听端口 CenterClient = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) #中心节点用于发送网络消息的socket TaskQueue = 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, SocketServer from apscheduler.scheduler import Scheduler CenterIP = '127.0.0.1' #中心节点IP CenterListenPort = 9999 #中心节点监听端口 NodeIP = socket.gethostbyname(socket.gethostname()) #任务节点自身IP NodeClient = 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 사용자 정의 마스터-슬레이브 분산 아키텍처 예제 분석과 관련된 더 많은 기사를 보려면 PHP 중국어 웹사이트를 주목하세요!