-
C#网络编程之gRPC实战(Protocol Buffers、双向流、拦截器)
第26章 gRPC实战
26.1 gRPC实战(Protocol Buffers、双向流、拦截器)
一、我踩过的gRPC坑:从“REST轮询卡顿”到“PB版本兼容事故”
刚做实时监控系统时,用REST API轮询服务器获取监控数据,延迟高达1秒,服务器CPU跑满(每秒处理1000次请求);后来转用gRPC的服务器流RPC,延迟直接降到100毫秒,CPU占用减少70%——但踩了一堆坑:一开始PB(Protocol Buffers)定义没加默认值,客户端升级后旧服务器返回的数据缺失字段,导致客户端崩溃;双向流聊天功能没处理线程安全,多个客户端同时发消息时服务器抛出并发异常;拦截器顺序配置错误,认证拦截器在日志拦截器之后执行,导致非法请求也被记录日志,浪费磁盘空间。这节我把这些踩坑经验揉进去,用大白话讲透gRPC的核心原理,结合C#实战vb.net教程C#教程python教程SQL教程access 2010教程代码逐行讲解Protocol Buffers、四种RPC类型、拦截器的使用,以及版本兼容、性能优化的最佳实践,让你的微服务通信既高效又可靠。
二、Protocol Buffers(PB):gRPC的“灵魂”,比JSON快5倍
Protocol Buffers是Google开发的一种轻量级、高效的序列化协议,是gRPC的默认序列化方式——就像你把数据压缩成一个极小的包,传输速度快,解析速度也快。它的核心优势是体积小、速度快、强类型,比JSON、XML高效得多。
核心优势(拓展知识:PB vs JSON对比)
| 特性 | Protocol Buffers | JSON |
| ---- | ---- | ---- | ---- |
| 体积大小 | 小(比如一个对象占100字节) | 大(同样对象占500字节) |
| 序列化速度 | 快(比JSON快5-10倍) | 慢(解析需要字符串处理) |
| 类型安全 | 强类型(编译时检查错误) | 弱类型(运行时发现错误) |
| 版本兼容 | 支持(新增字段不影响旧客户端) | 差(新增字段可能导致旧客户端崩溃) |
| 可读性 | 差(二进制格式,需要工具解析) | 好(人类可读的字符串) |
类比:PB就像你把文件压缩成ZIP包,体积小、传输快;JSON就像你把文件存成TXT,体积大、传输慢,但容易读。
实战1:编写第一个PB文件并生成C#代码
步骤1:编写.proto文件(定义服务和消息)
创建monitor.proto文件,定义监控数据的消息和gRPC服务:
proto
syntax = "proto3"; // 使用proto3语法,比proto2更简洁
option csharp_namespace = "GrpcMonitor"; // 生成的C#代码的命名空间
// 定义监控数据消息
message MonitorData {
int32 id = 1; // 字段编号,必须唯一,不能重复,版本兼容的关键
string service_name = 2; // 服务名
float cpu_usage = 3; // CPU使用率(0-100)
float memory_usage = 4; // 内存使用率(0-100)
int64 timestamp = 5; // 时间戳(毫秒)
}
// 定义获取监控数据的请求消息
message GetMonitorRequest {
string service_name = 1; // 要监控的服务名,空表示所有服务
}
// 定义gRPC服务
service MonitorService {
// 简单RPC:客户端发请求,服务器返回响应
rpc GetSingleMonitorData(GetMonitorRequest) returns (MonitorData);
// 服务器流RPC:客户端发请求,服务器返回多个响应(比如实时推送监控数据)
rpc StreamMonitorData(GetMonitorRequest) returns (stream MonitorData);
// 客户端流RPC:客户端发多个请求,服务器返回一个响应(比如上传批量数据)
rpc BatchUploadMonitorData(stream MonitorData) returns (UploadResponse);
// 双向流RPC:客户端和服务器互相发多个响应(比如聊天、游戏实时交互)
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
// 上传批量数据的响应消息
message UploadResponse {
int32 success_count = 1; // 成功上传的数量
int32 fail_count = 2; // 失败上传的数量
}
// 聊天消息
message ChatMessage {
string username = 1; // 用户名
string content = 2; // 消息内容
int64 timestamp = 3; // 时间戳
}
PB语法逐行讲解:
1.syntax = "proto3":使用proto3语法,比proto2更简洁,默认值更合理(比如int32默认0,string默认空);
2.option csharp_namespace:生成的C#代码的命名空间,避免和其他代码冲突;
3.message:定义数据结构,每个字段有字段编号(比如id=1),这是版本兼容的关键——字段编号不能修改,新增字段用新的编号,旧客户端会忽略新增的字段;
4.service:定义gRPC服务,里面的方法是RPC类型,支持四种类型:简单RPC、服务器流、客户端流、双向流;
5.stream:标记流类型,比如returns (stream MonitorData)表示服务器返回多个MonitorData消息。
我踩过的坑:一开始修改了字段编号,比如把id=1改成id=6,导致旧客户端解析数据时字段错误,崩溃了——字段编号是PB版本兼容的核心,一旦发布就不能修改,只能新增字段!
步骤2:生成C#代码
用Visual Studio的gRPC插件自动生成代码,或者用protoc命令:
bash
# 安装protoc和C#插件
dotnet tool install --global Grpc.Tools
# 生成C#代码(包含服务和消息)
protoc --csharp_out=. --grpc_out=. --plugin=protoc-gen-grpc=`which grpc_csharp_plugin` monitor.proto
生成的文件有两个:Monitor.cs(消息类)和MonitorGrpc.cs(服务基类和客户端类)。
三、gRPC四种服务类型实战:从简单请求到实时交互
gRPC支持四种服务类型,每种适合不同的场景,下面逐个讲解并给出C#实战代码。
实战2:简单RPC(客户端请求→服务器响应)
适用场景:普通的请求响应,比如获取单个监控数据、查询用户信息。
服务器代码
csharp
using Grpc.Core;
using GrpcMonitor;
using System;
using System.Threading.Tasks;
namespace GrpcServer;
// 实现MonitorService的基类
public class MonitorServiceImpl : MonitorService.MonitorServiceBase
{
// 实现简单RPC方法
public override Task<MonitorData> GetSingleMonitorData(GetMonitorRequest request, ServerCallContext context)
{
// 模拟从数据库获取监控数据
var monitorData = new MonitorData
{
Id = 1,
ServiceName = request.ServiceName ?? "OrderService",
CpuUsage = 45.2f,
MemoryUsage = 67.8f,
Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()
};
Console.WriteLine($"Received GetSingleMonitorData request for service: {request.ServiceName}");
return Task.FromResult(monitorData);
}
}
// 启动gRPC服务器
class Program
{
static async Task Main(string[] args)
{
// 监听50051端口(gRPC默认端口)
var server = new Server
{
Services = { MonitorService.BindService(new MonitorServiceImpl()) },
Ports = { new ServerPort("localhost", 50051, ServerCredentials.Insecure) } // 开发环境用不安全的凭证,生产环境用TLS
};
server.Start();
Console.WriteLine("gRPC server listening on port 50051");
Console.WriteLine("Press any key to stop...");
Console.ReadKey();
await server.ShutdownAsync();
}
}
代码逐行讲解:
1.MonitorServiceImpl:继承自动生成的MonitorService.MonitorServiceBase,实现RPC方法;
2.GetSingleMonitorData:简单RPC方法,接收GetMonitorRequest,返回MonitorData;
3.Server:gRPC服务器,绑定服务和端口,开发环境用ServerCredentials.Insecure(不安全),生产环境必须用TLS(SslServerCredentials);
4.Start:启动服务器,监听50051端口。
客户端代码
csharp
using Grpc.Core;
using GrpcMonitor;
using System;
using System.Threading.Tasks;
namespace GrpcClient;
class Program
{
static async Task Main(string[] args)
{
// 连接gRPC服务器
var channel = new Channel("localhost:50051", ChannelCredentials.Insecure);
var client = new MonitorService.MonitorServiceClient(channel);
// 调用简单RPC方法
var request = new GetMonitorRequest { ServiceName = "ProductService" };
var response = await client.GetSingleMonitorDataAsync(request);
Console.WriteLine($"Received Monitor Data:");
Console.WriteLine($"ID: {response.Id}");
Console.WriteLine($"Service Name: {response.ServiceName}");
Console.WriteLine($"CPU Usage: {response.CpuUsage}%");
Console.WriteLine($"Memory Usage: {response.MemoryUsage}%");
Console.WriteLine($"Timestamp: {DateTimeOffset.FromUnixTimeMilliseconds(response.Timestamp)}");
await channel.ShutdownAsync();
}
}
代码逐行讲解:
1.Channel:gRPC客户端连接通道,指定服务器地址和凭证;
2.MonitorServiceClient:自动生成的客户端类,调用RPC方法;
3.GetSingleMonitorDataAsync:异步调用简单RPC方法,返回响应。
实战3:服务器流RPC(客户端请求→服务器持续推送)
适用场景:实时推送数据,比如监控数据、股票行情、日志推送。
服务器代码(在MonitorServiceImpl中添加)
csharp
// 实现服务器流RPC方法
public override async Task StreamMonitorData(GetMonitorRequest request, IServerStreamWriter<MonitorData> responseStream, ServerCallContext context)
{
var serviceName = request.ServiceName ?? "AllServices";
Console.WriteLine($"Starting stream for service: {serviceName}");
// 每秒推送一次监控数据,直到客户端断开连接
while (!context.CancellationToken.IsCancellationRequested)
{
var monitorData = new MonitorData
{
Id = new Random().Next(1000),
ServiceName = serviceName,
CpuUsage = new Random().Next(0, 100) + (float)new Random().NextDouble(),
MemoryUsage = new Random().Next(0, 100) + (float)new Random().NextDouble(),
Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()
};
// 发送监控数据到客户端
await responseStream.WriteAsync(monitorData);
Console.WriteLine($"Sent monitor data: {monitorData.CpuUsage}% CPU");
// 等待1秒
await Task.Delay(1000, context.CancellationToken);
}
Console.WriteLine($"Stream for service {serviceName} stopped");
}
代码逐行讲解:
1.IServerStreamWriter:服务器流的写入器,用于向客户端发送多个响应;
2.context.CancellationToken:客户端断开连接时触发的取消令牌,用于停止推送;
3.WriteAsync:向客户端发送单个监控数据;
4.Task.Delay:每秒推送一次,模拟实时监控。
客户端代码
csharp
// 调用服务器流RPC方法
var streamRequest = new GetMonitorRequest { ServiceName = "OrderService" };
using var streamCall = client.StreamMonitorData(streamRequest);
// 循环接收服务器推送的消息
await foreach (var data in streamCall.ResponseStream.ReadAllAsync())
{
Console.WriteLine($"[{DateTimeOffset.FromUnixTimeMilliseconds(data.Timestamp)}] {data.ServiceName}: CPU {data.CpuUsage:F1}%, Memory {data.MemoryUsage:F1}%");
}
代码逐行讲解:
1.StreamMonitorData:调用服务器流方法,返回AsyncServerStreamingCall;
2.ResponseStream.ReadAllAsync():异步遍历服务器推送的所有消息,直到服务器停止推送或客户端断开连接。
我踩过的坑:一开始没处理context.CancellationToken,客户端断开连接后服务器还在推送,导致内存泄漏——必须检查取消令牌,及时停止推送!
实战4:双向流RPC(客户端和服务器互相推送)
适用场景:实时交互,比如聊天、游戏、协作编辑。
服务器代码(在MonitorServiceImpl中添加)
csharp
using System.Collections.Concurrent;
// 保存所有客户端的流,线程安全
private static readonly ConcurrentDictionary<string, IServerStreamWriter<ChatMessage>> _chatClients = new ConcurrentDictionary<string, IServerStreamWriter<ChatMessage>>();
// 实现双向流RPC方法
public override async Task Chat(IAsyncStreamReader<ChatMessage> requestStream, IServerStreamWriter<ChatMessage> responseStream, ServerCallContext context)
{
// 获取用户名(从请求的第一个消息中获取)
ChatMessage firstMessage = null;
try
{
firstMessage = await requestStream.MoveNext(context.CancellationToken);
if (!requestStream.Current.HasUsername)
{
await responseStream.WriteAsync(new ChatMessage { Content = "Error: Username is required", Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() });
return;
}
}
catch
{
return;
}
var username = requestStream.Current.Username;
Console.WriteLine($"User {username} joined chat");
// 把当前客户端的流加入字典
_chatClients.TryAdd(username, responseStream);
try
{
// 广播用户加入的消息
await BroadcastMessage(new ChatMessage { Username = "System", Content = $"{username} joined the chat", Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() });
// 循环接收客户端发送的消息
while (await requestStream.MoveNext(context.CancellationToken))
{
var message = requestStream.Current;
message.Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds();
Console.WriteLine($"Received message from {username}: {message.Content}");
// 广播消息给所有客户端
await BroadcastMessage(message);
}
}
finally
{
// 用户离开,从字典中删除
_chatClients.TryRemove(username, out _);
await BroadcastMessage(new ChatMessage { Username = "System", Content = $"{username} left the chat", Timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() });
Console.WriteLine($"User {username} left chat");
}
}
// 广播消息给所有客户端
private async Task BroadcastMessage(ChatMessage message)
{
foreach (var client in _chatClients.Values)
{
try
{
await client.WriteAsync(message);
}
catch
{
// 客户端断开连接,忽略
}
}
}
代码逐行讲解:
1.IAsyncStreamReader:客户端流的读取器,用于接收客户端发送的多个消息;
2.ConcurrentDictionary:线程安全的字典,保存所有在线客户端的流,避免并发问题;
3.BroadcastMessage:把消息广播给所有在线客户端;
4.finally块:用户离开时从字典中删除,避免内存泄漏。
客户端代码
csharp
// 调用双向流RPC方法
using var chatCall = client.Chat();
// 启动一个线程发送消息
var sendTask = Task.Run(async () =>
{
Console.WriteLine("Enter your username:");
var username = Console.ReadLine();
while (true)
{
Console.Write("Enter message (or 'exit' to quit): ");
var content = Console.ReadLine();
if (content == "exit")
break;
var message = new ChatMessage { Username = username, Content = content };
await chatCall.RequestStream.WriteAsync(message);
}
await chatCall.RequestStream.CompleteAsync();
});
// 接收服务器广播的消息
var receiveTask = Task.Run(async () =>
{
await foreach (var message in chatCall.ResponseStream.ReadAllAsync())
{
Console.WriteLine($"[{DateTimeOffset.FromUnixTimeMilliseconds(message.Timestamp)}] {message.Username}: {message.Content}");
}
});
await Task.WhenAny(sendTask, receiveTask);
代码逐行讲解:
1.Chat():调用双向流方法,返回AsyncDuplexStreamingCall;
2.RequestStream.WriteAsync:向服务器发送消息;
3.ResponseStream.ReadAllAsync():接收服务器广播的消息;
4.Task.WhenAny:等待发送或接收任务完成,用户输入exit时退出。
四、gRPC拦截器:给所有RPC方法加“通用逻辑”
gRPC拦截器是gRPC的高级功能,用于给所有RPC方法添加通用逻辑,比如日志、认证、限流、重试——就像你给所有接口加了一个统一的“过滤器”,不需要在每个方法里重复写代码。
拦截器的两种类型(拓展知识)
1.客户端拦截器:在客户端调用RPC方法前/后执行逻辑,比如添加认证头、记录请求时间;
2.服务器拦截器:在服务器处理RPC方法前/后执行逻辑,比如验证认证头、记录响应时间。
实战5:日志拦截器(服务器端)
csharp
using Grpc.Core;
using Grpc.Core.Interceptors;
using System;
using System.Diagnostics;
using System.Threading.Tasks;
namespace GrpcServer;
// 日志拦截器:记录每个RPC请求的时间、方法名、状态
public class LoggingInterceptor : Interceptor
{
public override async Task<TResponse> UnaryServerHandler<TRequest, TResponse>(TRequest request, ServerCallContext context, UnaryServerMethod<TRequest, TResponse> continuation)
{
var stopwatch = Stopwatch.StartNew();
try
{
// 调用实际的RPC方法
var response = await continuation(request, context);
stopwatch.Stop();
Console.WriteLine($"[INFO] {context.Method} succeeded in {stopwatch.ElapsedMilliseconds}ms");
return response;
}
catch (RpcException ex)
{
stopwatch.Stop();
Console.WriteLine($"[ERROR] {context.Method} failed with status {ex.StatusCode} in {stopwatch.ElapsedMilliseconds}ms: {ex.Message}");
throw;
}
}
// 重写其他类型的拦截方法:服务器流、客户端流、双向流
public override async Task ServerStreamingServerHandler<TRequest, TResponse>(TRequest request, IServerStreamWriter<TResponse> responseStream, ServerCallContext context, ServerStreamingServerMethod<TRequest, TResponse> continuation)
{
var stopwatch = Stopwatch.StartNew();
try
{
await continuation(request, responseStream, context);
stopwatch.Stop();
Console.WriteLine($"[INFO] {context.Method} stream completed in {stopwatch.ElapsedMilliseconds}ms");
}
catch (RpcException ex)
{
stopwatch.Stop();
Console.WriteLine($"[ERROR] {context.Method} stream failed with status {ex.StatusCode} in {stopwatch.ElapsedMilliseconds}ms: {ex.Message}");
throw;
}
}
}
代码逐行讲解:
1.LoggingInterceptor:继承Interceptor,重写对应类型的拦截方法;
2.UnaryServerHandler:简单RPC的拦截方法,continuation是实际的RPC方法;
3.Stopwatch:记录RPC方法的执行时间;
4.try-catch:捕获RPC异常,记录错误日志,然后重新抛出异常。
服务器注册拦截器
csharp
// 启动服务器时注册拦截器
var server = new Server
{
Services = { MonitorService.BindService(new MonitorServiceImpl()).Intercept(new LoggingInterceptor()) },
Ports = { new ServerPort("localhost", 50051, ServerCredentials.Insecure) }
};
我踩过的坑:一开始拦截器顺序配置错误,认证拦截器在日志拦截器之后执行,导致非法请求也被记录日志,浪费磁盘空间——服务器拦截器的执行顺序是先添加的先执行,所以要把认证拦截器放在日志拦截器之前,先验证再记录日志!
实战6:JWT认证拦截器(服务器端)
csharp
using Grpc.Core;
using Grpc.Core.Interceptors;
using System;
using System.IdentityModel.Tokens.Jwt;
using System.Security.Claims;
using System.Text;
using Microsoft.IdentityModel.Tokens;
namespace GrpcServer;
// JWT认证拦截器:验证客户端的JWT令牌
public class JwtAuthInterceptor : Interceptor
{
private readonly string _secretKey = "your-secret-key-1234567890"; // 生产环境要从配置文件读取
public override async Task<TResponse> UnaryServerHandler<TRequest, TResponse>(TRequest request, ServerCallContext context, UnaryServerMethod<TRequest, TResponse> continuation)
{
// 从请求头获取JWT令牌
if (!context.RequestHeaders.TryGetValue("Authorization", out var authHeader))
{
throw new RpcException(new Status(StatusCode.Unauthenticated, "Authorization header missing"));
}
var token = authHeader.Value;
if (!token.StartsWith("Bearer "))
{
throw new RpcException(new Status(StatusCode.Unauthenticated, "Invalid token format"));
}
token = token.Substring(7); // 去掉"Bearer "前缀
try
{
// 验证JWT令牌
var tokenHandler = new JwtSecurityTokenHandler();
var key = Encoding.ASCII.GetBytes(_secretKey);
tokenHandler.ValidateToken(token, new TokenValidationParameters
{
ValidateIssuerSigningKey = true,
IssuerSigningKey = new SymmetricSecurityKey(key),
ValidateIssuer = false, // 生产环境要验证Issuer
ValidateAudience = false, // 生产环境要验证Audience
ClockSkew = TimeSpan.Zero // 不允许时间偏移
}, out SecurityToken validatedToken);
// 获取用户信息,比如用户名
var jwtToken = (JwtSecurityToken)validatedToken;
var username = jwtToken.Claims.First(x => x.Type == ClaimTypes.Name).Value;
Console.WriteLine($"Authenticated user: {username}");
// 把用户信息存入context,供RPC方法使用
context.UserState["Username"] = username;
}
catch
{
throw new RpcException(new Status(StatusCode.Unauthenticated, "Invalid token"));
}
// 调用实际的RPC方法
return await continuation(request, context);
}
}
代码逐行讲解:
1.RequestHeaders.TryGetValue:从gRPC请求头中获取Authorization头;
2.JWT验证:用JwtSecurityTokenHandler验证令牌的签名和有效期;
3.context.UserState:把用户信息存入上下文,供后续的RPC方法使用;
4.RpcException:验证失败时抛出Unauthenticated状态码,客户端会收到对应的错误。
客户端添加JWT令牌
csharp
// 客户端调用时添加Authorization头
var headers = new Metadata();
headers.Add("Authorization", "Bearer your-jwt-token-here");
var request = new GetMonitorRequest { ServiceName = "ProductService" };
var response = await client.GetSingleMonitorDataAsync(request, headers);
五、总结:gRPC最佳实践
-
版本兼容
永远不要修改PB字段的编号,只能新增字段;
新增字段要设置默认值(proto3自动设置,比如int32默认0,string默认空);
旧客户端会忽略新增的字段,新客户端会用默认值读取旧服务器的字段。 -
性能优化
生产环境必须用TLS加密,保证数据安全;
启用压缩:gRPC支持gzip压缩,减少传输体积;
配置连接池:客户端用Channel的连接池,避免频繁创建连接;
设置超时:每个RPC方法设置超时,避免客户端一直等待。 -
拦截器使用
拦截器要尽量轻量,不要在拦截器中执行耗时操作;
服务器拦截器顺序:先认证,再日志,最后限流;
客户端拦截器顺序:先日志,再认证,最后重试。 -
适用场景
适合:微服务内部通信、实时系统(监控、聊天)、大数据传输;
不适合:浏览器直接调用(需要gRPC Web)、简单的HTTP请求(用REST更简单)。
下一节我们会学习gRPC的生态扩展:gRPC Gateway(把gRPC转成REST API)、gRPC Web(浏览器调用gRPC),让gRPC的适用场景更广泛。
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49555.html










