English | 简体中文 | 繁體中文 | Русский язык | Français | Español | Português | Deutsch | 日本語 | 한국어 | Italiano | بالعربية

Python으로自定义 주-기 병합 아키텍처 예제 분석

本文实例讲述了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)  #업무 노드가 네트워크 메시지를 전송하는 소켓
#업무: 네트워크 정보 전송
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()

파이썬 관련 내용에 더 관심이 있는 독자는 다음 주제를 확인할 수 있습니다: 《파이썬 URL 작업 기술 요약》、《파이썬 이미지 작업 기술 요약》、《파이썬 데이터 구조와 알고리즘 가이드》、《파이썬 소켓 프로그래밍 기술 요약》、《파이썬 함수 사용 기술 요약》、《파이썬 문자열 작업 요약》、《파이썬 입문 및 고급 가이드》 및 《파이썬 파일 및 디렉토리 작업 기술 요약》

이 문서에서 설명한 내용이 여러분의 파이썬 프로그래밍에 도움이 되길 바랍니다.

성명: 본문은 인터넷에서 가져왔으며, 저작권은 원작자에게 있으며, 인터넷 사용자가 자발적으로 기여하고 업로드한 내용입니다. 웹사이트는 소유권을 가지지 않으며, 인공적인 편집 처리를 하지 않았으며, 관련 법적 책임도 부담하지 않습니다. 저작권 문제가 있으면, notice#w로 이메일을 보내 주세요.3codebox.com에 신고를 보내는 경우, #을 @으로 변경하고 관련 증거를 제공해 주세요. 사실이 확인되면, 사이트는 즉시 저작권 침해 내용을 삭제합니다.

Elasticsearch 가이드