開始使用 gRPC-Rust - Streaming

1. 簡介

在本程式碼研究室中,您將使用 gRPC-Rust 建立用戶端和伺服器,做為以 Rust 編寫的路線對應應用程式基礎。

在本教程結束時,您將擁有一個客戶端,該客戶端使用 gRPC 連接到遠端伺服器,以獲取有關客戶端路由上特徵的信息,創建客戶端路由的摘要,並與伺服器和其他客戶端交換路由信息,例如流量更新。

服務是在 Protocol Buffers 檔案中定義,這個檔案會用於產生用戶端和伺服器的樣板程式碼,讓兩者可以互相通訊,節省您實作這項功能的時間和精力。

產生的程式碼不僅處理伺服器和客戶端之間通訊的複雜性,還處理資料序列化和反序列化。

課程內容

  • 如何使用通訊協定緩衝區定義服務 API。
  • 如何使用自動化程式碼產生技術,從 Protocol Buffers 定義建置基於 gRPC 的客戶端和伺服器。
  • 瞭解如何使用 gRPC 進行用戶端與伺服器之間的串流通訊。

本程式碼實驗室適合剛接觸 gRPC 或想複習 gRPC 的 Rust 開發人員,以及有興趣建構分散式系統的任何人。無需具備 gRPC 使用經驗。

2. 事前準備

必要條件

請確認已安裝下列項目:

  • GCC。請按照這裡的指示操作。
  • Git:安裝說明請參閱這裡
  • Rust 1.88.0 版。請按照這裡的安裝說明操作。

取得程式碼

為避免您必須從頭開始,本程式碼研究室提供應用程式的原始碼架構,供您完成。以下步驟將向您展示如何完成應用程序,包括使用協定緩衝區編譯器插件產生樣板 gRPC 程式碼。

首先,建立程式碼研究室目前使用的目錄,然後 cd 到該目錄:

mkdir streaming-grpc-rust-getting-started && cd streaming-grpc-rust-getting-started

下載並擷取程式碼研究室:

curl -sL https://github.com/grpc-ecosystem/grpc-codelabs/archive/refs/heads/2026.tar.gz \
  | tar xvz --strip-components=4 \
  grpc-codelabs-2026/codelabs/grpc-rust-streaming/start_here

或者,您也可以下載只包含 Codelab 目錄的 .zip 檔案,然後手動解壓縮。

如要略過實作的輸入作業,可以在 GitHub 上取得完整的原始碼

3. 定義訊息和服務

首先,請使用通訊協定緩衝區定義應用程式的 gRPC 服務、RPC 方法,以及要求和回應訊息類型。您的服務將提供:

  • 伺服器實作的 RPC 方法稱為 ListFeaturesRecordRouteRouteChat,客戶端呼叫這些方法。
  • 訊息類型 PointFeatureRectangleRouteNoteRouteSummary,這是呼叫上述方法時,用戶端與伺服器之間交換的資料結構。

這些 RPC 方法及其訊息類型都會在提供的原始碼 proto/routeguide.proto 檔案中定義。

協定緩衝區通常被稱為 protobuf。如要進一步瞭解 gRPC 術語,請參閱 gRPC 的「核心概念、架構和生命週期」。

定義訊息類型

首先,請定義 RPC 會使用的訊息。在原始碼的 proto/routeguide.proto 檔案中,請先定義 Point 訊息型別。Point 代表地圖上的經緯度座標組合。在本程式碼研究室中,請使用整數做為座標:

message Point {
  int32 latitude = 1;
  int32 longitude = 2;
}

數字 12message 結構中每個欄位的專屬 ID 號碼。

接著定義 Feature 訊息類型。Feature 會使用 string 欄位,為 Point 指定位置的項目提供名稱或郵寄地址:

message Feature {
  // The name or address of the feature.
  string name = 1;

  // The point where the feature is located.
  Point location = 2;
}

