【问题标题】:Broadcasting message from grpc server to all/some connected clients in python将消息从 grpc 服务器广播到 python 中的所有/一些连接的客户端
【发布时间】: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,其中包含客户端的contextrequest_iterator 对象,并且每次新客户端加入服务器时都会更新。但是在发送传出消息之前,如何从clientList 设置这些对象?我尝试迭代列表,但服务器多次将消息发送到同一个客户端(最后一个客户端听到),而不是一次将其发送给所有客户端。

非常感谢任何帮助!

【问题讨论】:

    标签: client-server grpc broadcast bidirectional grpc-python


    【解决方案1】:

    这当然是可能的。您在这里遇到的问题是服务器端对Scenario.Chat 的每次调用都对应于一个客户端连接。也就是说,这个函数在流式 RPC 启动时被调用,一旦函数退出,RPC 就结束。

    因此,如果您想要n 连接的客户端,您将需要nScenario.Chat 实例同时运行,每个实例都在自己的线程上。这确实意味着并发连接的客户端的数量受到您实例化服务器的线程池大小的限制。

    因此,假设您的服务器进程中有n 线程专门用于维护客户端连接。然后你需要另一个n+1th 线程(可能是主线程)来确定服务器何时向所有客户端广播消息(可能是通过查找来自 STDIN 的输入?)。当这个额外的线程确定应该广播一条消息时,它需要将此意图传达给所有维护与客户端连接的线程。有很多方法可以实现这一点。 threading.Condition 和全局 collections.deque,或每个客户端连接的 collections.deque(有点像 goroutine 之间的通道)将是两种方式。这里的棘手之处在于确保每个客户端连接都将接收到消息,而不管客户端连接线程需要多长时间才能唤醒以及n+1th 线程决定在此期间发送多少消息。

    如果这仍然不清楚,我可以跟进一些实际代码来展示这个想法。

    【讨论】:

    • 感谢@Richard Belleville 的解释!我能够使用单独的队列构建一个,但我很想了解他们为什么以这种方式构建 grpc。因为这样,你总是需要一些第三方同步机制。为什么不简单地允许用户查询所有上下文或简单地从工作线程中获取/设置上下文。
    • 有几个原因: 1. 可以使用 Python 标准库结构在数十行代码中构建这种同步。 2)这不是最常见的用例。通常,服务线程彼此独立运行,并不关心其他客户端连接的是什么。
    • @RichardBelleville 你能提供代码吗?对第二个选项(每个客户端连接的双端队列)特别感兴趣 - 目前尚不清楚如何实现这一点。在这种情况下应该使用拦截器吗? P.S.:对于第一种方法,是否可以用 threading.Queue 更改 threading.Condition + collections.deque?
    • This example code 非常相似。这是一个内存键值存储。不是向所有客户端广播消息,而是在对密钥 K 进行更新时向观看密钥 K 的客户端广播消息。这是一个更具体的问题,但同步方面应该是相同的。
    • @RichardBelleville 谢谢,您的实施对我帮助很大。我最近遇到了一个类似的问题:我需要在主线程中让一个消费者从 Kafka 获取消息(我不得不放弃 wait_for_termination),并以某种方式将这些消息发送到通过 grpc-streams 连接到处理程序的特定客户端。
    【解决方案2】:

    您可以在一个应用程序中启动多个端口。

    gRPC 可以在端口50011 中运行,带有socket.io 的flask 可以在端口8080 中运行

    使用python,您可以在server.py 中使用flask 框架和flask_socketio

    例如server.py

    from flask import Flask
    from flask_socketio import SocketIO, emit
    
    app = Flask(__name__)
    socketio = SocketIO(app)
    
    @app.route('/')
    def index():
        return "Hello, World!"
    
    if __name__ == '__main__':
        app.run(port=8080)
        app.run(debug=True)
        socketio.run(app)
    

    不使用 gRPC 流 API,而是使用 WebSocket 向所有连接的客户端和使用房间的特定/选定客户端进行广播。

    例如

    @socketio.on('message')
    def handle_message(data):
        // logic to send large data in chunks the logic should call the 
        // emit function in socket.io and emit an event that send the large
        // data in chunks eg emit('my response', chunkData)
    

    gRPC 主要是为一个客户端请求和响应而构建的,而 WebSocket 是为多个客户端构建的。

    【讨论】:

      猜你喜欢
      • 2019-01-06
      • 1970-01-01
      • 2016-02-24
      • 2023-04-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-07-15
      相关资源
      最近更新 更多