python调用zookeeper

ZooKeeper

1. 简介

ZooKeeper是一种分布式协调服务,用于管理大型主机。在分布式环境中协调和管理服务是一个复杂的过程。

ZooKeeper通过其简单的架构和API解决了这个问题。ZooKeeper允许开发人员专注于核心应用程序逻辑,而不必担心应用程序的分布式特性。

ZooKeeper框架最初是在“Yahoo!"上构建的,用于以简单而稳健的方式访问他们的应用程序。 后来,Apache ZooKeeper成为Hadoop,HBase和其他分布式框架使用的有组织服务的标准。 例如,Apache HBase使用ZooKeeper跟踪分布式数据的状态。

avatar

2. 概念知识

层次命名空间

下图描述了用于内存表示的ZooKeeper文件系统的树结构(ZooKeeper的数据保存形式)。ZooKeeper节点称为 znode 。每个znode由一个名称标识,并用路径(/)序列分隔。

每个znode最多可存储1MB的数据。

avatar

Znode的类型

Znode被分为持久(persistent)节点,顺序(sequential)节点和临时(ephemeral)节点。

  • 持久节点 - 即使在创建该特定znode的客户端断开连接后,持久节点仍然存在。默认情况下,除非另有说明,否则所有znode都是持久的。
  • 临时节点 - 客户端活跃时,临时节点就是有效的。当客户端与ZooKeeper集合断开连接时,临时节点会自动删除。因此,只有临时节点不允许有子节点。如果临时节点被删除,则下一个合适的节点将填充其位置。临时节点在leader选举中起着重要作用。
  • 顺序节点 - 顺序节点可以是持久的或临时的。当一个新的znode被创建为一个顺序节点时,ZooKeeper通过将10位的序列号附加到原始名称来设置znode的路径。例如,如果将具有路径 /myapp 的znode创建为顺序节点,则ZooKeeper会将路径更改为 /myapp0000000001 ,并将下一个序列号设置为0000000002。如果两个顺序节点是同时创建的,那么ZooKeeper不会对每个znode使用相同的数字。顺序节点在锁定和同步中起重要作用。

Watches(监视)

监视是一种简单的机制,使客户端收到关于ZooKeeper集合中的更改的通知。客户端可以在读取特定znode时设置Watches。Watches会向注册的客户端发送任何znode(客户端注册表)更改的通知。

Znode更改是与znode相关的数据的修改或znode的子项中的更改。只触发一次watches。如果客户端想要再次通知,则必须通过另一个读取操作来完成。当连接会话过期时,客户端将与服务器断开连接,相关的watches也将被删除。

ZooKeeper安装

在安装ZooKeeper之前,请确保你的系统是在以下任一操作系统上运行:

1任意Linux OS - 支持开发和部署。适合演示应用程序。 2 3Windows OS - 仅支持开发。 4 5Mac OS - 仅支持开发。

ZooKeeper服务器是用Java创建的,它在JVM上运行。你需要使用JDK 6或更高版本。

现在,按照以下步骤在你的机器上安装ZooKeeper框架。

步骤1:验证Java安装

相信你已经在系统上安装了Java环境。现在只需使用以下命令验证它。

$ java -version

如果你在机器上安装了Java,那么可以看到已安装的Java的版本。否则,请按照以下简单步骤安装最新版本的Java。

步骤1.1:下载JDK

通过访问链接下载最新版本的JDK,并下载最新版本的Java

步骤1.2:提取文件

通常,文件会下载到download文件夹中。验证并使用以下命令提取tar设置。

1$ cd /path/to/download/ 2$ tar -zxvf jdk-8u181-linux-x64.gz

步骤1.3:移动到/usr/local/jdk目录

要使Java对所有用户可用,请将提取的Java内容移动到“/usr/local/jdk"文件夹。

1$ sudo mkdir /usr/local/jdk 2$ sudo mv jdk1.8.0_181 /usr/local/jdk

步骤1.4:设置路径

要设置路径和JAVA_HOME变量,请将以下命令添加到〜/.bashrc文件中。

1export JAVA_HOME=/usr/local/jdk/jdk1.8.0_181 2export PATH=$PATH:$JAVA_HOME/bin

现在,将所有更改应用到当前运行的系统中。

$ source ~/.bashrc

步骤1.5

