0Pricing
gRPC & High Performance APIs · レッスン

サーバーストリーミングの解説

1つのクライアントリクエストに対してサーバーが複数のレスポンスを送信する、サーバーサイドストリーミングを理解・実装します。

「サーバーストリーミングの解説」はCoddyKit上の無料gRPC & High Performance APIsレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはgRPC & High Performance APIs学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 gRPC & High Performance APIsコースには全4レッスンが含まれています。

このレッスンの一部はまだ翻訳されておらず、英語で表示されています。

Server Streaming Basics

In gRPC, a server-side streaming call is when a client sends a single request, but the server responds with a sequence of messages.

Think of it like subscribing to a newsletter: you send one request (subscribe), and the server sends you many updates over time (newsletters).

Use Cases for Streaming

Server streaming is perfect for scenarios where the server needs to push updates or a large amount of data to the client over time. Common uses include:

  • Real-time data feeds: Stock prices, sensor readings.
  • Notifications: Alerts, chat messages.
  • Large data downloads: Breaking a big file into smaller chunks.

Protobuf for Server Streaming

To define a server-side streaming method in your .proto file, you simply add the stream keyword before the response type.

This tells gRPC that the server will send multiple messages for each client request, not just one.

Streaming Protobuf Example

Here's how you define a service that streams messages from the server:

syntax = "proto3";

option java_package = "com.coddykit.grpc";
option java_outer_classname = "StreamingProto";

service NotifierService {
  rpc SubscribeToNotifications (SubscriptionRequest) returns (stream Notification);
}

message SubscriptionRequest {
  string userId = 1;
}

message Notification {
  string message = 1;
  int64 timestamp = 2;
}

Implementing the Server Stream

On the server side, your streaming method will receive a single request object, just like a unary call. However, instead of returning a single response, you'll use a StreamObserver to send multiple responses back to the client.

You'll typically loop and send messages, then call onCompleted() when done.

Server Method Structure

The server method for a streaming call takes the request and a StreamObserver. You send responses via responseObserver.onNext() and signal completion with responseObserver.onCompleted().

// Example (Java)
public void subscribeToNotifications(SubscriptionRequest request,
    io.grpc.stub.StreamObserver<Notification> responseObserver) {

  String userId = request.getUserId();
  System.out.println("Client " + userId + " subscribed.");

  // Simulate sending multiple notifications
  for (int i = 0; i < 3; i++) {
    Notification notification = Notification.newBuilder()
        .setMessage("Update " + (i + 1) + " for " + userId)
        .setTimestamp(System.currentTimeMillis())
        .build();
    responseObserver.onNext(notification); // Send a message
    try {
      Thread.sleep(1000); // Wait a bit
    } catch (InterruptedException e) { /* handle */ }
  }

  responseObserver.onCompleted(); // Signal completion
  System.out.println("Finished sending notifications to " + userId);
}

Receiving Streamed Responses

The client makes a single call, but then it needs to wait and process multiple responses. It provides a StreamObserver to handle the incoming messages, errors, and the completion signal from the server.

This observer will have onNext(), onError(), and onCompleted() methods.

Client Stream Observer

The client's StreamObserver defines how it reacts to each event from the server stream. It processes each onNext message until onCompleted is called.

// Example (Java)
StreamObserver<Notification> responseObserver = new StreamObserver<Notification>() {
  @Override
  public void onNext(Notification notification) {
    System.out.println("Received: " + notification.getMessage());
  }

  @Override
  public void onError(Throwable t) {
    System.err.println("Error: " + t.getMessage());
  }

  @Override
  public void onCompleted() {
    System.out.println("Server stream completed.");
  }
};

// Call the streaming method
// asyncStub.subscribeToNotifications(request, responseObserver);

Complete Server Stream Service

Here's a complete gRPC server that implements the SubscribeToNotifications server-side streaming method. Run this first, then the client!

import io.grpc.Server;
import io.grpc.ServerBuilder;
import io.grpc.stub.StreamObserver;
import com.coddykit.grpc.StreamingProto.SubscriptionRequest;
import com.coddykit.grpc.StreamingProto.Notification;
import com.coddykit.grpc.NotifierServiceGrpc.NotifierServiceImplBase;

public class StreamingServer {

  private Server server;