接著是 Rectangle 訊息,代表經緯度矩形,以兩個對角點「lo」和「hi」表示。

message Rectangle {
  // One corner of the rectangle.
  Point lo = 1;

  // The other corner of the rectangle.
  Point hi = 2;
}

還有一條 RouteNote 訊息,表示在給定時間點發送的訊息。

message RouteNote {
  // The location from which the message is sent.
  Point location = 1;

  // The message to be sent.
  string message = 2;
}

我們也需要 RouteSummary 訊息。這則訊息是針對 RecordRoute RPC 收到的回應,下一節將說明這項 RPC。其中包含收到的個別點數、偵測到的特徵數量,以及涵蓋的總距離 (各點之間距離的累計總和)。

message RouteSummary {
  // The number of points received.
  int32 point_count = 1;

  // The number of known features passed while traversing the route.
  int32 feature_count = 2;

  // The distance covered in metres.
  int32 distance = 3;

  // The duration of the traversal in seconds.
  int32 elapsed_time = 4;
}

定義服務方法

我們先定義服務,稍後再定義訊息。如要定義服務,請在 .proto 檔案中指定具名服務。proto/routeguide.proto 檔案具有名為 RouteGuideservice 結構,可定義應用程式服務提供的一或多個方法。

在服務定義中定義 RPC 方法,並指定其請求和回應類型。 在本程式碼研究室的這一節中,請定義:

ListFeatures

取得指定 Rectangle 內可用的 Feature。由於矩形可能涵蓋大範圍並包含大量特徵,因此系統會串流處理結果,而非一次傳回 (例如在含有重複欄位的訊息中)。

這個 RPC 的適當類型是伺服器端串流 RPC:用戶端會將要求傳送至伺服器,並取得串流來讀取一系列訊息。客戶端會從傳回的資料流讀取訊息,直到沒有更多訊息為止。如範例所示,您可以在回應型別前加上 stream 關鍵字,指定伺服器端串流方法。

rpc ListFeatures(Rectangle) returns (stream Feature) {}

RecordRoute

接受所遍歷路徑上的 Point 串流,並在遍歷完成時傳回 RouteSummary

在這種情況下,客戶端流 RPC 似乎比較合適:客戶端寫入一系列訊息,並再次使用提供的流將它們傳送到伺服器。客戶端完成訊息寫入後,會等待伺服器讀取所有訊息並回傳回應。您可以透過在請求類型前放置 stream 關鍵字來指定用戶端串流傳輸方法。

rpc RecordRoute(stream Point) returns (RouteSummary) {}

RouteChat

接受在路線遍歷期間傳送的 RouteNote 串流,同時接收其他 RouteNote (例如來自其他使用者)。

這正是雙向串流的適用用途。雙向串流 RPC 的兩端都會使用讀寫串流傳送一連串訊息。這兩個串流各自獨立運作,因此用戶端和伺服器可以依任意順序讀取及寫入資料。

舉例來說,伺服器可以等待接收所有用戶端訊息,再撰寫回覆內容;也可以讀取訊息,然後撰寫訊息;或是讀取和撰寫訊息的某種組合。

每個資料流中的消息順序都得以保留。如要指定這類方法,請在要求和回應前加上 stream 關鍵字。

rpc RouteChat(stream RouteNote) returns (stream RouteNote) {}

4. 產生用戶端和伺服器程式碼

我們已在 generated/ 目錄中提供 .proto 檔案產生的程式碼,包括您在上方所做的所有新增項目。不過,我們想花點時間說明程式碼生成功能的運作方式。

我們的 .proto 檔案描述了客戶端或伺服器使用的所有結構和函數。我們使用 Cargo 建置腳本 (build.rs) 和 grpc-protobuf-build crate 來自動產生此程式碼。

Cargo.toml 中,我們已經將 grpc-protobuf-build 新增為建置依賴項。

build.rs 中,我們配置 grpc_protobuf_build::CodeGenproto/routeguide.proto 編譯到 generated/ 目錄中。關鍵語句如下:

grpc_protobuf_build::CodeGen::new()
    .include("proto")
    .input("routeguide.proto")
    .output_dir("generated")
    .compile()
    .unwrap();

這會呼叫 grpc_protobuf_build Crate 的程式碼產生作業,並將 routeguide.proto 傳遞給該作業。我們已將這項功能包裝在程式碼中,只會在傳遞功能旗標時執行,因此只會在您想要時重新產生。我們已為您產生程式碼,因此您不必立即執行這項操作。

cargo build --bin routeguide-server --features regenerate_proto

執行 cargo build 時,build.rs 會將通訊協定緩衝區定義編譯到 generate/ 目錄,包括:

  • 訊息型別 PointFeature 的結構體定義。
  • 我們需要為伺服器實作的 Tonic 服務特徵是 route_guide_server::RouteGuide
  • 我們將用來呼叫伺服器的 gRPC-Rust 用戶端類型:route_guide_client::RouteGuideClient<T>

詳情請參閱 protoc-gen-rust-grpc 指南

接著,我們會在伺服器上實作服務方法。

5. 實作服務

首先,我們來看看如何建立 RouteGuide 伺服器。讓 RouteGuide 服務正常運作的程序分為兩個部分:

  • 實作從服務定義產生的服務介面:執行服務的實際「工作」。
  • 執行 gRPC 伺服器,監聽來自用戶端的要求,並將要求分派至正確的方法實作。

src/server/server.rs 中,我們可以透過 gRPC 的 include_generated_proto! 巨集,將產生的程式碼帶入範圍,並匯入 RouteGuide 特徵和 Point

mod grpc_pb {
    grpc::include_generated_proto!("generated", "routeguide");
}

pub use grpc_pb::{
    route_guide_server::{RouteGuideServer, RouteGuide},
    Point, Feature, Rectangle, RouteNote, RouteSummary
};

我們可以先定義一個結構體來表示我們的服務。目前我們可以在 src/server/server.rs 上執行此操作:

#[derive(Debug)]
pub struct RouteGuideService {
    features: Vec<Feature>,
}

現在,我們需要從產生的程式碼實作 route_guide_server::RouteGuide 特徵。

實作 RouteGuide

我們需要實作產生的 RouteGuide 介面。實作方式如下所示。模板裡已經包含了這部分內容。

#[tonic::async_trait]
impl RouteGuide for RouteGuideService {
    async fn list_features(
        &self,
        request: Request<Rectangle>,
    ) -> Result<Response<ListFeaturesStream>, Status> {
        ...
    }

    async fn record_route(
        &self,
        request: Request<tonic::Streaming<Point>>,
    ) -> Result<Response<RouteSummary>, Status> {
        ...
    }

    async fn route_chat(
        &self,
        request: Request<tonic::Streaming<RouteNote>>,
    ) -> Result<Response<RouteChatStream>, Status> {
        ...
    }
}

讓我們詳細瞭解每種 RPC 實作。

伺服器端串流 RPC:ListFeatures

首先播放《ListFeatures》。這是一個伺服器端串流 RPC(客戶端將發送一條訊息,伺服器將回應多個訊息),因此我們需要向客戶端發送多個 Feature

async fn list_features(
        &self,
        request: Request<Rectangle>,
    ) -> Result<Response<ListFeaturesStream>, Status> {
    println!("ListFeatures = {:?}", request);

    let (tx, rx) = mpsc::channel(4);
    let features = self.features.clone();

    tokio::spawn(async move {
        for feature in &features[..] {
            if in_range(&feature.location().to_owned(), request.get_ref()) {
                println!("  => send {feature:?}");
                tx.send(Ok(feature.clone())).await.unwrap();
            }
        }
        println!(" /// done sending");
    });

    let output_stream = ReceiverStream::new(rx);
    Ok(Response::new(Box::pin(output_stream)))
}

