0Pricing
C# Academy · 课时

客户端与双向流式传输

构建客户端流式传输和全双工双向流式传输通道,以应对高吞吐量场景。

客户端与双向流式传输 是 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. gRPC 与 Protobuf 基础
  2. 一元与服务器流式 RPC
  3. 客户端与双向流式传输
  4. 截止时间、取消与拦截器
← 返回 C# Academy