  private void start() throws Exception {
    int port = 50051;
    server = ServerBuilder.forPort(port)
        .addService(new NotifierServiceImpl())
        .build()
        .start();
    System.out.println("Server started, listening on " + port);

    Runtime.getRuntime().addShutdownHook(new Thread() {
      @Override
      public void run() {
        System.err.println("*** shutting down gRPC server since JVM is shutting down");
        StreamingServer.this.stop();
        System.err.println("*** server shut down");
      }
    });
  }

  private void stop() {
    if (server != null) {
      server.shutdown();
    }
  }

  private void blockUntilShutdown() throws InterruptedException {
    if (server != null) {
      server.awaitTermination();
    }
  }

  public static void main(String[] args) throws Exception {
    final StreamingServer server = new StreamingServer();
    server.start();
    server.blockUntilShutdown();
  }

  static class NotifierServiceImpl extends NotifierServiceImplBase {
    @Override
    public void subscribeToNotifications(SubscriptionRequest request,
        StreamObserver<Notification> responseObserver) {

      String userId = request.getUserId();
      System.out.println("Server received subscription from: " + userId);

      for (int i = 0; i < 3; i++) {
        Notification notification = Notification.newBuilder()
            .setMessage("Update " + (i + 1) + " for " + userId)
            .setTimestamp(System.currentTimeMillis())
            .build();
        responseObserver.onNext(notification);
        try {
          Thread.sleep(1000); // Simulate some work
        } catch (InterruptedException e) {
          Thread.currentThread().interrupt();
          responseObserver.onError(e);
          return;
        }
      }
      responseObserver.onCompleted();
      System.out.println("Server finished sending notifications to: " + userId);
    }
  }
}

Complete Client Stream Receiver

Now, run this client code. It will connect to the server and receive the stream of notifications.

import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import com.coddykit.grpc.StreamingProto.SubscriptionRequest;
import com.coddykit.grpc.StreamingProto.Notification;
import com.coddykit.grpc.NotifierServiceGrpc;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public class StreamingClient {

  public static void main(String[] args) throws InterruptedException {
    ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 50051)
        .usePlaintext() // For local testing, no TLS
        .build();

    NotifierServiceGrpc.Stub asyncStub = NotifierServiceGrpc.newStub(channel);

    CountDownLatch latch = new CountDownLatch(1);

    System.out.println("Client sending subscription request...");
    SubscriptionRequest request = SubscriptionRequest.newBuilder()
        .setUserId("user123")
        .build();

    asyncStub.subscribeToNotifications(request, new StreamObserver<Notification>() {
      @Override
      public void onNext(Notification notification) {
        System.out.println("Client received notification: " + notification.getMessage());
      }

      @Override
      public void onError(Throwable t) {
        System.err.println("Client received error: " + t.getMessage());
        latch.countDown();
      }

      @Override
      public void onCompleted() {
        System.out.println("Client stream completed.");
        latch.countDown();
      }
    });

    latch.await(5, TimeUnit.SECONDS); // Wait for stream to complete
    System.out.println("Client finished.");

    channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
  }
}

Stream Method Check

You're building a gRPC service where a client requests a list of recent log entries, and the server continuously sends new entries as they occur. Which Protobuf definition correctly sets up the GetLogStream method for this?

Streaming Recap

Great job! In this lesson, you've learned about server-side streaming in gRPC.

  • It allows a server to send multiple responses for a single client request.
  • It's defined using the stream keyword on the response type in Protobuf.
  • You implemented both server and client logic to handle these continuous data flows.

Next, we'll explore client-side streaming, where the client sends multiple requests!

よくある質問

「サーバーストリーミングの解説」レッスンは無料ですか?

はい。「サーバーストリーミングの解説」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、gRPC & High Performance APIsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 gRPC & High Performance APIsコースには全4レッスンが含まれています。

「サーバーストリーミングの解説」で何を学びますか?

1つのクライアントリクエストに対してサーバーが複数のレスポンスを送信する、サーバーサイドストリーミングを理解・実装します。 ブラウザで直接実行するハンズオンコードでgRPC & High Performance APIsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

gRPC & High Performance APIsを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのgRPC & High Performance APIsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。

「サーバーストリーミングの解説」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このgRPC & High Performance APIsレッスンでコードを書いて実行できますか?

はい。すべてのgRPC & High Performance APIsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. サーバーストリーミングの解説
  2. クライアントストリーミングの解説
  3. 双方向ストリーミング
  4. ストリーミングのフロー制御とバックプレッシャー
← gRPC & High Performance APIsに戻る