[Python] RPC实现
2021-04-27 14:29
标签:处理 进程同步 bin ada display listen 查找 client print 单线程同步 消息协议 客户端 client.py 服务端 blocking_single.py 多线程同步 服务端 multithread.py
多进程同步 服务端 multiprocess.py PreForking同步 单进程异步 PreForking异步 参考 Python多线程和多进程 https://www.cnblogs.com/yssjun/p/11302500.html [Python] RPC实现 标签:处理 进程同步 bin ada display listen 查找 client print 原文地址:https://www.cnblogs.com/cxc1357/p/13197183.html
1 // 输入
2 {
3 in: "ping",
4 params: "ireader 0"
5 }
6
7 // 输出
8 {
9 out: "pong",
10 result: "ireader 0"
11 }
1 # coding: utf-8
2 # client.py
3
4 import json
5 import time
6 import struct
7 import socket
8
9
10 def rpc(sock, in_, params):
11 response = json.dumps({"in": in_, "params": params}) # 请求消息体
12 length_prefix = struct.pack("I", len(response)) # 请求长度前缀
13 sock.sendall(length_prefix)
14 sock.sendall(response)
15 length_prefix = sock.recv(4) # 响应长度前缀
16 length, = struct.unpack("I", length_prefix)
17 body = sock.recv(length) # 响应消息体
18 response = json.loads(body)
19 return response["out"], response["result"] # 返回响应类型和结果
20
21 if __name__ == ‘__main__‘:
22 s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
23 s.connect(("localhost", 8080))
24 for i in range(10): # 连续发送10个rpc请求
25 out, result = rpc(s, "ping", "ireader %d" % i)
26 print out, result
27 time.sleep(1) # 休眠1s,便于观察
28 s.close() # 关闭连接
1 # coding: utf8
2 # blocking_single.py
3
4 import json
5 import struct
6 import socket
7
8
9 def handle_conn(conn, addr, handlers):
10 print addr, "comes"
11 while True: # 循环读写
12 length_prefix = conn.recv(4) # 请求长度前缀
13 if not length_prefix: # 连接关闭了
14 print addr, "bye"
15 conn.close()
16 break # 退出循环,处理下一个连接
17 length, = struct.unpack("I", length_prefix)
18 body = conn.recv(length) # 请求消息体
19 request = json.loads(body)
20 in_ = request[‘in‘]
21 params = request[‘params‘]
22 print in_, params
23 handler = handlers[in_] # 查找请求处理器
24 handler(conn, params) # 处理请求
25
26
27 def loop(sock, handlers):
28 while True:
29 conn, addr = sock.accept() # 接收连接
30 handle_conn(conn, addr, handlers) # 处理连接
31
32
33 def ping(conn, params):
34 send_result(conn, "pong", params)
35
36
37 def send_result(conn, out, result):
38 response = json.dumps({"out": out, "result": result}) # 响应消息体
39 length_prefix = struct.pack("I", len(response)) # 响应长度前缀
40 conn.sendall(length_prefix)
41 conn.sendall(response)
42
43
44 if __name__ == ‘__main__‘:
45 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) # 创建一个TCP套接字
46 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) # 打开reuse addr选项
47 sock.bind(("localhost", 8080)) # 绑定端口
48 sock.listen(1) # 监听客户端连接
49 handlers = { # 注册请求处理器
50 "ping": ping
51 }
52 loop(sock, handlers) # 进入服务循环