如您所見,我們取得要求物件 (用戶端要尋找的 Features 中的 Rectangle)。這次我們需要傳回值串流。我們會建立管道並產生新的非同步工作,在其中執行查閱作業,然後將符合限制的特徵傳送至管道。系統會將管道的 Stream 半部包裝在 tonic::Response 中,然後傳回給呼叫端。

用戶端串流 RPC:RecordRoute

現在來看看稍微複雜一點的內容:用戶端串流方法 RecordRoute,我們會從用戶端取得 Points 串流,並傳回包含行程資訊的單一 RouteSummary。它接收一個資料流作為輸入,伺服器可以使用該資料流來讀取和寫入訊息。它可以透過 next() 方法疊代處理用戶端訊息,並傳回單一回應。

async fn record_route(
        &self,
        request: Request<tonic::Streaming<Point>>,
    ) -> Result<Response<RouteSummary>, Status> {
    println!("RecordRoute");
    let mut stream = request.into_inner();
    let mut summary = RouteSummary::default();
    let mut last_point = None;
    let now = Instant::now();

    while let Some(point) = stream.next().await {
        let point = point?;
        println!("  ==> Point = {point:?}");

        // Increment the point count
        summary.set_point_count(summary.point_count() + 1);

        // Find features
        for feature in &self.features[..] {
            if feature.location().latitude() == point.latitude() {
                if feature.location().longitude() == point.longitude(){
                    summary.set_feature_count(summary.feature_count() + 1);
                }
            }
        }

        // Calculate the distance
        if let Some(ref last_point) = last_point {
            let new_dist = summary.distance() + calc_distance(last_point, &point);
            summary.set_distance(new_dist);
        }
        last_point = Some(point);
    }
    summary.set_elapsed_time(now.elapsed().as_secs() as i32);
    Ok(Response::new(summary))
}

在方法體中,我們使用流的 next() 方法重複讀取客戶端對請求物件(在本例中為 Point)的請求,直到沒有更多訊息為止。如果為 None,表示串流仍正常,可以繼續讀取。

雙向串流遠端程序呼叫:RouteChat

最後,我們來看看雙向串流 RPC RouteChat()

async fn route_chat(
        &self,
        request: Request<tonic::Streaming<RouteNote>>,
    ) -> Result<Response<RouteChatStream>, Status> {
    println!("RouteChat");

    let mut notes: HashMap<(i32, i32), Vec<RouteNote>> = HashMap::new();
    let mut stream = request.into_inner();

    let output = async_stream::try_stream! {
        while let Some(note) = stream.next().await {
            let note = note?;
            let location = note.location();
            let key = (location.latitude(), location.longitude());
            let location_notes = notes.entry(key).or_insert(vec![]);
            location_notes.push(note);
            for note in location_notes {
                yield note.clone();
            }
        }
    };
    Ok(Response::new(Box::pin(output)))
}

這次我們得到的是一個流,就像我們在客戶端串流範例中一樣,可以用來讀取和寫入訊息。但是,這次,我們透過方法的流傳回值,而客戶端仍在向其訊息流寫入訊息。這裡的讀取和寫入語法與用戶端串流方法非常相似,但伺服器會傳回 RouteChatStream。雖然雙方一律會依撰寫順序收到對方的訊息,但用戶端和伺服器可以依任意順序讀取及寫入資料,因為資料串是完全獨立運作。

我們使用 try_stream! 建立輸出串流,表示串流可能會傳回錯誤。

啟動伺服器

實作這個方法後,我們也需要啟動 gRPC 伺服器,讓用戶端實際使用我們的服務。填入 main()

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let addr = "[::1]:10000".parse().unwrap();
    println!("RouteGuideServer listening on: {addr}");
    let route_guide = RouteGuideService {
        features: load(),
    };
    let svc = RouteGuideServer::new(route_guide);
    Server::builder().add_service(svc).serve(addr).await?;
    Ok(())
}