使用步骤1中说明的验证命令(java -version)验证Java安装。

步骤2:ZooKeeper框架安装

步骤2.1:下载ZooKeeper

要在你的计算机上安装ZooKeeper框架,请访问以下链接并下载最新版本的ZooKeeper。

http://zookeeper.apache.org/releases.html

到目前为止,最新版本的ZooKeeper是3.4.12(ZooKeeper-3.4.12.tar.gz)。

步骤2.2:提取tar文件

使用以下命令提取tar文件

1$ cd /path/to/download/ 2$ tar -zxvf zookeeper-3.4.12.tar.gz 3$ cd zookeeper-3.4.12 4$ mkdir data

步骤2.3:创建配置文件

使用命令 vi conf/zoo.cfg 和所有以下参数设置为起点,打开名为 conf/zoo.cfg 的配置文件。

1$ vi conf/zoo.cfg 2 3tickTime = 2000 4dataDir = /path/to/zookeeper/data 5clientPort = 2181

一旦成功保存配置文件,再次返回终端。你现在可以启动zookeeper服务器。

步骤2.4:启动ZooKeeper服务器

执行以下命令

$ bin/zkServer.sh start

执行此命令后,你将收到以下响应

1$ JMX enabled by default 2$ Using config: /Users/../zookeeper-3.4.12/bin/../conf/zoo.cfg 3$ Starting zookeeper ... STARTED

步骤2.5:启动CLI

键入以下命令

$ bin/zkCli.sh

键入上述命令后,将连接到ZooKeeper服务器,你应该得到以下响应。

1Connecting to localhost:2181 2................ 3................ 4................ 5Welcome to ZooKeeper! 6................ 7................ 8WATCHER:: 9WatchedEvent state:SyncConnected type: None path:null 10[zk: localhost:2181(CONNECTED) 0]

停止ZooKeeper服务器

连接服务器并执行所有操作后,可以使用以下命令停止zookeeper服务器。

1$ bin/zkServer.sh stop 2

Kazoo

kazoo是Python连接操作ZooKeeper的客户端库。我们可以通过kazoo来使用ZooKeeper。

1. 安装

pip install kazoo

2. 使用

连接ZooKeeper

1from kazoo.client import KazooClient 2 3zk = KazooClient(hosts='127.0.0.1:2181') 4 5# 启动连接 6zk.start() 7 8# 停止连接 9zk.stop()

创建节点

1# 创建节点路径,但不能设置节点数据值 2zk.ensure_path("/my/favorite") 3 4# 创建节点,并设置节点保存数据,ephemeral表示是否是临时节点,sequence表示是否是顺序节点 5zk.create("/my/favorite/node", b"a value", ephemeral=True, sequence=True)

读取节点

1# 获取子节点列表 2children = zk.get_children("/my/favorite") 3 4# 获取节点数据data 和节点状态stat 5data, stat = zk.get("/my/favorite")

设置监视

1def my_func(event): 2 # 检查最新的节点数据 3 4# 当子节点发生变化的时候,调用my_func 5children = zk.get_children("/my/favorite/node", watch=my_func)

server端

1import threading 2from kazoo.client import KazooClient 3 4class ThreadServer(object): 5 def __init__(self, host, port, handlers): 6 self.sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 7 self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) 8 self.host = host 9 self.port = port 10 self.sock.bind((host, port)) 11 self.handlers = handlers 12 13 def serve(self): 14 """ 15 开始服务 16 """ 17 self.sock.listen(128) 18 self.register_zk() 19 print("开始监听") 20 while True: 21 conn, addr = self.sock.accept() 22 print("建立链接%s" % str(addr)) 23 t = threading.Thread(target=self.handle, args=(conn,)) 24 t.start() 25 26 def handle(self, client): 27 stub = ServerStub(client, self.handlers) 28 try: 29 while True: 30 stub.process() 31 except EOFError: 32 print("客户端关闭连接") 33 34 client.close() 35 36 def register_zk(self): 37 """ 38 注册到zookeeper 39 """ 40 self.zk = KazooClient(hosts='127.0.0.1:2181') 41 self.zk.start() 42 self.zk.ensure_path('/rpc') # 创建根节点 43 value = json.dumps({'host': self.host, 'port': self.port) 44 # 创建服务子节点 45 self.zk.create('/rpc/server', value.encode(), ephemeral=True, sequence=True) 46

