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

如何配置tonic::transport::Endpoint实现GRPC客户端断连自动恢复?

gRPC客户端断连后无法自动重连的问题解决

问题描述

使用gRPC客户端(基于tonic框架)时,遇到以下问题:当客户端与服务器连接断开(如设备长时间休眠或物理断连)后,无法自动恢复连接,程序无法继续使用gRPC服务,仅返回超时或断连错误。已通过搭建简单的tonic服务端与客户端复现该问题,服务端部署在远程机器以模拟物理断连场景。

错误日志

物理断连后报错

[2023-01-28T15:11:55Z DEBUG tower::buffer::worker] service.ready=true message=processing request
[2023-01-28T15:11:55Z DEBUG h2::codec::framed_write] send frame=Headers { stream_id: StreamId(19), flags: (0x4: END_HEADERS) }
[2023-01-28T15:11:55Z DEBUG h2::codec::framed_write] send frame=Data { stream_id: StreamId(19) }
[2023-01-28T15:11:55Z DEBUG h2::codec::framed_write] send frame=Data { stream_id: StreamId(19), flags: (0x1: END_STREAM) }
[2023-01-28T15:12:01Z DEBUG hyper::proto::h2::server] stream error: connection error: broken pipe
[2023-01-28T15:12:01Z DEBUG h2::codec::framed_write] send frame=Reset { stream_id: StreamId(19), error_code: CANCEL }

超时报错

[2023-01-28T16:11:55Z DEBUG client] send: interval(40)
[2023-01-28T16:11:55Z DEBUG h2::codec::framed_write] send frame=Reset { stream_id: StreamId(25), error_code: CANCEL }
[2023-01-28T16:11:55Z DEBUG tower::buffer::worker] service.ready=true message=processing request
[2023-01-28T16:11:55Z DEBUG h2::codec::framed_write] send frame=Headers { stream_id: StreamId(27), flags: (0x4: END_HEADERS) }
[2023-01-28T16:11:55Z DEBUG h2::codec::framed_write] send frame=Data { stream_id: StreamId(27) }
[2023-01-28T16:11:55Z DEBUG h2::codec::framed_write] send frame=Data { stream_id: StreamId(27), flags: (0x1: END_STREAM) }
[2023-01-28T16:11:56Z ERROR client] status: Cancelled, message: "Timeout expired", details: [], metadata: MetadataMap { headers: {} }

复现代码

服务端代码(Rust + tonic)

struct Service {}
#[tonic::async_trait]
impl test_grpc::say_server::Say for Service {
    async fn hello(&self, request: Request<RequestSay>) -> Result<Response<ResponseSay>, Status> {
        let r = request.into_inner().text;
        debug!("in request: {}", r);
        Ok(Response::new(ResponseSay {
            text: format!("hello {r}"),
        }))
    }
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error + Send + Sync>> {
    env_logger::Builder::new()
        .filter_level(log::LevelFilter::from_str("debug").unwrap())
        .init();
    let s = Service {};
    let key = "secret token";
    let svc = test_grpc::say_server::SayServer::with_interceptor(
        s,
        move |req: Request<()>| -> Result<Request<()>, Status> {
            let token: MetadataValue<_> = key.parse().unwrap();
            match req.metadata().get("authorization") {
                Some(t) if token == t => Ok(req),
                _ => Err(Status::unauthenticated("No valid auth token")),
            }
        },
    );
    let addr = "0.0.0.0:8804".parse::<SocketAddr>().unwrap();
    Server::builder()
        .add_service(svc)
        .serve(addr)
        .await
        .unwrap();
    Ok(())
}

客户端代码(Rust + tonic)

async fn tester_client(sleep: Duration, uri: &str, key: &str) {
    let uri = uri.parse().unwrap();
    debug!("create connect");
    let chan = tonic::transport::Channel::builder(uri)
        .timeout(Duration::from_secs(20))
        .connect_timeout(Duration::from_secs(20))
        //.http2_keep_alive_interval(Duration::from_secs(5))
        //.keep_alive_while_idle(true)
        .connect_lazy();

    let key = key.parse::<tonic::metadata::MetadataValue<_>>().unwrap();
    let mut key = Some(key);
    let mut service = test_grpc::say_client::SayClient::with_interceptor(
        chan,
        move |mut req: tonic::Request<()>| {
            if let Some(secret) = &mut key {
                req.metadata_mut().insert("authorization", secret.clone());
            }
            Ok(req)
        },
    );
    loop {
        let send_text = format!("interval({})", sleep.as_secs_f32() / 60.0);
        debug!("send: {send_text}");
        let res = match service
            .hello(tonic::Request::new(test_grpc::RequestSay {
                text: send_text.clone(),
            }))
            .await
        {
            Ok(r) => r,
            Err(e) => {
                error!("{e:#}");
                continue;
            }
        };
        debug!("recv: {}", res.into_inner().text);
        time::sleep(sleep).await;
        println!();
    }
}

解决方案

当前客户端使用connect_lazy()创建连接,但默认情况下,tonic的Channel不会自动在连接断开后重建。需要结合以下配置和代码调整实现自动重连:

  1. 启用重试依赖
    在Cargo.toml中添加tonic的retry特性及tower相关依赖:

    tonic = { version = "0.9", features = ["retry", "tls"] }
    tower = { version = "0.4", features = ["retry", "util"] }
    rand = "0.8"
    
  2. 修改客户端代码实现自动重连
    调整为双层循环结构,内层处理正常请求,外层负责重建连接;同时配置keep-alive参数提前检测连接状态:

    use tokio::time;
    use tonic::{Status, Request, Response};
    use std::time::Duration;
    use rand::Rng;
    
    async fn tester_client(sleep: Duration, uri: &str, key: &str) {
        let uri = uri.parse().unwrap();
        let key = key.parse::<tonic::metadata::MetadataValue<_>>().unwrap();
        
        loop {
            debug!("create new connection");
            // 构建带keep-alive配置的Channel
            let chan = tonic::transport::Channel::builder(uri.clone())
                .timeout(Duration::from_secs(20))
                .connect_timeout(Duration::from_secs(20))
                .http2_keep_alive_interval(Duration::from_secs(5))
                .keep_alive_while_idle(true)
                .connect_lazy();
    
            let mut service = test_grpc::say_client::SayClient::with_interceptor(
                chan,
                move |mut req: Request<()>| {
                    req.metadata_mut().insert("authorization", key.clone());
                    Ok(req)
                },
            );
    
            // 内层循环处理请求,遇致命错误则跳出重建连接
            let mut connection_valid = true;
            while connection_valid {
                let send_text = format!("interval({})", sleep.as_secs_f32() / 60.0);
                debug!("send: {send_text}");
                match service
                    .hello(Request::new(test_grpc::RequestSay {
                        text: send_text.clone(),
                    }))
                    .await
                {
                    Ok(res) => {
                        debug!("recv: {}", res.into_inner().text);
                        time::sleep(sleep).await;
                    }
                    Err(e) => {
                        error!("{e:#}");
                        // 判定为无法恢复的连接错误,标记重建
                        if matches!(e.code(), tonic::Code::Unavailable | tonic::Code::Cancelled | tonic::Code::DeadlineExceeded) {
                            connection_valid = false;
                        }
                        // 短暂退避后重试当前请求
                        time::sleep(Duration::from_secs(1)).await;
                    }
                }
                println!();
            }
            
            // 重建连接前随机退避,避免服务端压力
            let mut rng = rand::thread_rng();
            let backoff = Duration::from_secs(rng.gen_range(2..5));
            time::sleep(backoff).await;
        }
    }
    
  3. 关键配置说明

    • http2_keep_alive_interval + keep_alive_while_idle:保持空闲连接活跃,提前发现连接失效
    • 双层循环:内层处理正常请求,外层负责在连接失效时重建客户端实例
    • 随机退避:避免大量客户端同时重建连接导致服务端雪崩

补充说明

如果希望更优雅的重试逻辑,可以使用tower::retry中间件包装Channel,通过自定义策略捕获连接错误并自动重试,无需手动维护双层循环。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 20:05:31