You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

请求基于gRPC实现多Client B与Server、Client A的流式通信方案及示例

完全可以用gRPC实现你想要的架构!gRPC对流式通信的支持非常友好,刚好能满足你这种服务器主动向Client A推送Client B列表的场景。下面我给你一套基础的实现示例,涵盖proto定义、服务器端和客户端代码,你可以直接参考修改。

一、先定义gRPC的Proto文件

这是gRPC通信的核心,用来约定服务接口和数据结构。我们需要定义Client B的信息格式、注册接口,以及给Client A的流式推送接口:

syntax = "proto3";

package client_manager;

// 单个Client B的信息结构
message ClientBInfo {
  int32 client_id = 1;  // 你提到的字典中的整数ID
  string ip_address = 2; // Client B的IP地址
}

// Client A请求订阅列表的空请求(可扩展过滤条件)
message ClientListRequest {}

// 服务器返回的Client B列表响应
message ClientListResponse {
  repeated ClientBInfo clients = 1; // 批量Client B信息
}

// Client B注册自身的请求
message ClientRegisterRequest {
  int32 client_id = 1;
  string ip_address = 2;
}

// 注册结果响应
message ClientRegisterResponse {
  bool success = 1;
  string message = 2;
}

// 核心服务定义
service ClientManager {
  // Client A订阅列表的流式推送(服务器主动推送更新)
  rpc SubscribeClientList(ClientListRequest) returns (stream ClientListResponse) {};
  
  // Client B向服务器注册自身
  rpc RegisterClient(ClientRegisterRequest) returns (ClientRegisterResponse) {};
}
二、编译Proto文件

在终端执行以下命令(需要先安装grpc-tools),生成对应语言的代码(这里以Python为例):

python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. client_manager.proto
三、服务器端实现(Python)

服务器需要维护Client B的字典,同时给所有订阅的Client A推送最新列表:

import grpc
from concurrent import futures
import time
import client_manager_pb2
import client_manager_pb2_grpc

# 存储所有在线Client B的信息:key=client_id,value=ip_address
client_b_store = {}
# 存储所有订阅列表的Client A的流式连接
subscribers = []

class ClientManagerServicer(client_manager_pb2_grpc.ClientManagerServicer):
    def RegisterClient(self, request, context):
        # 注册Client B并更新存储
        client_b_store[request.client_id] = request.ip_address
        print(f"Client B [{request.client_id}] 已注册,IP: {request.ip_address}")
        
        # 向所有订阅的Client A推送最新列表
        self._broadcast_client_list()
        
        return client_manager_pb2.ClientRegisterResponse(
            success=True,
            message="注册成功"
        )
    
    def SubscribeClientList(self, request, context):
        # 将当前Client A的连接加入订阅列表
        subscribers.append(context)
        print("新的Client A已订阅列表更新")
        
        # 先推送一次当前的全量列表
        yield self._build_client_list_response()
        
        # 保持连接,等待后续更新(实际可添加退出逻辑)
        try:
            while True:
                time.sleep(1)
        except grpc.RpcError:
            # Client A断开连接,移除订阅
            subscribers.remove(context)
            print("Client A已取消订阅")
    
    def _build_client_list_response(self):
        # 构造Client B列表的响应消息
        client_list = []
        for client_id, ip in client_b_store.items():
            client_list.append(client_manager_pb2.ClientBInfo(
                client_id=client_id,
                ip_address=ip
            ))
        return client_manager_pb2.ClientListResponse(clients=client_list)
    
    def _broadcast_client_list(self):
        # 向所有订阅的Client A广播最新列表
        latest_response = self._build_client_list_response()
        # 遍历副本避免修改列表时出错
        for subscriber in subscribers[:]:
            try:
                subscriber.send(latest_response)
            except grpc.RpcError:
                # 客户端已断开,移除无效连接
                subscribers.remove(subscriber)

def start_server():
    server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
    client_manager_pb2_grpc.add_ClientManagerServicer_to_server(
        ClientManagerServicer(), server)
    server.add_insecure_port('[::]:50051')
    server.start()
    print("服务器已启动,监听端口50051")
    try:
        while True:
            time.sleep(86400)
    except KeyboardInterrupt:
        server.stop(0)

if __name__ == '__main__':
    start_server()
四、Client B实现(Python)

Client B启动时向服务器注册自身:

import grpc
import client_manager_pb2
import client_manager_pb2_grpc

def run_client_b(client_id, ip_address):
    with grpc.insecure_channel('localhost:50051') as channel:
        stub = client_manager_pb2_grpc.ClientManagerStub(channel)
        response = stub.RegisterClient(
            client_manager_pb2.ClientRegisterRequest(
                client_id=client_id,
                ip_address=ip_address
            )
        )
        print(f"注册结果: {response.message}")

if __name__ == '__main__':
    # 示例:替换为实际的Client B ID和IP(可自动获取本地IP)
    run_client_b(1, "192.168.1.100")
五、Client A实现(Python)

Client A订阅服务器的流式推送,实时获取Client B列表:

import grpc
import client_manager_pb2
import client_manager_pb2_grpc

def run_client_a():
    with grpc.insecure_channel('localhost:50051') as channel:
        stub = client_manager_pb2_grpc.ClientManagerStub(channel)
        print("正在订阅Client B列表更新...")
        # 监听流式响应
        for response in stub.SubscribeClientList(client_manager_pb2.ClientListRequest()):
            print("\n当前Client B列表:")
            for client in response.clients:
                print(f"ID: {client.client_id} | IP: {client.ip_address}")

if __name__ == '__main__':
    run_client_a()
额外说明
  • 这个示例用Python实现,gRPC支持Go、Java、C#等多种语言,核心逻辑完全通用。
  • 生产环境中建议添加:
    • TLS加密认证,避免未授权客户端连接
    • Client B心跳检测,自动移除下线的客户端
    • 优化推送逻辑(仅在列表变化时推送,而非定时)
    • 完善的错误处理和重试机制

内容的提问来源于stack exchange,提问作者mdlejtecole

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 07:24:00