client端

1from services import ThreadServer 2from services import InvalidOperation 3import sys 4 5 6class Handlers: 7 @staticmethod 8 def divide(num1, num2=1): 9 """ 10 除法 11 :param num1: 12 :param num2: 13 :return: 14 """ 15 if num2 == 0: 16 raise InvalidOperation() 17 val = num1 / num2 18 return val 19 20 21if __name__ == '__main__': 22 if len(sys.argv) < 3: 23 print("usage:python server.py [host] [port]") 24 exit(1) 25 host = sys.argv[1] 26 port = sys.argv[2] 27 server = ThreadServer(host, int(port), Handlers) 28 server.serve()

server改写

1import random 2import time 3 4class DistributedChannel(object): 5 def __init__(self): 6 self._zk = KazooClient(hosts='127.0.0.1:2181') 7 self._zk.start() 8 self._get_servers() 9 10 def _get_servers(self, event=None): 11 """ 12 从zookeeper获取服务器地址信息列表 13 """ 14 servers = self._zk.get_children('/rpc', watch=self._get_servers) 15 print(servers) 16 self._servers = [] 17 for server in servers: 18 data = self._zk.get('/rpc/' + server)[0] 19 addr = json.loads(data) 20 self._servers.append(addr) 21 22 def _get_server(self): 23 """ 24 随机选出一个可用的服务器 25 """ 26 return random.choice(self._servers) 27 28 def get_connection(self): 29 """ 30 提供一个可用的tcp连接 31 """ 32 while True: 33 server = self._get_server() 34 print(server) 35 try: 36 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 37 sock.connect((server['host'], server['port'])) 38 except ConnectionRefusedError: 39 time.sleep(1) 40 continue 41 else: 42 break 43 return sock

client改写

1from services import ClientStub 2from services import DistributedChannel 3from services import InvalidOperation 4import time 5 6 7channel = DistributedChannel() 8 9for i in range(50): 10 try: 11 stub = ClientStub(channel) 12 val = stub.divide(i) 13 except InvalidOperation as e: 14 print(e.message) 15 else: 16 print(val) 17 time.sleep(1)

测试完整代码

1import json 2from kazoo.client import KazooClient 3 4zk = KazooClient(hosts="192.168.218.136:2181") 5 6# 启动连接 7zk.start() 8# 创建节点路径,但不能设置节点数据值 9zk.ensure_path("/rpc") 10addr1 = {"host": "127.0.0.1", "port": 8001} 11addr1_str = json.dumps(addr1) 12# 创建节点,并设置节点保存数据,ephemeral表示是否是临时节点,sequence表示是否是顺序节点 13zk.create("/rpc/server", addr1_str.encode(), ephemeral=True, sequence=True) 14 15addr2 = {"host": "127.0.0.1", "port": 8002} 16addr2_str = json.dumps(addr2) 17zk.create("/rpc/server", addr2_str.encode(), ephemeral=True, sequence=True) 18 19 20# 获取子节点列表 21children = zk.get_children("/rpc") 22 23# 获取节点数据data 和节点状态stat 24for i in zk.get_children("/rpc"): 25 print(zk.get("/rpc/"+i)[0]) 26 27 28def on_change(event): 29 print(event) 30 31 32# 设置监视, 该监视只触发一次 33servers = zk.get_children("/rpc", watch=on_change) 34 35# 停止连接 36zk.stop()

service.py