以下是 main() 的運作方式:

  1. 指定要用來監聽用戶端要求的通訊埠
  2. 建立載入功能的 RouteGuideService
  3. 使用我們建立的服務,透過 RouteGuideServer::new() 建立 gRPC 伺服器實例。
  4. 向 gRPC 伺服器註冊服務實作。
  5. 使用我們的連接埠詳細資訊在伺服器上呼叫 serve(),以執行阻塞等待,直到進程被終止。

6. 創建客戶端

在本節中,我們將探討如何在 src/client/client.rs 中為我們的 RouteGuide 服務建立一個 Rust 用戶端。

首先,請將產生的程式碼納入範圍。

mod grpc_pb {
    grpc::include_generated_proto!("generated", "routeguide");
}

use grpc_pb::route_guide_client::RouteGuideClient;
use grpc_pb::{Point, Rectangle, RouteNote};

呼叫服務方法

現在來看看如何呼叫服務方法。在 gRPC-Rust 中,流式 RPC 是非同步的、非阻塞的,使用 Rust 的 async/await 語法和 Tokio 流。

伺服器端串流 RPC:PrintFeatures

在伺服器串流 RPC 中,用戶端會將單一要求訊息傳送至伺服器,並取得回應訊息串流。以下是我們在 client.rs 中呼叫伺服器端串流方法 list_features() 的位置 (對應於 proto 中找到的 ListFeatures rpc 宣告)。伺服器接著會傳回一連串 Feature 訊息:

async fn print_features(client: &RouteGuideClient<Channel>) -> Result<(), Box<dyn Error>> {
    let rectangle = proto!(Rectangle {
        lo: proto!(Point {
            latitude: 400_000_000,
            longitude: -750_000_000,
        }),
        hi: proto!(Point {
            latitude: 420_000_000,
            longitude: -730_000_000,
        }),
    });

    let mut stream = client.list_features(rectangle).await;

    while let Some(feature) = stream.recv().await {
        println!(
            "FEATURE: Name = \"{}\", Lat = {}, Lon = {}",
            feature.name(),
            feature.location().latitude(),
            feature.location().longitude()
        );
    }
    let status = stream.status().await;
    assert!(status.is_ok(), "{:?}", status);
    Ok(())
}

用戶端串流 RPC:RecordRoute

使用用戶端串流時,用戶端會開啟連往伺服器的串流,並傳送一連串訊息。串流結束時,該物件會收到單一回應訊息。

在此,我們使用 client.record_route().await 啟動呼叫,透過 stream.send(point).await 在串流中逐一傳送多個產生的 Point 座標,然後使用 stream.close_and_recv().await 關閉串流,接收單一伺服器 RouteSummary 訊息。

async fn run_record_route(client: &RouteGuideClient<Channel>) -> Result<(), Box<dyn Error>> {
    let mut rng = rand::rng();
    let point_count: i32 = rng.random_range(2..100);

    let mut points = vec![];
    for _ in 0..=point_count {
        points.push(random_point(&mut rng));
    }

    println!("Traversing {} points", points.len());
    let mut stream = client.record_route().await;

    for point in &points {
        if stream.send(point).await.is_err() {
            break;
        }
    }

    match stream.close_and_recv().await {
        Ok(response) => {
            println!(
                "SUMMARY: Feature Count = {}, Distance = {}",
                response.feature_count(),
                response.distance()
            );
        }
        Err(e) => println!("something went wrong: {e:?}"),
    }
    Ok(())
}

雙向串流遠端程序呼叫:RouteChat

最後,我們來看看雙向串流 RPC RouteChat()。在這裡,用戶端和伺服器都會來回傳遞一連串訊息。我們會產生 tokio 工作,透過 tx.send(note).await.is_err() 持續將訊息傳送至伺服器。同時,rx.recv().await 會監聽伺服器的回應訊息,並在收到時列印這些訊息。

