クライアントストリーミングの解説
クライアントが一連のメッセージをサーバーへ送信できる、クライアントサイドストリーミングの実装方法を学習します。
「クライアントストリーミングの解説」はCoddyKit上の無料gRPC & High Performance APIsレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはgRPC & High Performance APIs学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 gRPC & High Performance APIsコースには全4レッスンが含まれています。
このレッスンの一部はまだ翻訳されておらず、英語で表示されています。
What is Client Streaming?
Welcome to client streaming! In gRPC, client streaming is a communication pattern where the client sends a sequence of messages to the server.
Unlike a simple unary RPC (request-response), the client doesn't just send one message. Instead, it sends a stream of messages, and the server processes them, then sends back a single response at the end.
How Client Streaming Works
Imagine uploading a large file by sending it in many small chunks. The server collects all chunks, rebuilds the file, and then sends a single "upload complete" confirmation.
- The client initiates the RPC.
- The client sends multiple messages asynchronously.
- The server receives and processes these messages.
- Once the client finishes sending (signals completion), the server sends back a single response.
Defining Client Stream in Protobuf
To define a client-streaming method in your .proto file, you use the stream keyword for the request type, but not for the response type.
Here's an example for a log upload service:
syntax = "proto3";
package client_streaming;
service LogService {
rpc UploadLogs (stream LogEntry) returns (UploadSummary);
}
message LogEntry {
string message = 1;
int64 timestamp = 2;
}
message UploadSummary {
int32 uploaded_count = 1;
string status_message = 2;
}Server: Handling the Client Stream
On the server side, your method will receive a StreamObserver for the client's incoming messages and will use another StreamObserver to send its single response.
The server's StreamObserver will have onNext() for each incoming message, onError() for errors, and onCompleted() when the client finishes sending.
import io.grpc.stub.StreamObserver;
import io.grpc.Server;
import io.grpc.ServerBuilder;
import client_streaming.LogEntry;
import client_streaming.LogServiceGrpc;
import client_streaming.UploadSummary;
public class LogServer {
private Server server;
private void start() throws Exception {
int port = 50051;
server = ServerBuilder.forPort(port)
.addService(new LogServiceImpl())
.build()
.start();
System.out.println("Server started, listening on " + port);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
System.err.println("*** shutting down gRPC server since JVM is shutting down");
LogServer.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 LogServer logServer = new LogServer();
logServer.start();
logServer.blockUntilShutdown();
}
static class LogServiceImpl extends LogServiceGrpc.LogServiceImplBase {
@Override
public StreamObserver<LogEntry> uploadLogs(StreamObserver<UploadSummary> responseObserver) {
return new StreamObserver<LogEntry>() {
private int logCount = 0;
@Override
public void onNext(LogEntry log) {
// Process each log entry as it arrives
System.out.println("Received log: " + log.getMessage() + " at " + log.getTimestamp());
logCount++;
}
@Override
public void onError(Throwable t) {
System.err.println("UploadLogs cancelled or failed: " + t.getMessage());
responseObserver.onError(t);
}
@Override
public void onCompleted() {
// After all logs are received, send a single summary response
UploadSummary summary = UploadSummary.newBuilder()
.setUploadedCount(logCount)
.setStatusMessage("Successfully processed " + logCount + " log entries.")
.build();
responseObserver.onNext(summary);
responseObserver.onCompleted();
System.out.println("Finished processing client stream. Sent summary.");
}
};
}
}
}Client: Sending the Stream
On the client side, you'll get a StreamObserver to send your messages. You call onNext() for each message you want to send and finally onCompleted() to signal the end of the stream.
The server's single response will be handled by a separate StreamObserver you provide.
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import client_streaming.LogEntry;
import client_streaming.LogServiceGrpc;
import client_streaming.UploadSummary;
import java.util.concurrent.TimeUnit;
public class LogClient {
private final LogServiceGrpc.LogServiceStub asyncStub;
private final ManagedChannel channel;
public LogClient(String host, int port) {
channel = ManagedChannelBuilder.forAddress(host, port)
.usePlaintext() // For demonstration, use plaintext
.build();
asyncStub = LogServiceGrpc.newStub(channel);
}
public void shutdown() throws InterruptedException {
channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
}
public void uploadMultipleLogs() throws InterruptedException {
StreamObserver<UploadSummary> responseObserver = new StreamObserver<UploadSummary>() {
@Override
public void onNext(UploadSummary summary) {
System.out.println("Server Response: " + summary.getStatusMessage() + " (" + summary.getUploadedCount() + " logs)");
}
@Override
public void onError(Throwable t) {
System.err.println("UploadLogs failed: " + t.getMessage());
}
@Override
public void onCompleted() {
System.out.println("Server has completed processing.");
}
};
StreamObserver<LogEntry> requestObserver = asyncStub.uploadLogs(responseObserver);
try {
// Send multiple log entries
LogEntry log1 = LogEntry.newBuilder().setMessage("User login attempt").setTimestamp(System.currentTimeMillis()).build();
LogEntry log2 = LogEntry.newBuilder().setMessage("Database query executed").setTimestamp(System.currentTimeMillis() + 100).build();
LogEntry log3 = LogEntry.newBuilder().setMessage("API call completed").setTimestamp(System.currentTimeMillis() + 200).build();
requestObserver.onNext(log1);
System.out.println("Client sent log 1");
Thread.sleep(100); // Simulate some delay
requestObserver.onNext(log2);
System.out.println("Client sent log 2");
Thread.sleep(100);
requestObserver.onNext(log3);
System.out.println("Client sent log 3");
// Mark the end of the client stream
requestObserver.onCompleted();
System.out.println("Client finished sending logs.");
// Wait for server response (handled by responseObserver)
Thread.sleep(1000); // Give time for server to respond
} catch (RuntimeException e) {
requestObserver.onError(e);
throw e;
}
}
public static void main(String[] args) throws Exception {
LogClient client = new LogClient("localhost", 50051);
try {
client.uploadMultipleLogs();
} finally {
client.shutdown();
}
}
}Running the Example
To see client streaming in action:
- First, compile your
.protofile to generate the necessary Java classes. - Run the LogServer application. It will start listening for requests.
- Then, run the LogClient application. It will send three log entries and wait for the server's summary.
Observe the console output from both the client and server to understand the flow of messages.
Key StreamObserver Methods
The StreamObserver interface is crucial for handling streaming RPCs. Both the client and server use implementations of this interface.
onNext(T value): Called for each message received in the stream. The client uses this to send messages, the server uses it to receive.onError(Throwable t): Called if an RPC fails or is cancelled.onCompleted(): Called when the stream has finished. The client calls this after sending all messages; the server calls it after sending its final response.
When to Use Client Streaming
Client streaming is ideal for scenarios where a client needs to send a large amount of data or a series of related messages to a server, and only cares about a single final result.
- Large File Uploads: Sending a file chunk by chunk.
- Log Aggregation: A client sending many log entries to a central logging service.
- Batch Operations: Sending a list of items to be processed as a single batch, receiving a summary.
- Sensor Data Collection: Continuously sending readings from a sensor device.
Quick Check
You're designing a gRPC service where a client needs to send a series of sensor readings to a server, and the server will process them and return a single summary report.
Recap: Client Streaming
You've learned about client-side streaming in gRPC!
- Client streaming allows a client to send a sequence of messages.
- The server processes these messages and sends a single response back.
- In Protobuf, you define it by using the
streamkeyword for the request type. - Both client and server use
StreamObserverto manage the flow of messages (onNext(),onError(),onCompleted()). - It's great for tasks like uploading large data or sending continuous log entries.
Next, we'll explore server-side streaming!
よくある質問
「クライアントストリーミングの解説」レッスンは無料ですか?
はい。「クライアントストリーミングの解説」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、gRPC & High Performance APIsコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 gRPC & High Performance APIsコースには全4レッスンが含まれています。
「クライアントストリーミングの解説」で何を学びますか?
クライアントが一連のメッセージをサーバーへ送信できる、クライアントサイドストリーミングの実装方法を学習します。 ブラウザで直接実行するハンズオンコードでgRPC & High Performance APIsを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
gRPC & High Performance APIsを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのgRPC & High Performance APIsは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。
「クライアントストリーミングの解説」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このgRPC & High Performance APIsレッスンでコードを書いて実行できますか?
はい。すべてのgRPC & High Performance APIsレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- サーバーストリーミングの解説
- クライアントストリーミングの解説
- 双方向ストリーミング
- ストリーミングのフロー制御とバックプレッシャー