客户端与双向流式传输
构建客户端流式传输和全双工双向流式传输通道,以应对高吞吐量场景。
客户端与双向流式传输 是 CoddyKit 上的免费 C# Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 C# Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 C# Academy 课程共包含 4 节课。
Client 流式传输与双向流式传输
除了单向和服务端流式传输之外,gRPC 还支持客户端流式传输(client 发送多条消息,服务器回复一次)和双向流式传输(双方通过一个连接同时发送多条消息)。
Client 流式传输:Proto 定义
在请求类型前添加 stream 关键字。client 会发送一系列消息;完成后,服务器返回单个响应。
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; }在服务器上实现 Client 流式传输
使用 IAsyncStreamReader<T> 读取到达的 client 消息。调用 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 };
}从 Client 调用 Client 流式传输
打开流式 call,使用 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,以实现全双工通信。client 和服务器可以随时独立发送消息。
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);
}从 Client 进行双向流式传输
启动两个任务来同时发送和接收:一个用于写入消息,另一个用于读取回复。
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 能够发出“发送完成”信号,同时继续等待服务器响应。
// 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);
}使用通道实现线程安全的双向流式传输
将 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 });
}
}快速检查
在客户端流式传输中,服务器何时发送单个响应?
回顾:Client 与双向流式传输
要点:
- 客户端流式传输:client 发送多条消息,服务器发送一条 reply,适合文件上传
- 双向流式传输:双方独立发送,适合聊天和实时数据流
- 使用
ReadAllAsync()和await foreach消费流 - 调用
RequestStream.CompleteAsync()半关闭 client 的写入端 - HTTP/2 流控制会自动处理背压,请不要缓冲整个流
- System.Threading.Channels 与双向流式传输结合良好,适用于消息分发场景
常见问题解答
「客户端与双向流式传输」课时是免费的吗?
是的 — 「客户端与双向流式传输」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 C# Academy 课程的其余内容,请升级到 CoddyKit PRO。 C# Academy 课程共包含 4 节课。
「客户端与双向流式传输」这节课中我会学到什么?
构建客户端流式传输和全双工双向流式传输通道,以应对高吞吐量场景。 你通过在浏览器中直接运行的动手代码来练习 C# Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 C# Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 C# Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「客户端与双向流式传输」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 C# Academy 课中编写并运行代码吗?
能。每节 C# Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。