【发布时间】:2020-07-30 16:59:22
【问题描述】:
我正在学习如何在 python 中使用 grpc 流在客户端和服务器之间交换消息。我找到了一个基本示例,可以在服务器和客户端之间进行简单的消息发送。我正在尝试对其进行修改,以便我可以跟踪连接到 grpc 服务器(在服务器端)的所有客户端,并且可以做两件事:1)从服务器广播到所有客户端,2)向特定连接的客户端发送消息客户。
这是.proto 文件
syntax = 'proto3';
service Scenario {
rpc Chat(stream DPong) returns (stream DPong) {}
}
message DPong {
string name = 1;
}
这里是client.py,它创建一个守护进程来监听传入消息并等待任何传出消息的标准输入
import threading
import grpc
import time
import scenario_pb2_grpc, scenario_pb2
# new changes
msgQueue = queue.Queue()
def run():
channel = grpc.insecure_channel('localhost:50052')
stub = scenario_pb2_grpc.ScenarioStub(channel)
print('client connected')
global queue
def inputStream():
while 1:
msg = input('>>Enter message\n>>')
yield scenario_pb2.DPong(name=msg)
input_stream = stub.Chat(inputStream())
def read_incoming():
while 1:
print('receivedFromServer: {}\n>>'.format(next(input_stream).name))
thread = threading.Thread(target=read_incoming)
thread.daemon = True
thread.start()
while 1:
time.sleep(1)
if __name__ == '__main__':
print('client starting ...')
run()
下面是server.py
import random
import string
import threading
import grpc
import scenario_pb2_grpc
import scenario_pb2
import time
from concurrent import futures
clientList = []
class Scenario(scenario_pb2_grpc.ScenarioServicer):
def Chat(self, request_iterator, context):
clients = []
def stream():
while 1:
time.sleep(1)
msg = input('>>Enter message\n>>')
for i in clientList:
yield msg
output_stream = stream()
def read_incoming():
while 1:
received = next(request_iterator).name
if (context,request_iterator) not in clientList:
clientList.append((context, request_iterator))
print('receivedFromClient: {}'.format(received), len(clientList))
thread = threading.Thread(target=read_incoming)
thread.daemon = True
thread.start()
while 1:
msg = output_stream
yield scenario_pb2.DPong(name=next(msg))
if __name__ == '__main__':
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
scenario_pb2_grpc.add_ScenarioServicer_to_server(
Scenario(), server)
server.add_insecure_port('[::]:50052')
server.start()
print('listening ...')
while 1:
time.sleep(1)
到目前为止,我已经尝试维护一个列表对象clientList,其中包含客户端的context 和request_iterator 对象,并且每次新客户端加入服务器时都会更新。但是在发送传出消息之前,如何从clientList 设置这些对象?我尝试迭代列表,但服务器多次将消息发送到同一个客户端(最后一个客户端听到),而不是一次将其发送给所有客户端。
非常感谢任何帮助!
【问题讨论】:
标签: client-server grpc broadcast bidirectional grpc-python