1#! /usr/bin/env python 2# -*- coding: utf-8 -*- 3# Author: Wjy 4import struct 5from io import BytesIO 6import socket 7import threading 8import json 9import random 10import time 11from kazoo.client import KazooClient 12 13 14class InvalidOperation(Exception): 15 16 def __init__(self, message=None): 17 self.message = message or "invalid operation" 18 19 20class MethodProtocol(object): 21 """ 22 解读方法名 23 """ 24 def __init__(self, connection): 25 self.conn = connection 26 27 def _read_all(self, size): 28 """ 29 帮助我们读取二进制数据 30 :param size: 想要读取的二进制数据大小 31 :return: 二进制数据 bytes 32 """ 33 # self.conn 34 # 读取二进制数据 35 # socket.recv(4) => ?4 36 # BytesIO.read 37 if isinstance(self.conn, BytesIO): 38 buff = self.conn.read(size) 39 return buff 40 else: 41 # socket 42 have = 0 43 buff = b"" 44 while have < size: 45 chunk = self.conn.recv(size - have) 46 buff += chunk 47 l = len(chunk) 48 have += l 49 50 if l == 0: 51 # 表示客户端socket关闭了 52 raise EOFError() 53 return buff 54 55 def get_method_name(self): 56 """ 57 提供方法名 58 :return: str 方法名 59 """ 60 # 读取字符串长度 61 buff = self._read_all(4) 62 length = struct.unpack("!I", buff)[0] 63 64 # 读取字符串 65 buff = self._read_all(length) 66 name = buff.decode() 67 return name 68 69 70class DivideProtocol(object): 71 """ 72 divide过程消息协议转换工具 73 """ 74 75 def args_encode(self, num1, num2=1): 76 """ 77 将原始的调用请求参数转换打包成二进制消息数据 78 :param num1: int 79 :param num2: int 80 :return: bytes 二进制消息shuju 81 """ 82 name = "divide" 83 84 # 处理方法的名字 字符串 85 # 处理字符串的长度 86 buff = struct.pack("!I", 6) 87 # 处理字符 88 buff += name.encode() 89 90 # 处理参数1 91 # 处理序号 92 buff2 = struct.pack("!B", 1) 93 # 处理参数值 94 buff2 += struct.pack("!i", num1) 95 96 # 处理参数2 97 if num2 != 1: 98 # 处理序号 99 buff2 += struct.pack("!B", 2) 100 # 处理参数值 101 buff2 += struct.pack("!i", num2) 102 103 # 处理消息长度,边界固定 104 length = len(buff2) 105 buff += struct.pack("!I", length) 106 107 buff += buff2 108 109 return buff 110 111 def _read_all(self, size): 112 """ 113 帮助我们读取二进制数据 114 :param size: 想要读取的二进制数据大小 115 :return: 二进制数据 bytes 116 """ 117 # self.conn 118 # 读取二进制数据 119 # socket.recv(4) => ?4 120 # BytesIO.read 121 if isinstance(self.conn, BytesIO): 122 buff = self.conn.read(size) 123 return buff 124 else: 125 # socket 126 have = 0 127 buff = b"" 128 while have < size: 129 chunk = self.conn.recv(size - have) 130 buff += chunk 131 l = len(chunk) 132 have += l 133 134 if l == 0: 135 # 表示客户端socket关闭了 136 raise EOFError() 137 return buff 138 139 def args_decode(self, connection): 140 """ 141 接受调用请求消息数据并进行解析 142 :param connection: 连接对象 socket BytesIO 143 :return: dict 包含了解析之后的参数 144 """ 145 param_len_map = { 146 1: 4, 147 2: 4, 148 } 149 param_fmt_map = { 150 1: "!i", 151 2: "!i", 152 } 153 param_name_map = { 154 1: "num1", 155 2: "num2" 156 } 157 158 # 保存用来返回的参数 159 # args = {"num1": xxx, "num2": xxx} 160 args = { 161 162 } 163 164 self.conn = connection 165 # 处理方法的名已经提前被处理(稍后实现) 166 167 # 处理消息边界 168 # 读取二进制数据 169 # socket.recv(4) => ?4 170 # BytesIO.read 171 buff = self._read_all(4) 172 # 将二进制数据转换为python数据类型 173 length = struct.unpack("!I", buff)[0] 174 175 # 已经读取处理的字节数 176 have = 0 177 178 # 处理第一个参数 179 # 处理参数序号 180 buff = self._read_all(1) 181 have += 1 182 param_seg = struct.unpack("!B", buff)[0] 183 184 # 处理参数值 185 param_len = param_len_map[param_seg] 186 buff = self._read_all(param_len) 187 have += param_len 188 param_fmt = param_fmt_map[param_seg] 189 param = struct.unpack(param_fmt, buff)[0] 190 191 param_name = param_name_map[param_seg] 192 args[param_name] = param 193 194 if have >= length: 195 return args 196 197 # 处理第二个参数 198 # 处理参数序号 199 buff = self._read_all(1) 200 param_seg = struct.unpack("!B", buff)[0] 201 202 # 处理参数值 203 param_len = param_len_map[param_seg] 204 buff = self._read_all(param_len) 205 param_fmt = param_fmt_map[param_seg] 206 param = struct.unpack(param_fmt, buff)[0] 207 208 param_name = param_name_map[param_seg] 209 args[param_name] = param 210 211 return args 212 213 def result_encode(self, result): 214 """ 215 将原始结果数据转换为消息协议二进制数据 216 :param result: 原始结果数据 float InvalidOperation 217 :return: bytes 消息协议二进制数据 218 """ 219 # 正常 220 if isinstance(result, float): 221 # 处理返回值类型 222 buff = struct.pack("!B", 1) 223 buff += struct.pack("!f", result) 224 return buff 225 # 异常 226 else: 227 # 处理返回值类型 228 buff = struct.pack("!B", 2) 229 # 处理返回值 230 length = len(result.message) 231 # 处理字符串长度 232 buff += struct.pack("!I", length) 233 # 处理字符 234 buff += result.message.encode() 235 return buff 236 237 def result_decode(self, connection): 238 """ 239 将返回值消息数据转换为原始返回值 240 :param connection: socket BytesIO 241 :return: float InvalidOperation对象 242 """ 243 self.conn = connection 244 245 # 处理返回值类型 246 buff = self._read_all(1) 247 result_type = struct.unpack("!B", buff)[0] 248 249 if result_type == 1: 250 # 正常 251 # 读取float数量 252 buff = self._read_all(4) 253 val = struct.unpack("!f", buff)[0] 254 return val 255 else: 256 # 异常 257 # 读取字符串的长度 258 buff = self._read_all(4) 259 length = struct.unpack("!I", buff)[0] 260 261 # 读取字符串 262 buff = self._read_all(length) 263 message = buff.decode() 264 265 return InvalidOperation(message) 266 267 268class Channel(object): 269 """ 270 用户客户端建立网络连接 271 """ 272 def __init__(self, host, port): 273 """ 274 275 :param host: 服务器地址 276 :param port: 服务器端口号 277 """ 278 self.host = host 279 self.port = port 280 281 def get_connection(self): 282 """ 283 获取连接对象 284 :return: 与服务器通讯的socket 285 """ 286 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 287 sock.connect((self.host, self.port)) 288 return sock 289 290 291class DistributedChannel(object): 292 """ 293 支持分布式的zookeeper的RPC客户端通讯连接工具 294 """ 295 def __init__(self): 296 # 创建kazoo对象,用来跟zookeeper连接,获取信息 297 self.zk = KazooClient("192.168.218.136:2181") 298 self.zk.start() 299 self._servers = [] 300 self._get_servers() 301 302 def _get_servers(self, event=None): 303 """ 304 从zookeeper中获取所有可用的RPC服务器地址信息 305 :return: 306 """ 307 self._servers = [] 308 # 从zookeeper中获取/rpc节点下所有可用的rpc服务器节点 309 servers = self.zk.get_children("/rpc", watch=self._get_servers) 310 # 遍历节点,获取服务器的地址信息 311 for server in servers: 312 addr_data = self.zk.get("/rpc/" + server)[0] 313 addr = json.loads(addr_data) 314 self._servers.append(addr) 315 316 def _get_server(self): 317 """ 318 从可用的服务器列表中选出一台服务器 319 :return: 320 """ 321 return random.choice(self._servers) 322 323 def get_connection(self): 324 """ 325 提供一个具体的与RPC服务器的连接socket 326 :return: 327 """ 328 while True: 329 addr = self._get_server() 330 print(addr) 331 try: 332 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 333 sock.connect((addr["host"], addr["port"])) 334 except ConnectionRefusedError: 335 time.sleep(1) 336 continue 337 else: 338 return sock 339 340 341 342class Server(object): 343 """ 344 RPC服务器 345 """ 346 def __init__(self, host, port, handlers): 347 self.host = host 348 self.port = port 349 # 创建socket的工具对象 350 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 351 352 # 设置socket 353 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) 354 355 # 绑定地址 356 sock.bind((self.host, self.port)) 357 self.sock = sock 358 self.handlers = handlers 359 360 def serve(self): 361 """ 362 开启服务器运行,提供RPC服务 363 :return: 364 """ 365 # 开启服务器的监听,等待客户端的连接请求 366 self.sock.listen(128) 367 print("服务器开始监听") 368 369 # 接收客户端的连接请求 370 while True: 371 client_sock, client_addr = self.sock.accept() 372 print("与客户端%s建立了连接" % str(client_addr)) 373 374 # 交给ServerStub,完成客户端的具体的RPC调用请求 375 stub = ServerStub(client_sock, self.handlers) 376 try: 377 while True: 378 stub.process() 379 except EOFError: 380 # 表示客户端关闭了连接 381 print("客户端关闭了连接") 382 client_sock.close() 383 384 385class ThreadServer(object): 386 """ 387 多线程RPC服务器 388 """ 389 def __init__(self, host, port, handlers): 390 # 创建socket的工具对象 391 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) 392 393 # 设置socket 394 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) 395 396 # 绑定地址 397 sock.bind((host, port)) 398 self.host = host 399 self.port = port 400 self.sock = sock 401 self.handlers = handlers 402 403 def serve(self): 404 """ 405 开启服务器运行,提供RPC服务 406 :return: 407 """ 408 # 开启服务器的监听,等待客户端的连接请求 409 self.sock.listen(128) 410 print("服务器开始监听") 411 412 # 注册到zookeeper 413 self.regoster_zookeeper() 414 415 # 接收客户端的连接请求 416 while True: 417 client_sock, client_addr = self.sock.accept() 418 print("与客户端%s建立了连接" % str(client_addr)) 419 420 # 创建子线程处理这个客户端 421 t = threading.Thread(target=self.handle, args=(client_sock,)) 422 # 开启子线程执行 423 t.start() 424 425 def regoster_zookeeper(self): 426 """ 427 在zookeeper中心中注册本服务器的地址信息 428 :return: 429 """ 430 # 创建kazoo客户端 431 zk = KazooClient("192.168.218.136:2181") 432 # 建立与zookeeper的连接 433 zk.start() 434 # 在zookeeper中创建节点保存数据 435 zk.ensure_path("/rpc") 436 data = json.dumps({"host": self.host, "port": self.port}) 437 # 在zookeeper上存储服务器地址及端口,设置成临时节点,顺序节点 438 zk.create("/rpc/server", data.encode(), ephemeral=True, sequence=True) 439 440 def handle(self, client_sock): 441 """ 442 子线程调用的方法,用来处理一个客户端的请求 443 :return: 444 """ 445 # 交给ServerStub,完成客户端的具体的RPC调用请求 446 stub = ServerStub(client_sock, self.handlers) 447 try: 448 while True: 449 stub.process() 450 except EOFError: 451 # 表示客户端关闭了连接 452 print("客户端关闭了连接") 453 client_sock.close() 454 455 456class ClientStub(object): 457 """ 458 用来帮助客户端完成远程过程调用 RPC调用 459 460 stub = ClientStub() 461 stib.divide(200) 462 """ 463 def __init__(self, channel): 464 self.channel = channel 465 self.conn = self.channel.get_connection() 466 467 def divide(self, num1, num2=1): 468 # 将调用的参数打包成消息协议的数据 469 proto = DivideProtocol() 470 args = proto.args_encode(num1, num2) 471 472 # 将消息数据通过网络发送给服务器 473 self.conn.sendall(args) 474 475 # 接收服务器返回的返回值消息数据,并进行解析 476 result = proto.result_decode(self.conn) 477 478 # 将结果值(正常float 或 异常InvalidOperation)返回给客户端 479 if isinstance(result, float): 480 # 正常 481 return result 482 else: 483 # 异常 484 raise result 485 486 def add(self): 487 pass 488 489 490class ServerStub: 491 """ 492 帮助服务端完成远程过程调用 493 """ 494 def __init__(self, connection, handlers): 495 """ 496 497 :param connection: 与客户端的连接 498 :param handlers: 真正本地被调用方法(函数 过程) 499 class Handlers: 500 501 @staticmethod 502 def divide(num1, num2=1): 503 pass 504 505 def add(): 506 pass 507 """ 508 self.conn = connection 509 self.method_proto = MethodProtocol(self.conn) 510 self.process_map = { 511 "divide": self._process_divide, 512 "add": self._process_add, 513 } 514 self.handlers = handlers 515 516 def process(self): 517 """ 518 当服务端接收了一个客户端的连接,建立好连接后,完成远端调用处理 519 :return: 520 """ 521 # 接收消息数据,并解析方法的名字 522 name = self.method_proto.get_method_name() 523 524 # 根据机械获得的方法(过程)名,调用相应的过程协议,接收并解析消息数据 525 # self.process_map[name]() 526 _process = self.process_map[name] 527 _process() 528 529 def _process_divide(self): 530 """ 531 处理除法过程调用 532 :return: 533 """ 534 # 创建用于除法过程调用参数协议解析的工具 535 proto = DivideProtocol() 536 # 解析调用消息参数 537 args = proto.args_decode(self.conn) 538 # args = {"num1": xxx, "num2": xxx} 539 540 # 进行除法的本地过程调用 541 # 将本地调用过程的返回值(包括可能的异常)打包成消息协议数据,通过网络返回给客户端 542 try: 543 val = self.handlers.divide(**args) 544 except InvalidOperation as e: 545 ret_message = proto.result_encode(e) 546 else: 547 ret_message = proto.result_encode(val) 548 549 self.conn.sendall(ret_message) 550 551 def _process_add(self): 552 pass 553 554 555if __name__ == '__main__': 556 # 狗贼消息数据 557 proto = DivideProtocol() 558 # divide(200, 100) 559 # message = proto.args_encode(200, 100) 560 # divide(200) 561 message = proto.args_encode(200) 562 conn = BytesIO() 563 conn.write(message) 564 conn.seek(0) 565 566 # 解析消息数据 567 method = MethodProtocol(conn) 568 name = method.get_method_name() 569 print(name) 570 571 args = proto.args_decode(conn) 572 print(args)

