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

Python gRPC客户端双向流中,等价于Go的stream.Send()的方法是什么?

Python gRPC双向流客户端实现同流回发请求的方法

在Python的gRPC中,双向流的客户端实现可以通过调用RouteChat返回的流对象的send()方法,在同一条流里主动向服务器发送消息,无需依赖一次性生成器传参。

改造思路

原代码通过生成器一次性发送所有请求后被动接收响应,要实现动态回发需调整为:

  1. 调用stub.RouteChat()时传入空生成器(如iter([])),获取双向流对象。
  2. 单独启动线程监听服务器响应,避免阻塞发送操作。
  3. 主线程通过流对象的send()方法,随时发送新的请求消息。

改造后的示例代码

import threading
import time
from your_proto_module import RouteNote, make_route_note  # 替换为你的proto实际模块

def receive_responses(stream):
    """监听服务器响应的线程函数"""
    try:
        for response in stream:
            print(f"Received message '{response.message}' at ({response.location.latitude}, {response.location.longitude})")
    except Exception as e:
        print(f"接收响应出错: {e}")

def guide_route_chat(stub):
    # 初始化双向流,传入空生成器
    stream = stub.RouteChat(iter([]))
    
    # 启动线程异步监听响应
    threading.Thread(target=receive_responses, args=(stream,), daemon=True).start()
    
    # 发送初始消息
    initial_messages = [
        make_route_note("First message", 0, 0),
        make_route_note("Second message", 0, 1),
    ]
    for msg in initial_messages:
        print(f"Sending '{msg.message}' at ({msg.location.latitude}, {msg.location.longitude})")
        stream.send(msg)
        time.sleep(1)
    
    # 模拟业务场景下的动态回发(比如收到响应后、用户输入触发等)
    time.sleep(3)
    new_msg = make_route_note("Dynamic reply message", 1, 1)
    print(f"Sending dynamic reply: '{new_msg.message}' at ({new_msg.location.latitude}, {new_msg.location.longitude})")
    stream.send(new_msg)
    
    # 根据业务需求决定何时关闭流
    time.sleep(2)
    stream.close()

关键说明

  • 双向流对象为grpc._channel._MultiThreadedRendezvous类型,自带send()方法用于发送消息,close()方法用于关闭客户端发送流。
  • 必须用独立线程处理响应接收,因为for response in stream是阻塞操作,主线程执行会无法同时发送消息。
  • 发送消息前需确保流未关闭,否则会触发异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 17:02:56