クライアントと双方向ストリーミング
高スループットのシナリオに向けて、クライアントストリーミングと全二重の双方向ストリーミングチャネルを構築します。
「クライアントと双方向ストリーミング」はCoddyKit上の無料C# Academyレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはC# Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 C# Academyコースには全4レッスンが含まれています。
クライアントストリーミングと双方向ストリーミング
Unary とサーバーストリーミングに加えて、gRPC はクライアントストリーミング(クライアントが複数送信し、サーバーが 1 回応答する)と双方向ストリーミング(1 つの接続上で両側が同時に複数のメッセージを送信する)をサポートしています。
クライアントストリーミング: Proto 定義
リクエスト型の前に stream キーワードを追加します。クライアントは一連のメッセージを送信し、完了するとサーバーが 1 つのレスポンスを返します。
service UploadService {
// Client streams chunks, server returns a summary
rpc UploadFile (stream FileChunk) returns (UploadSummary);
}
message FileChunk { bytes data = 1; string filename = 2; }
message UploadSummary { int64 bytes_received = 1; string checksum = 2; }サーバーでクライアントストリーミングを実装する
到着したクライアントメッセージを読み取るには IAsyncStreamReader<T> を使用します。MoveNext() を呼び出すか、ReadAllAsync() で反復処理します。
public override async Task<UploadSummary> UploadFile(
IAsyncStreamReader<FileChunk> requestStream,
ServerCallContext context)
{
long totalBytes = 0;
using var ms = new MemoryStream();
await foreach (var chunk in requestStream.ReadAllAsync(context.CancellationToken))
{
await ms.WriteAsync(chunk.Data.Memory, context.CancellationToken);
totalBytes += chunk.Data.Length;
}
var checksum = ComputeMd5(ms.ToArray());
return new UploadSummary { BytesReceived = totalBytes, Checksum = checksum };
}クライアントからクライアントストリーミングを呼び出す
ストリーミング呼び出しを開始し、RequestStream.WriteAsync() でメッセージを書き込み、CompleteAsync() で完了を通知します。その後、レスポンスを待機します。
using var call = client.UploadFile();
var fileBytes = await File.ReadAllBytesAsync("large-file.bin");
const int chunkSize = 64 * 1024; // 64 KB
for (int offset = 0; offset < fileBytes.Length; offset += chunkSize)
{
var chunk = fileBytes.Skip(offset).Take(chunkSize).ToArray();
await call.RequestStream.WriteAsync(new FileChunk
{
Filename = "large-file.bin",
Data = Google.Protobuf.ByteString.CopyFrom(chunk)
});
}
await call.RequestStream.CompleteAsync(); // signal end
var summary = await call; // await server response
Console.WriteLine($"Uploaded {summary.BytesReceived} bytes");双方向ストリーミング: Proto 定義
全二重通信を行うには、両側に stream を追加します。クライアントとサーバーは、互いに独立して、いつでもメッセージを送信できます。
service ChatService {
rpc Chat (stream ChatMessage) returns (stream ChatMessage);
}
message ChatMessage {
string user = 1;
string content = 2;
int64 sent_at = 3;
}サーバーで双方向ストリーミングを実装する
IAsyncStreamReader から読み取り、IServerStreamWriter に同時に書き込みます。Task.WhenAll を使用して、両方のループを同時に実行します。
public override async Task Chat(
IAsyncStreamReader<ChatMessage> requestStream,
IServerStreamWriter<ChatMessage> responseStream,
ServerCallContext context)
{
// Broadcast to all connected clients
var readTask = Task.Run(async () =>
{
await foreach (var msg in requestStream.ReadAllAsync(context.CancellationToken))
{
_chatHub.Broadcast(msg);
}
});
var writeTask = Task.Run(async () =>
{
await foreach (var msg in _chatHub.GetMessagesAsync(context.CancellationToken))
{
await responseStream.WriteAsync(msg);
}
});
await Task.WhenAll(readTask, writeTask);
}クライアントから双方向ストリーミングを行う
2 つのタスク(メッセージを書き込むタスクと返信を読み取るタスク)を起動して、送信と受信を同時に行います。
using var call = client.Chat();
// Read task
var readTask = Task.Run(async () =>
{
await foreach (var msg in call.ResponseStream.ReadAllAsync())
Console.WriteLine($"{msg.User}: {msg.Content}");
});
// Write task
while (Console.ReadLine() is string text && text != "exit")
{
await call.RequestStream.WriteAsync(
new ChatMessage { User = "Alice", Content = text,
SentAt = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() });
}
await call.RequestStream.CompleteAsync();
await readTask;フロー制御とバックプレッシャー
HTTP/2 にはフロー制御が組み込まれています。コンシューマーの処理が遅い場合、gRPC は送信側を自動的に一時停止し、メモリのオーバーフローを防ぎます。メッセージを手動でバッファーせず、ストリームに任せてください。
// DO: Let gRPC flow control handle backpressure
await foreach (var chunk in requestStream.ReadAllAsync(ct))
{
await ProcessChunkAsync(chunk); // naturally paced
}
// DON'T: Buffering all chunks defeats flow control
var allChunks = new List<FileChunk>();
await foreach (var chunk in requestStream.ReadAllAsync(ct))
allChunks.Add(chunk); // may OOM on large uploadsハーフクローズとストリームの完了
双方向ストリーミングでは、一方から読み取りを続けながら、どちら側でも書き込み側を閉じることができます。これにより、クライアントはサーバーからのレスポンスを待ちながら「送信完了」を通知できます。
// Client signals it's done sending
await call.RequestStream.CompleteAsync(); // half-close
// Continue reading server responses after half-close
await foreach (var reply in call.ResponseStream.ReadAllAsync())
{
Console.WriteLine(reply.Content);
}スレッドセーフな双方向ストリーミングに Channels を使用する
System.Threading.Channels と双方向ストリーミングを組み合わせると、多数のリスナーにメッセージを安全に分配できる、すっきりとしたプロデューサー・コンシューマーパターンを実現できます。
private readonly Channel<ChatMessage> _broadcast =
Channel.CreateUnbounded<ChatMessage>();
// Producer: receives from each client stream
public async Task ReceiveLoopAsync(
IAsyncStreamReader<ChatMessage> reader, CancellationToken ct)
{
await foreach (var msg in reader.ReadAllAsync(ct))
await _broadcast.Writer.WriteAsync(msg, ct);
}
// Consumer: writes to each client's response stream
public async Task SendLoopAsync(
IServerStreamWriter<ChatMessage> writer, CancellationToken ct)
{
await foreach (var msg in _broadcast.Reader.ReadAllAsync(ct))
await writer.WriteAsync(msg);
}実例: リアルタイム分析パイプライン
データ取り込みサービスが IoT デバイスからテレメトリイベントのストリームを受け取り、集計した統計情報をリアルタイムでストリーミングして返します。これは双方向ストリーミングに最適なユースケースです。
service TelemetryService {
rpc StreamTelemetry(stream TelemetryEvent) returns (stream AggregatedStats);
}
// Server implementation sends a rolling aggregate every 100 events:
public override async Task StreamTelemetry(
IAsyncStreamReader<TelemetryEvent> requests,
IServerStreamWriter<AggregatedStats> responses,
ServerCallContext context)
{
int count = 0;
double total = 0;
await foreach (var e in requests.ReadAllAsync(context.CancellationToken))
{
total += e.Value;
if (++count % 100 == 0)
await responses.WriteAsync(
new AggregatedStats { Count = count, Average = total / count });
}
}確認問題
クライアントストリーミングでは、サーバーはいつ 1 つのレスポンスを送信しますか?
まとめ: クライアントストリーミングと双方向ストリーミング
重要なポイント:
- クライアントストリーミング: クライアントが複数のメッセージを送信し、サーバーが 1 つの返信を送信する。ファイルアップロードに適している
- 双方向ストリーミング: 両側が独立して送信する。チャットやライブデータフィードに適している
await foreachとReadAllAsync()を使用してストリームを利用するRequestStream.CompleteAsync()を呼び出してクライアントの書き込み側をハーフクローズする- HTTP/2 のフロー制御がバックプレッシャーを自動的に処理するため、ストリーム全体をバッファーしない
- System.Threading.Channels は、ファンアウトが必要なシナリオで双方向ストリーミングと相性がよい
よくある質問
「クライアントと双方向ストリーミング」レッスンは無料ですか?
はい。「クライアントと双方向ストリーミング」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、C# Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 C# Academyコースには全4レッスンが含まれています。
「クライアントと双方向ストリーミング」で何を学びますか?
高スループットのシナリオに向けて、クライアントストリーミングと全二重の双方向ストリーミングチャネルを構築します。 ブラウザで直接実行するハンズオンコードでC# Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
C# Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのC# Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。
「クライアントと双方向ストリーミング」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このC# Academyレッスンでコードを書いて実行できますか?
はい。すべてのC# Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- gRPCとProtobufの基礎
- UnaryとサーバーストリーミングRPC
- クライアントと双方向ストリーミング
- デッドライン、キャンセルとインターセプター