client.py

1from service import ClientStub 2from service import Channel 3from service import InvalidOperation 4from service import DistributedChannel 5import time 6 7# 创建与服务器的连接 8# channel = Channel("127.0.0.1", 8000) 9channel = DistributedChannel() 10 11# 运行调用 12for i in range(50): 13 try: 14 # 创建用于RPC调用的工具 15 stub = ClientStub(channel) 16 17 val = stub.divide(i * 100, 50) 18 except InvalidOperation as e: 19 print(e.message) 20 except Exception as e: 21 print(e) 22 else: 23 print(val) 24 25 time.sleep(1)

使用.py

1from service import InvalidOperation 2from service import Server, ThreadServer 3import sys 4 5 6class Handlers: 7 8 @staticmethod 9 def divide(num1, num2=1): 10 """ 11 除法 12 :param num1: int 13 :param num2: int 14 :return: 15 """ 16 if num2 == 0: 17 raise InvalidOperation() 18 val = num1 / num2 19 return val 20 21 22if __name__ == '__main__': 23 # 开启服务器 24 # _server = Server("127.0.0.1", 8000, Handlers) 25 # _server.serve() 26 27 # 从启动命令中提取服务器运行的ip地址和端口号 28 host = sys.argv[1] 29 port = sys.argv[2] 30 31 _server = ThreadServer(host, int(port), Handlers) 32 _server.serve() 33
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