async fn run_route_chat(client: &RouteGuideClient<Channel>) -> Result<(), Box<dyn Error>> {
    let (mut tx, mut rx) = client.route_chat().await;

    let start = time::Instant::now();
    tokio::spawn(async move {
        let mut interval = time::interval(Duration::from_millis(50));
        for _ in 0..10 {
            let time = interval.tick().await;
            let elapsed = time.duration_since(start);
            let note = proto!(RouteNote {
                location: proto!(Point {
                    latitude: 409146138 + elapsed.as_millis() as i32,
                    longitude: -746188906,
                }),
                message: format!("at {elapsed:?}"),
            });
            if tx.send(note).await.is_err() {
                return;
            }
        }
        tx.close();
    });

    while let Some(note) = rx.recv().await {
        println!(
            "Note: Latitude = {}, Longitude = {}, Message = \"{}\"",
            note.location().latitude(),
            note.location().longitude(),
            note.message()
        );
    }
    let status = rx.status().await;
    assert!(status.is_ok(), "{:?}", status);
    Ok(())
}

雖然雙方一律會依撰寫順序收到對方的訊息,但用戶端和伺服器可以依任意順序讀取及寫入資料,因為資料串是完全獨立運作。

建立及傳遞用戶端

如要呼叫服務方法,我們必須先建立與伺服器通訊的管道。我們會先建立端點,然後連線至該端點,並在連線至 RouteGuideClient::new() 時傳遞所建立的管道,藉此建立這個項目,如下所示:

// Create channel to connect to server
let channel = Channel::builder(
    "dns:///[::1]:10000",
    Arc::new(LocalChannelCredentials::new()),
)
.build();

// Create a new client
let client = RouteGuideClient::new(channel);

建立這個用戶端後,我們就可以呼叫上述方法,並將用戶端傳遞至這些方法。我們將所有這些程式碼加入 main(),後者使用 Tokio 非同步執行階段。以下是完整程式碼:

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // Create channel to connect to server
    let channel = Channel::builder(
        "dns:///[::1]:10000",
        Arc::new(LocalChannelCredentials::new()),
    )
    .build();

    // Create a new client
    let client = RouteGuideClient::new(channel);

    println!("\n*** SERVER STREAMING ***");
    print_features(&client).await?;

    println!("\n*** CLIENT STREAMING ***");
    run_record_route(&client).await?;

    println!("\n*** BIDIRECTIONAL STREAMING ***");
    run_route_chat(&client).await?;

    Ok(())
}

7. 立即試用

如要執行用戶端和伺服器,請先確認 Cargo.toml 中定義了兩個二進位目標:

[[bin]]
name = "routeguide-server"
path = "src/server/server.rs"

[[bin]]
name = "routeguide-client"
path = "src/client/client.rs"

接著,從目前使用的目錄執行下列指令:

  1. 在一個終端機中執行伺服器:
cargo run --bin routeguide-server
  1. 從另一個終端機執行用戶端:
cargo run --bin routeguide-client

您會看到如下所示的輸出內容:

*** SERVER STREAMING ***
FEATURE: Name = "Patriots Path, Mendham, NJ 07945, USA", Lat = 407838351, Lon = -746143763
FEATURE: Name = "101 New Jersey 10, Whippany, NJ 07981, USA", Lat = 408122808, Lon = -743999179
FEATURE: Name = "U.S. 6, Shohola, PA 18458, USA", Lat = 413628156, Lon = -749015468
...
*** CLIENT STREAMING ***
Traversing 86 points
SUMMARY: Feature Count = 0, Distance = 803709356

*** BIDIRECTIONAL STREAMING ***
Note: Latitude = 409146138, Longitude = -746188906, Message = "at 112.45µs"
Note: Latitude = 409146139, Longitude = -746188906, Message = "at 1.00011245s"
Note: Latitude = 409146140, Longitude = -746188906, Message = "at 2.00011245s"

8. 後續步驟

9. 本程式碼研究室的協作者

  • Cathy Zhao
  • Arvind Bright
  • Nathaniel Ford