gRPC-Rust スタートガイド - ストリーミング

1. はじめに

この Codelab では、gRPC-Rust を使用して、Rust で記述されたルート マッピング アプリケーションの基盤となるクライアントとサーバーを作成します。

このチュートリアルを終えると、gRPC を使用してリモート サーバーに接続し、クライアントのルート上の機能に関する情報を取得し、クライアントのルートの概要を作成し、トラフィックの更新などのルート情報をサーバーや他のクライアントと交換するクライアントが完成します。

サービスはプロトコル バッファ ファイルで定義されます。このファイルを使用して、クライアントとサーバーが相互に通信できるように、クライアントとサーバーのボイラープレート コードを生成します。これにより、この機能を実装する時間と労力を節約できます。

生成されたコードは、サーバーとクライアント間の通信の複雑さだけでなく、データのシリアル化と逆シリアル化も処理します。

学習内容

  • プロトコル バッファを使用してサービス API を定義する方法。
  • 自動コード生成を使用して、プロトコル バッファ定義から gRPC ベースのクライアントとサーバーを構築する方法。
  • gRPC を使用したクライアントとサーバーのストリーミング通信について。

この Codelab は、gRPC を初めて使用する Rust デベロッパー、gRPC の復習を希望するデベロッパー、分散システムの構築に関心のある方を対象としています。gRPC の経験は必要ありません。

2. 始める前に

前提条件

以下がインストールされていることを確認してください。

  • GCC。こちらの手順を行います
  • Git: インストール方法については、こちらをご覧ください。
  • Rust バージョン 1.88.0。インストール方法については、こちらをご覧ください。

コードを取得する

この Codelab では、完全にゼロから始める必要がないように、アプリケーションのソースコードのスケルトンを提供しています。次の手順では、プロトコル バッファ コンパイラ プラグインを使用してボイラープレート gRPC コードを生成するなど、アプリケーションを完成させる方法について説明します。

まず、Codelab の作業ディレクトリを作成して、そのディレクトリに cd します。

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

Codelab をダウンロードして解凍します。

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 は、地図上の緯度と経度の座標ペアを表します。この Codelab では、座標に整数を使用します。

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

数値 12 は、message 構造体の各フィールドの一意の ID 番号です。

次に、Feature メッセージ タイプを定義します。Feature は、Point で指定された場所にあるものの名前または郵便番号に string フィールドを使用します。

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

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

次に、緯度と経度の長方形を表す Rectangle メッセージ。これは、対角線上の 2 つの点「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 へのレスポンスとして受信されます。受信した個々のポイントの数、検出された機能の数、各ポイント間の距離の累積合計としてカバーされた合計距離が含まれます。

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 ファイルには、アプリケーションのサービスによって提供される 1 つ以上のメソッドを定義する RouteGuide という名前の service 構造体があります。

サービス定義内で RPC メソッドを定義し、リクエストとレスポンスのタイプを指定します。 この Codelab のセクションでは、以下を定義します。

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 では、両側が読み取り / 書き込みストリームを使用して一連のメッセージを送信します。2 つのストリームは独立して動作するため、クライアントとサーバーは任意の順序で読み取りと書き込みを行うことができます。

たとえば、サーバーはクライアント メッセージをすべて受信してからレスポンスを書き込むことも、メッセージを読み取ってからメッセージを書き込むことも、読み取りと書き込みの組み合わせを行うこともできます。

各ストリーム内のメッセージの順序は保持されます。このタイプのメソッドを指定するには、リクエストとレスポンスの両方の前に stream キーワードを配置します。

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

4. クライアントとサーバーのコードを生成する

上記の追加を含め、.proto ファイルから生成されたコードは generated/ ディレクトリに用意されています。ただし、コード生成の仕組みについて説明します。

.proto ファイルには、クライアントまたはサーバーが使用するすべての構造体と関数が記述されています。このコードを自動的に生成するには、Cargo ビルドスクリプト(build.rs)と grpc-protobuf-build クレートを使用します。

Cargo.toml には、ビルド依存関係として grpc-protobuf-build がすでに追加されています。

build.rs で、grpc_protobuf_build::CodeGen を構成して、proto/routeguide.protogenerated/ ディレクトリにコンパイルします。重要な行は次のとおりです。

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

これにより、grpc_protobuf_build クレートのコード生成が呼び出され、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 サービスが機能するには、次の 2 つの要素が必要です。

  • サービス定義から生成されたサービス インターフェースを実装する: サービスの実際の「作業」を行う。
  • 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(クライアントが 1 つのメッセージを送信し、サーバーが多数のメッセージで応答する)であるため、複数の 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 の場合でも、ストリームは良好で、読み取りを続行できます。

双方向ストリーミング RPC: 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 で呼び出しを開始し、生成された Point 座標を stream.send(point).await を使用してストリーム経由で 1 つずつ送信し、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(())
}

双方向ストリーミング RPC: 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);

このクライアントを作成したら、上記のメソッドを呼び出して、クライアントを渡すことができます。このコードはすべて、Tokio 非同期ランタイムを使用する main() に追加します。完全なコードは次のとおりです。

#[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. 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. この Codelab の共同編集者

  • Cathy Zhao
  • Arvind Bright
  • Nathaniel Ford