Hadoop 2.6.0 HA高可用集群配置详解(二)

Zookeeper集群安装Zookeeper是一个开源分布式协调服务,其独特的LeaderFollower集群结构,很好的解决了分布式单点问题。目前主要用于诸如:统一命名服务、配置管理、锁服务、集群管理等场景。大数据应用中主要使用Zookeeper的集群管理功能。本集群使用zookeeper3.4.5cdh5.7.1版本。首先在Hado

4cast

4castpackageloadcsv.KumarAwanish发布:2020122117:43:04.501348作者:KumarAwanish作者邮箱:awanish00@gmail.com首页:

Twitter的分布式自增ID算法snowflake (Java版)

概述分布式系统中,有一些需要使用全局唯一ID的场景,这种时候为了防止ID冲突可以使用36位的UUID,但是UUID有一些缺点,首先他相对比较长,另外UUID一般是无序的。有些时候我们希望能使用一种简单一些的ID,并且希望ID能够按照时间有序生成。而twitter的snowflake解决了这种需求,最初Twitter把存储系统从MySQL迁移

mysql设置时区

mysql设置时区mysql\_query("SETtime\_zone'8:00'")ordie('时区设置失败,请联系管理员!');中国在东8区所以加8方法二:selectcount(user\_id)asdevice,CONVERT\_TZ(FROM\_UNIXTIME(reg\_time),'08:00','0