从RPC到gRPC:了解远程过程调用、Protocol Buffers以及现代分布式系统的通信机制
任何应用程序在某些时候都需要与其他系统进行交互。移动应用会与后端服务进行通信;后端服务又会与支付网关对接;认证服务需要与用户服务进行交互;数据传输流程则要与存储系统相连。 问题不在于系统是否需要相互沟通,而在于应该如何实现这种沟通。 多年来,基于HTTP和JSON的REST架构一直是人们的首选方案。它运行稳定、使用简单,而且相关的开发工具也随处可见。然而,随着系统规模的扩大——无论是参与交互的服务数量增加、数据交换量增大,还是对实时通信的需求提高——REST架构逐渐暴露出了其局限性。 这时,远程过程调用、Protocol Buffers以及gRPC这些技术应运而生了。 通过这本手册,你将了解什
任何应用程序在某些时候都需要与其他系统进行交互。移动应用会与后端服务进行通信;后端服务又会与支付网关对接;认证服务需要与用户服务进行交互;数据传输流程则要与存储系统相连。
问题不在于系统是否需要相互沟通,而在于应该如何实现这种沟通。
多年来,基于HTTP和JSON的REST架构一直是人们的首选方案。它运行稳定、使用简单,而且相关的开发工具也随处可见。然而,随着系统规模的扩大——无论是参与交互的服务数量增加、数据交换量增大,还是对实时通信的需求提高——REST架构逐渐暴露出了其局限性。
这时,远程过程调用、Protocol Buffers以及gRPC这些技术应运而生了。
通过这本手册,你将了解什么是远程过程调用以及它被设计出来是为了解决什么问题;同时也会理解Protocol Buffers的含义、存在理由及其工作原理。
接下来,你会看到谷歌是如何将这些技术结合在一起,从而创造出gRPC的——这是现代分布式系统中最为强大的通信框架之一。你将学习gRPC的四种通信模式,会看到如何通过一个契约文件生成多种语言版本的代码,还会了解如何使用Flutter实现一个包含认证、错误处理和超时机制等生产级功能的完整系统。
读完这本书后,你不仅会知道gRPC是什么,还会明白在什么情况下应该使用它,在什么情况下不适合使用它,以及作为系统工程师应该如何思考服务之间的通信问题。
目录
什么是远程过程调用?
要理解远程过程调用,首先需要了解什么是过程调用。
过程调用指的是调用某个程序或函数,使其代码得以执行。例如,在Dart中:
double calculateTax(double amount) {
return amount * 0.075;
}
final tax = calculateTax(50000); // 这是一个本地过程调用
你调用`calculateTax`函数,传递一个参数,然后得到结果。这个函数位于同一台机器上、同一个进程中、同一个内存空间中。这就是本地过程调用。
而远程过程调用则将这种机制应用到网络环境中。被调用的函数可能位于另一台机器上、另一个进程中,甚至位于另一个国家。但从调用者的角度来看,这种操作感觉与调用本地函数完全一样。
// 这看起来像是一个本地函数调用
final tax = await taxService.calculateTax(amount: 50000);
// 但实际上,其内部流程如下:
// 1. 将参数序列化为二进制格式
// 2. 通过网络连接发送到远程服务器
// 3> 服务器使用这些参数执行calculateTax函数
// 4> 将结果序列化回来
// 5> 通过网络将结果传回
// 6> 将序列化后的结果反序列化为Dart对象
// 7> 最后将结果返回给调用者,就像它是本地调用的结果一样
所有的网络相关操作都被隐藏了起来。你只需调用一个函数,就能得到结果,而中间的所有细节都由远程过程调用框架来处理。
这就是远程过程调用的核心理念:让调用远程函数变得与调用本地函数一样自然。
远程过程调用解决的问题
要理解为什么远程过程调用如此重要,就需要了解没有它时会遇到哪些问题。
如果没有远程过程调用,调用远程服务的方式就会是这样的:
Future calculateTax(double amount) async {
final response = await http.post(
Uri.parse('https://tax-service.internal/api/v1/calculate'),
headers: {
'Content-Type': 'application/json',
'Authorization': 'Bearer $token',
},
body: jsonEncode({'amount': amount}),
);
if (response.statusCode != 200) {
throw Exception('Tax calculation failed: ${response.statusCode}`);
}
final data = jsonDecode(response.body);
return (data['tax'] as num).toDouble();
}
每次调用远程服务时,你都需要完成以下步骤:
知道端点URL的地址,并将其硬编码到代码中
确定应该使用哪种HTTP方法
手动将请求数据序列化为JSON格式
自行处理HTTP状态码
手动将服务器返回的响应数据反序列化为JSON格式
将动态类型的数据转换为你期望得到的具体类型
确保响应中的字段名称与你的预期一致
使用RPC时:
这就像调用一个本地函数一样简单:
final tax = await taxService.calculateTax(
TaxRequest(amount: 50000),
);
// “tax”变量现在就是一个类型明确的TaxResponse对象
// 不需要URL,不需要使用HTTP方法,也不需要进行JSON解析或类型转换。
框架会处理所有相关细节。函数签名被定义在一份契约文件中,客户端和服务器都会使用这份文件。类型检查会在编译阶段进行;如果服务器更改了响应数据的结构,客户端在代码投入生产环境之前就会无法正常编译。
RPC的作用就在于消除网络通信中可能出现的复杂性,让你能够专注于真正需要完成的任务。
为什么选择gRPC而非REST:实际应用案例
在介绍Protocol Buffers和gRPC的技术细节之前,首先有必要明确在特定场景下选择gRPC而非REST的原因。这并不是说gRPC总是更优秀,而是要清楚地了解它在哪些方面确实具有优势。
许多内部系统需要处理大型数据包
以一家大型电信企业的内部API为例:对于计划注册或激活流程而言,一个请求和响应数据包中可能包含上千个字段。而会有数十个内部应用程序调用这个接口,包括计费系统、CRM平台、面向客户的移动应用、内部监控工具以及合作伙伴门户等。
使用REST和JSON技术时,所有这些应用程序都会将包含上千个字段的数据以文本形式进行传输。subscription_activation_status、rate_plan_identifier、network_provisioning_reference这类字段名称在每次请求中都会作为字符串被发送出去。实际上,数据包中很大一部分内容并不是真正的数据,而是用于描述这些数据的标签。
而使用Protocol Buffers时,字段名称根本不会出现在数据包中,只有字段编号和实际值会被传输。因此,那个包含上千个字段的数据包体积会大幅减小。对于每天被数十个系统调用数百万次的接口来说,这种带宽节省的效果是非常显著的,直接体现在基础设施成本的降低上。
除了数据量之外,生成的客户端代码的质量也同样重要。使用REST时,每个应用程序都需要自行阅读API文档来理解契约内容;当后端修改了字段名称或类型时,并不是所有应用程序都能立即察觉到这些变化,有些应用程序直到出现故障时才会发现问题。
而使用共享的.proto文件时,所有应用程序都可以从同一个源代码生成类型明确的客户端代码。一旦契约内容发生变化,所有应用程序都需要重新生成客户端代码。编译器会立即指出哪些地方发生了变动,从而确保没有任何有问题的代码被投入生产环境。
低带宽及远程网络环境
在网络质量波动较大的市场中,这是gRPC最容易被忽视的优势之一。
在许多地区,大量移动用户仍然使用2G或3G网络。在2G连接下,带宽可能仅为50到100千比特每秒。对于大小为80千字节的REST JSON响应,通过2G网络下载需要超过6秒钟;而同样大小的protobuf二进制文件,其传输时间则不到2秒钟——后者의大小通常只有前者的3到10分之一。
这种差异绝非技术上的琐事,而是决定一个应用程序是否对用户来说可用、是否正常运行的关键因素。对于任何在网络环境不稳定的市场中运行的产品而言,protobuf的二进制编码效率无疑构成了其重要的竞争优势。
除了文件大小之外,gRPC还通过HTTP/2协议进行通信,该协议会维持一个持久的连接,而不会为每个请求都建立新的连接。在网络速度较慢的环境中,由于建立连接所需的时间可能长达数百毫秒,因此重复使用同一个连接可以显著节省传输时间。
微服务之间的通信
当两个内部服务需要相互通信时,有多种选择。如果不需要立即得到响应,操作可以异步进行,或者需要将某些信息发送给多个接收方,那么使用Kafka或RabbitMQ这样的消息队列系统会非常合适。
然而,许多服务之间的交互本质上是同步的。例如,认证服务必须在请求处理之前验证令牌;欺诈检测服务需要在支付授权之前评估交易情况;定价服务也必须在生成报价之前计算出价格。这类操作无法等待事件发生后再进行处理。
对于高频进行的同步服务间通信来说,使用gRPC通过HTTP/2协议并结合protobuf编码的方式,其效率远高于REST。持久的复用连接避免了每次请求时都进行连接建立的操作;二进制编码方式也消除了在数据传输过程中进行JSON序列化与反序列化的开销;此外,由于两个服务都是基于相同的协议接口进行开发的,因此编译过程也会更加统一。
当大规模应用中两个服务每秒要相互调用数千次时,这些效率差异就会转化为实际的性能和成本差距。
跨多个团队管理API接口协议
在大型工程组织中,通常会有多个团队开发其他团队所依赖的服务。REST API的接口协议往往仅保存在文档中,但这类文档很容易过时。如果后端团队修改了某个字段的名称,移动端团队可能要等到用户反馈程序崩溃时才会发现这一变化;数据团队也可能要在凌晨2点的时候,因为数据处理流程出现错误才意识到问题。
而gRPC采用的protobuf仓库管理机制,将接口协议的管理从文档依赖转变为代码驱动的方式。任何对接口协议的修改都会通过拉取请求的形式提交,所有相关团队都会对其进行审核。这种机制能够确保在编译阶段就发现那些可能引发问题的修改,从而避免在生产环境中出现意外情况。
这种治理机制带来的好处会随着团队规模的扩大而增强。组织规模越大,这种机制的价值也就越明显。
实时通信
REST是一种请求-响应模型。客户端发出请求,服务器给出响应,对话就此结束。而对于需要实现实时功能的情况,要么采用轮询的方式(这种方式效率较低),要么在REST API之外另行搭建WebSocket服务器来处理实时通信(这样就需要维护两个不同的系统)。
gRPC的流式通信机制能够让你在用于普通请求的同一框架内自然地实现实时通信功能。无论是实时的余额更新、交易通知,还是双向聊天会话,所有这些功能都使用相同的客户端代码、相同的连接方式,以及相同的protobuf编码格式,与普通的请求处理过程完全一致。
一个框架就能满足所有的通信需求,无需额外搭建任何基础设施。
Protocol Buffers:一种用于数据交换的新技术
RPC本身只是一个概念。要实现它,你需要两样东西:一是定义客户端与服务器之间交互规则的方法,二是能够高效地将数据序列化以便在网络上传输的技术。
而Protocol Buffers正是为了解决这些问题而诞生的。
Protocol Buffers,通常被称为protobuf,是一种与具体语言和平台无关、且具有高度扩展性的数据序列化技术。它由谷歌在2001年开发出来,内部使用多年后于2008年被开源。
大规模应用中JSON带来的问题
对于Web API来说,JSON是目前最主流的数据格式。它易于人类阅读、灵活性强,而且得到了广泛的支持,在很多场景下都是理想的选择。
然而,当数据量达到一定规模时,JSON的结构缺陷就会暴露出来,从而影响系统的性能。
以用户信息响应为例:
{
"id": "usr_001",
"first_name": "John",
"last_name": "Smith",
"email": "john@example.com",
"phone_number": "+2348012345678",
"account_type": "savings",
"balance": 500000.00,
"currency": "NGN",
"is_verified": true,
"is_active": true,
"kyc_level": 3,
"created_at": "2024-01-15T10:30:00Z",
"last_login": "2026-07-20T09:15:00Z"
}
在每次数据传输过程中,所有的字段名称都会以字符串的形式被发送到网络中。"first_name"、"account_type"、"phone_number"这些其实并不是数据本身,而是数据的标签,但它们仍然会占用一定的网络带宽。
现在假设有一个企业内部API,其请求和响应的数据中包含了超过一千个字段,而且每天会有数十个内部应用程序数千次地调用这个API。那么,在这些数据中,有很大一部分实际上都是字段名称的字符串,并非真正有用的数据。这种结构上的缺陷会导致大量的网络带宽被浪费,同时也会增加数据处理的成本。
除了数据量大的问题之外,JSON还有一个更大的缺点:它在网络层面上并没有固定的数据结构规范。因此,后端开发人员可以在新的部署环境中将"first_name"字段的名称改成"firstName",而这样的变更在用户实际使用系统中并不会引发任何问题。但这种随意性会导致系统运行不稳定。
Protocol Buffers的不同之处
Protocol Buffers采用一种完全不同的数据编码方式,从而解决了上述两个问题。
它不会将数据编码成包含字段名称的、人类可读的文本形式,而是仅使用字段编号和数值将其编码为紧凑的二进制格式。在传输过程中,字段名称根本不会被发送到网络中。
以下是用Protocol Buffers定义的相同用户信息示例:
message UserProfile {
string id = 1;
string first_name = 2;
string last_name = 3;
string email = 4;
string phone_number = 5;
string account_type = 6;
double balance = 7;
string currency = 8;
bool isVerified = 9;
bool is_active = 10;
int32 kyc_level = 11;
string created_at = 12;
string last_login = 13;
}
当Protocol Buffers对这些数据进行处理时,生成的输出是人类无法理解的二进制格式。但对于机器而言,这种格式极其紧凑且便于解析。字段编号(如1、2、3等)用于标识各个字段,而字段名称在编码后的结果中根本不会出现。
这样一来,相同用户信息的JSON格式数据大小约为280字节,而用Protocol Buffers编码后仅为95字节——体积减少了三倍之多。对于包含上千个字段的企业级数据来说,这种差异显得尤为显著。
由于数据结构是通过`.proto`文件定义的,客户端和服务器在编译时都会依据这个文件进行代码生成,因此如果字段名称发生变化,错误会在编译阶段就被发现,而不会在运行时才被暴露出来。
Proto文件
这种文件的编写语言是Protocol Buffer Language(即proto3),而不是Go、Dart、Python或Java。你可以使用任何文本编辑器来编写它。安装了`vscode-proto3`扩展的VS Code可以提供语法高亮显示、自动完成功能以及实时验证功能。
以下是一个金融科技平台使用的完整`.proto`文件示例:
syntax = "proto3";
package banking;
option go_package = "./banking";
option java_package = "com.fintech.banking";
service BankingService {
// 单向请求:一次请求对应一个响应
rpc Login (LoginRequest) returns (LoginResponse);
// 单向请求:获取用户信息
rpc GetProfile (ProfileRequest) returns (UserProfile);
// 服务器端流式传输:实时更新余额信息
rpc WatchBalance (BalanceRequest) returns (stream BalanceResponse);
// 服务器端流式传输:实时推送交易记录
rpc StreamTransactions (TransactionRequest) returns (stream Transaction);
// 客户端流式传输:分块上传身份验证文件
rpc UploadDocument (stream DocumentChunk) returns (UploadResponse);
// 双向流式传输:支持实时聊天
rpc Chat (stream ChatMessage) returns (stream ChatMessage);
}
message LoginRequest {
string email = 1;
string password = 2;
}
message LoginResponse {
string token = 1;
string user_id = 2;
int64 expires_at = 3;
}
message ProfileRequest {
string user_id = 1;
}
message UserProfile {
string id = 1;
string first_name = 2;
string last_name = 3;
string email = 4;
string phone_number = 5;
string account_type = 6;
double balance = 7;
string currency = 8;
bool isVerified = 9;
int32 kyc_level = 10;
}
message BalanceRequest {
string user_id = 1;
}
message BalanceResponse {
double balance = 1;
string currency = 2;
int64 timestamp = 3;
}
message TransactionRequest {
string user_id = 1;
int32 limit = 2;
}
message Transaction {
string id = 1;
double amount = 2;
string description = 3;
string type = 4;
int64 timestamp = 5;
}
message DocumentChunk {
bytes data = 1;
string document_type = 2;
int32 chunk_index = 3;
bool is_last = 4;
}
message UploadResponse {
bool success = 1;
string document_id = 2;
string message = 3;
}
message ChatMessage {
string sender_id = 1;
string content = 2;
int64 timestamp = 3;
}
让我们仔细地分析这个文件中的每一个部分。
语法声明
syntax = "proto3";
这一行代码告诉protobuf编译器你正在使用哪种版本的Protocol Buffer语言。目前的标准版本是proto3,因此它必须是每个.proto文件中第一行非注释内容。
包声明
package banking;
包名的存在可以避免在不同服务中使用多个.proto文件时出现命名冲突。它的作用类似于命名空间——如果两个服务都定义了名为UserProfile的消息类型,那么包名就能区分它们:banking.UserProfile与auth.UserProfile就是不同的消息类型。
语言特定选项
option go_package = "./banking";
option java_package = "com.fintech.banking";
这些选项用于告诉编译器如何为特定的编程语言组织生成后的代码。它们不会影响原始的.proto文件本身,只会影响生成的代码。
服务定义
service BankingService {
rpc Login (LoginRequest) returns (LoginResponse);
rpc WatchBalance (BalanceRequest) returns (stream BalanceResponse);
}
service块用于定义RPC契约。可以将其理解为任何面向对象语言中的抽象类——它规定了哪些函数存在、这些函数接受什么参数以及返回什么结果。
每一条rpc语句都定义了一个远程过程。当类型前加上stream关键字时,表示会发送多条消息而不仅仅是一条。
消息定义
message LoginRequest {
string email = 1;
string password = 2;
}
message是一种数据结构,可以将其视为一个只有字段而没有方法或逻辑的类。每个字段由三部分组成:
类型可以是string、int32、int64、double、bool、bytes,或者另一种消息类型。
名称是该字段在生成后的代码中的显示形式,这只是为了便于人类阅读,并不会出现在二进制编码中。
字段编号是protobuf在二进制输出中使用的唯一标识符,而不是字段的名称。这一点非常重要:一旦为某个字段指定了编号,就绝对不能更改或重复使用这个编号。二进制编码过程中使用的是这些编号,而不是字段名称;如果更改了字段编号,之前编码好的数据就会变得无法读取。
你可以安全地添加新的字段并为其指定新的编号,也可以删除字段(被删除的字段编号会保持保留状态,不能再被重复使用),或者重新命名字段(不过字段名称在二进制编码中并不会出现)。但是绝对不能更改字段编号,也不能重复使用已删除字段的编号,更不能修改字段的类型。
JSON与Protocol Buffers的对比
既然你已经了解了这两种格式,那么我们就来直接进行比较吧。
大小对比
让我们来看一下用这两种格式表示的同一个登录请求:
JSON(文本格式):
{
"email": "john@example.com",
"password": "securepassword123"
}
大约55字节。
Protocol Buffers二进制格式:
字段1(email):标签 + 长度 + 值的字节长度;字段2(password)也是如此。
大约38字节。
对于一个只包含两个字段的简单消息来说,这两种格式之间的差异并不明显。但如果是包含上千个字段的企业级数据,情况就不同了:在JSON中,仅字段名称所占用的空间就会占到总数据大小的40%到60%;而在Protocol Buffers中,字段名称根本不会增加数据的体积。
在带宽仅为50千比特每秒的2G连接环境中,80字节的JSON响应与15字节的Protocol Buffers响应之间的差距,就相当于加载速度相差13秒与2秒。对于那些网络基础设施较为落后的地区来说,这种性能差异其实意味着用户体验上的巨大差别。
速度对比
Protocol Buffers的序列化与反序列化过程要比JSON解析快得多,因为二进制数据不需要进行字符串的分词处理、引号的处理、空白字符的忽略,也不需要推断数据类型。解析器只需读取字段编号、值的数据类型以及具体的数值,然后直接进入下一个字段——整个过程完全是二进制的操作。
而JSON解析则必须逐个字符地处理字符串,通过周围的引号和分隔符来识别键值对,从数据的格式中推断出数据类型,最后再根据这些信息构建对象。
对于那些在每次会话中需要处理大量响应信息的移动设备来说,这种解析速度上的差异会直接体现在CPU使用率和电池消耗上。
架构与类型安全性
JSON在网络层并不具备任何格式约束机制。后端可以将"balance"字段的名称改为"current_balance",而客户端只有在应用程序在实际环境中崩溃时才会发现这一变化。
Protocol Buffers的架构则在编译阶段就会被强制执行。如果.proto文件发生了任何会影响客户端正常运行的更改,客户端在尝试编译时就会失败,这样问题就能在到达任何用户之前就被及时发现并解决。
公正的对比
| JSON | Protocol Buffers | |
|---|---|---|
| 编码方式 | 文本格式(UTF-8) | 二进制格式 |
| 是否便于人类阅读 | 是 | 否 |
| 数据大小 | 较大(包含字段名称) | 小3到10倍 |
| 解析速度 | 较慢(需要文本分词处理) | 较快(直接进行二进制读取) |
| 格式约束机制 | 网络层没有约束 | 编译时强制执行 |
| 代码生成方式 | 可选 | 必选且自动完成 |
| 适用场景 | 公共API、人工审核场景 | 内部服务、高性能应用场景 |
Protocol编译器与代码生成
protoc编译器会读取你的.proto文件,并根据你指定的语言生成相应的代码。正是这一机制使得“通用契约”成为现实。
生成Go语言代码(用于后端服务器):
protoc \
--go_out=. \
--go-grpc_out=. \
proto/banking.proto
生成的文件包括:
banking.pb.go <- 消息结构体
banking_grpc.pb.go <- 服务器接口代码
生成Dart语言代码(用于Flutter客户端):
protoc \
--dart_out=grpc:lib/generated \
proto/banking.proto
生成的文件包括:
lib-generated/
banking.pb.dart <- 消息类代码
banking_pb grpc.dart <- 客户端代理代码
生成Python语言代码(用于数据服务):
protoc \
--python_out=. \
--grpc_python_out=. \
proto/banking.proto
生成的文件包括:
banking_pb2.py <- 消息类代码
banking_pb2_grpc.py 客户端与服务器类代码
生成TypeScript语言代码(用于Web前端):
protoc \
--ts_out=. \
proto/banking.proto
生成的文件包括:
banking.ts <- 带类型定义的消息类及客户端代码
所有这些代码都是从同一个banking.proto文件中自动生成的。
Go后端工程师无需编写序列化代码,Flutter前端工程师也无需编写反序列化代码;Python数据工程师不必手动解析二进制数据,TypeScript Web开发人员更不需要手动构建HTTP请求。所有这些工作都由团队共同约定的“契约”自动完成。
下面是各种语言生成的代码示例,以便大家更直观地了解这一过程:
生成的Go语言服务器接口代码(后端实现部分):
// 生成代码,切勿修改
type BankingServiceServer interface {
Login(context.Context, *LoginRequest) (*LoginResponse, error)
WatchBalance(*BalanceRequest, BankingService_WatchBalanceServer) error
mustEmbedUnimplementedBankingServiceServer()
}
// Go后端工程师负责编写实现代码
type bankingServer struct {
pb.UnimplementedBankingServiceServer
}
func (s *bankingServer) Login(
ctx context.Context,
req *pb.LoginRequest,
) (*pbLoginResponse, error) {
token, err := authService.login(req.Email, req.Password)
if err != nil {
return nil, status.Errorf(codes.Unauthenticated, "无效的凭证")
}
return &pb LoginResponse{
Token: token,
UserId: user.Id,
}, nil
}
生成的Python客户端(数据团队使用这个版本):
import grpc
import banking_pb2
import banking_pb2_grpc
channel = grpc.secure_channel(
'api.fintech-platform.com:50051',
grpc.ssl_channel_credentials()
)
stub = banking_pb2_grpc.BankingServiceStub(channel)
response = stub.login(banking_pb2 LoginRequest(
email='john@example.com',
password='password123'
))
print(f"令牌:{response.token}")
print(f"用户ID:{response.user_id}")
生成的Dart客户端(你在Flutter项目中使用这个版本):
import 'package:grpc/grpc.dart';
import 'generated/banking.pb grpc.dart';
import 'generated/banking.pb.dart';
final channel = ClientChannel('api.fintech-platform.com', port: 50051);
final client = BankingServiceClient(channel);
final response = await client.login(
LoginRequest(email: 'john@example.com', password: 'password123'),
);
print('令牌:${response.token}`);
print('用户ID:${response.userId'));
三种不同的语言,三个不同的团队,但它们都是基于同一个`.proto`文件生成的。这些客户端都是强类型化的,并且能够确保与服务器端的数据保持同步。
什么是gRPC?
gRPC是谷歌开发的开源远程过程调用框架。它于2016年被开源,如今已经成为Cloud Native Computing Foundation(CNCF)认证的项目。这意味着它在云原生生态系统中已经经过了最高级别的生产环境验证。
gRPC结合了以下三种技术:
- 远程过程调用作为编程模型:可以像调用本地函数一样调用远程函数。
- Protocol Buffers作为接口定义语言和数据序列化格式:它支持强类型契约,并采用紧凑的二进制编码方式。
- HTTP/2作为传输协议:HTTP/2支持多路复用、持久连接,同时还采用了二进制帧格式进行数据传输。
这三项技术的结合使得gRPC成为一种比REST更快、结构更清晰、功能也更强大的通信框架。
在谷歌内部,几乎所有的服务间通信都使用gRPC。Netflix、Uber、Square、Dropbox、Lyft等数百家机构也都将gRPC用于它们的微服务通信系统。目前,Go、Java、Python、C++、C#、Ruby、Node.js、PHP、Dart、Kotlin等多种语言都支持使用gRPC。
为什么HTTP/2对gRPC如此重要
gRPC是完全基于HTTP/2协议构建的。了解HTTP/2所提供的功能对于理解gRPC的工作原理至关重要。
目前大多数REST API所使用的HTTP/1.1存在一些根本性的性能限制:在同一连接上,必须先完成一个请求后才能开始下一个请求;所有请求头信息都会以文本形式发送;服务器也不能主动向客户端发送数据,除非客户端首先发起请求。
而HTTP/2正是为了解决这些限制而设计的。
多路复用
HTTP/2在单个连接中实现了数据流的分割传输。多个独立的请求可以同时通过同一个TCP连接进行发送。
使用单个TCP连接访问api.fintech-platform.com
数据流1:登录请求 ---------> 登录响应
数据流2:个人资料请求 ------->> 个人资料响应
数据流3:余额查询请求 ------->> 余额更新信息(持续发送中)
数据流4:交易记录请求 --> 交易详情信息(持续发送中)
这四个数据流同时通过同一个连接进行传输
在HTTP/1.1中,需要建立四个独立的连接,或者等待前一个请求完成后再开始下一个请求。而HTTP/2则可以通过一个持久的连接同时处理所有这些请求,无需等待。
这就是gRPC具备流式传输功能的基础所在。正是这种持久的多路复用连接机制,使得服务器能够在客户端继续发送其他请求的同时,持续向客户端推送余额更新信息及交易通知。
二进制帧格式
HTTP/1.1将所有数据都以文本形式进行传输;而HTTP/2则将所有数据打包成二进制帧进行传输。二进制格式更加紧凑,机器解析起来也速度更快。
每个gRPC消息都会被分解成二进制帧,然后通过HTTP/2连接发送出去。结合protobuf的二进制编码技术,gRPC的数据在每一层传输时都采用最为紧凑的形式。
头部信息压缩(HPACK)
HTTP/1.1会在每个请求中都发送完整的头部信息。例如,用于携带JWT令牌的Authorization头部信息,其大小可能达到500字节甚至更多,并且会在每次请求中都被重复发送。
HTTP/2采用了HPACK压缩技术。之前请求中已经发送过的头部信息会被缓存起来,后续的请求只需发送那些发生了变化的部分即可。一旦Authorization头部信息被发送过,之后就会通过一个简短的索引来引用它,而无需再次完整地传输该头部信息。
对于那些在每次会话中都要发送大量经过身份验证的请求的移动应用来说,这种压缩机制能够显著节省带宽,尤其是在网络速度较慢的情况下,每一字节的数据都显得非常宝贵。
服务器主动推送
HTTP/2允许服务器在不需要客户端发起请求的情况下,主动向客户端发送数据。客户端打开一个数据流后,服务器就会在有新数据需要发送时,通过这个数据流将它们推送给客户端。
这就是gRPC服务器流式传输机制的运作原理。客户端只需发送一次WatchBalance请求,服务器就会在余额发生变化时立即推送新的BalanceResponse响应给客户端。无需进行轮询或重复请求,连接会始终保持开放状态,只要服务器有新内容要发送,就会立刻将其推送给客户端。
gRPC的四种通信模式
这是本文最重要的部分。实际上,gRPC并不只有一种通信模式,而是共有四种不同的通信方式。每种通信模式的具体实现都在.proto文件中进行了明确规定,它们各自适用于不同的使用场景。
模式1:单例请求/响应
客户端发送一个请求,服务器返回多个响应。就请求与响应的流程而言,这种模式与REST API调用是完全相同的。
rpc Login (LoginRequest) returns (LoginResponse);
客户端 ---- 发送登录请求 ----> 服务器
客户端 ---- 接收登录响应 ----> 服务器
操作完成。
何时使用单向RPC: 登录、获取个人资料信息、发起支付请求、创建数据、检索配置信息——任何遵循“询问-响应”模式的操作都适合使用单向RPC。
Future login(String email, String password) async {
try {
return await _client.login(
LoginRequest(email: email, password: password),
);
} on GrpcError catch (e) {
throw _mapGrpcError(e);
}
}
func (s *bankingServer) Login(
ctx context.Context,
req *pb.LoginRequest,
) (*pbLoginResponse, error) {
user, err := s.authService.login(req.Email, req.Password)
if err != nil {
return nil, status.Errorf(codes.Unauthenticated, "无效的凭证:%v", err)
}
token, _ := s.tokenService.Generate(user.Id)
return &pbLoginResponse{
Token: token,
UserId: user.Id,
}, nil
}
模式2:服务器端流式RPC
客户端发送一个请求,服务器会连续不断地返回多个响应。连接会保持开启状态,服务器会在数据准备好时立即将其推送给客户端。
rpc WatchBalance (BalanceRequest) returns (stream BalanceResponse);
rpc StreamTransactions (TransactionRequest) returns (stream Transaction);
客户端 ---- 发送余额查询请求 ----> 服务器
客户端 ---- 接收余额更新信息 ----> 服务器(余额:500000)
客户端 ---- 再次接收余额更新信息 ----> 服务器(余额:495000,因为发生了扣款操作)
客户端 ---- 最后接收余额更新信息 ----> 服务器(余额:995000,因为发生了存款操作)
[连接保持开启状态,服务器会随时推送最新的数据变化]
何时使用服务器端流式RPC: 账户实时余额、交易通知、股票实时价格、体育比赛比分、新闻推送、系统监控面板——任何需要服务器持续更新数据的场景都适合使用这种模式。
Stream watchBalance(String userId) {
return _client.watchBalance(
BalanceRequest(userId: userId),
);
}
在Flutter中,可以使用StreamBuilder来处理这类数据流:
StreamBuilder(
stream: _dataSource.watchBalance(currentUserId),
builder: (context, snapshot) {
if (snapshot.connectionState == ConnectionState.waiting) {
return const CircularProgressIndicator();
}
if (snapshot.hasError) {
return Text('错误:${snapshot.error}`);
}
if (!snapshot.hasData) {
return const Text('正在等待余额数据……');
}
final balance = snapshot.data!;
return Column(
children: [
Text(${balance.currency} ${balance.balance.toStringAsFixed(2)}'),
Text('更新时间:${DateTime.fromMillisecondsSinceEpoch(balance.timestampToInt())}',
],
);
},
) 每当服务器发送新的余额信息时,StreamBuilder会再次调用builder,于是小部件就会显示更新后的数值。整个过程中不存在任何轮询机制或手动刷新操作——完全由服务器主动发送数据,小部件则被动接收并展示这些数据。
Go语言服务器实现:
func (s *bankingServer) WatchBalance(
req *pb.BalanceRequest,
stream pb.BankingService_WatchBalanceServer,
) error {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case <-stream.Context().Done():
return nil
case <-ticker.C:
balance, err := s.accountService.GetBalance(reqUserId)
if err != nil {
return status.Errorf(codes.Internal, "failed to fetch balance: %v", err)
}
if err := stream.Send(&pb.BalanceResponse{
Balance: balance.Amount,
Currency: balance.Currency,
Timestamp: time.Now().UnixMilli(),
}); err != nil {
return err
}
}
}
}
模式3:客户端流式RPC通信
客户端向服务器发送一系列消息,服务器会逐一处理这些消息,最后再统一给出响应。
rpc UploadDocument (stream DocumentChunk) returns (UploadResponse);
客户端 ---- 数据块1(字节0-1024)-----> 服务器
客户端 ---- 数据块2(字节1024-2048)--> 服务器
客户端 ---- 数据块3(字节2048-3072)--> 服务器
客户端 ---- 数据块4(最后一个数据块)-------> 服务器
客户端 <---收到UploadResponse-------------- 服务器(文档ID:「doc_001」)
何时使用客户端流式通信:当需要分块上传大文件(如身份验证文件、个人资料照片),或者批量发送传感器检测数据,又将大量记录提交给服务器时,这种通信方式非常适用。
Dart语言实现:
Future<UploadResponse>> uploadDocument(
List<Uint8List>> chunks,
String documentType,
) async {
try {
Stream<DocumentChunk>> chunkStream() async* {
for (int i = 0; i < chunks.length; i++) {
yield DocumentChunk(
data: chunks[i],
documentType: documentType,
chunkIndex: i,
isLast: i == chunks.length - 1,
);
}
}
return await _client.uploadDocument(chunkStream());
} on GrpcError catch (e) {
throw _mapGrpcError(e);
}
}
chunkStream()是一个异步生成器函数。async*关键字表示该函数会逐步产生多个值,而不是直接返回一个值。每次调用yield时,都会生成一个DocumentChunk消息,由gRPC发送给服务器。服务器会依次接收这些消息并全部处理完毕,最后才返回一个UploadResponse作为最终结果。
Go语言服务器实现:
func (s *bankingServer) UploadDocument(
stream pb.BankingService_UploadDocumentServer,
) error {
var allData []byte
var documentType string
for {
chunk, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return status.Errorf(codes.Internal, "failed to receive chunk: %v", err)
}
allData = append(allData, chunk.Data...)
documentType = chunk.DocumentType
}
docId, err := s.documentService.Store(allData, documentType)
if err != nil {
return status.Errorf(codes INTERNAL, "failed to store document: %v", err)
}
return stream.SendAndClose(&pbUploadResponse{
Success: true,
DocumentId: docId,
Message: "Document uploaded successfully",
})
}
模式4:双向流式RPC
客户端和服务器会同时发送消息,任何一方都可以随时发送数据,另一方不需要等待对方。
rpc Chat (stream ChatMessage) returns (stream ChatMessage);
客户端 ----"Hello"---------------------------> 服务器
服务器 <---"Hi, how can I help?"-------------- 客户端
客户端 ----"What is my account balance?"-----> 服务器
服务器 <---"Your balance is NGN 500,000"------ 客户端
服务器 <---"New transaction alert: -5,000"---- 客户端(服务器主动发起通信)
客户端 ----"Thanks"--------------------------> 服务器
[双方可以自由且同时进行交流]
何时使用双向流式通信: 实时聊天、实时协作编辑文档、多人游戏状态同步、交互式交易终端、实时客户支持服务。
Dart语言实现示例:
void startChat(String userId) {
final outgoing = StreamController<ChatMessage>>();
final incoming = _client.chat(outgoing.stream);
incoming.listen(
(message) {
print '${message.senderId}: ${message.content}`);
},
onError: (error) {
print('聊天错误: $error');
},
onDone: () {
print('聊天会话结束');
},
);
outgoing.add(ChatMessage(
senderId: userId,
content: 'Hello, I need help with my account',
timestamp: DateTime.now().millisecondsSinceEpoch,
));
}
StreamController负责管理发送方的消息流。每当用户发送消息时,就可以将这些消息添加到outgoing中。而incoming流则会接收来自服务器的消息。这两条消息流会通过同一个HTTP/2连接同时传输。
func (s *bankingServer) Chat(stream pb.BankingService_ChatServer) error {
for {
msg, err := stream.Recv()
if err == io.EOF {
return nil
}
if err != nil {
return err
}
response := s.chatService.Process(msg)
if err := stream.Send(&pb.ChatMessage{
SenderId: "support_agent",
Content: response,
Timestamp: time.Now().UnixMilli(),
}); err != nil {
return err
}
}
}
protobuf仓库:组织级的最佳实践
在小型项目中,.proto文件可以直接保存在后端代码库中。移动开发人员可以通过克隆后端代码库来获取这些文件,这种做法在规模较小的项目中确实可行。
但对于任何规模较大的组织来说,这种做法就会出现问题。后端代码库会成为“事实的权威”,从而使后端团队对协议内容拥有单方面的控制权;其他团队只有在他们的开发环境出现故障时,才会发现这些变更。
业界公认的最佳实践是专门设立一个protobuf仓库:这个仓库属于所有团队,但没有任何一个团队能够单独拥有它。
fintech-api-contracts/
proto/
auth/
auth.proto
banking/
banking.proto
payments/
payments.proto
notifications/
notifications.proto
kyc/
kyc.proto
scripts/
generate_dart.sh
generate_go.sh
generate_python.sh
README.md
合同变更的运作机制
任何API的变更都会遵循相同的流程:
开发人员提出对banking.proto文件的修改建议 | 在fintech-api-contracts仓库中提交Pull Request | >Flutter团队的负责人进行审核: “这个修改会不会影响客户端的功能?我们需要进行更新吗?” >Go后端团队的负责人进行审核: “这个修改是否可行?是否符合我们的开发规范?” >React Web团队的负责人进行审核: “Web客户端是否需要做出相应的调整?” >Python数据团队的负责人进行审核: “这个修改会不会影响我们的数据处理流程?” | >所有团队都批准后, | >Pull Request会被合并,此时该变更就正式成为标准规范 | >每个团队都会运行他们的代码生成脚本 | >如果存在会导致故障的变更,构建过程就会失败 >这样,这些变更就能在编译阶段就被发现, >在任何代码被部署到生产环境之前,问题就能得到解决。这种流程能够带来REST架构与文档所无法实现的优势:通过编译器的强制检查,确保所有团队都能使用相同版本的协议规范,从而保证合同内容在各个团队之间保持一致。
使用Dart和Flutter构建完整的gRPC系统
现在,让我们将这些知识点应用到一个实际的生产环境示例中,来看看整个系统的构建过程。
项目设置
首先,在你的Flutter项目中添加gRPC依赖:
# pubspec.yaml dependencies: flutter: sdk: flutter grpc: ^3.2.4 protobuf: ^3.1.0 dev_dependencies: protoc_plugin: ^21.1.2接下来,安装protoc编译器以及Dart的protoc插件:
# macOS brew install protobuf # 安装Dart的protoc插件 dart pub global activate protoc_plugin然后,使用protoc工具从.proto文件生成Dart代码:
protoc \ --dart_out=grpc:lib/generated \ -I proto \ proto/banking/banking.proto数据源层
// lib/features/banking/data/datasources/banking_remote_datasource.dart import 'package:grpc/grpc.dart'; import '__generated/banking.pb.dart'; import '__generated/banking.pb grpc.dart'; import '__core/error/app_exception.dart'; class BankingRemoteDataSource { late final BankingServiceClient _client; late final ClientChannel _channel; BankingRemoteDataSource({ required String host, required int port, required String authToken, }) { _channel = ClientChannel( host, port: port, options: const ChannelOptions( credentials: ChannelCredentials.secure(), connectionTimeout: Duration(seconds: 10), ), ); _client = BankingServiceClient( _channel, options: CallOptions( metadata: {'authorization': 'Bearer $authToken'}, timeout: const Duration(seconds: 30), ), ); } Futurelogin(String email, String password) async { try { return await _client.login( LoginRequest(email: email, password: password), ); } on GrpcError catch (e) { throw _mapGrpcError(e); } } Future getProfile(String userId) async { try { return await _client.getProfile( ProfileRequest(userId: userId), ); } on GrpcError catch (e) { throw _mapGrpcError(e); } } Stream watchBalance(String userId) { return _client .watchBalance(BalanceRequest(userId: userId)) .handleError((error) { if (error is GrpcError) throw _mapGrpcError(error); throw error; }); } Stream streamTransactions(String userId, {int limit = 20}) { return _client .streamTransactions( TransactionRequest(userId: userId, limit: limit), ) .handleError((error) { if (error is GrpcError) throw _mapGrpcError(error); throw error; }); } Future uploadDocument( List Uint8List chunks, String documentType, ) async { try { Stream chunkStream() async* { for (int i = 0; i < chunks.length; i++) { yield DocumentChunk( data: chunks[i], documentType: documentType, chunkIndex: i, isLast: i == chunks.length - 1, ); } } return await _client.uploadDocument(chunkStream()); } on GrpcError catch (e) { throw _mapGrpcError(e); } } ResponseStream startChat(Stream outgoing) { return _client.chat(outgoing); } Future dispose() async { await _channel.shutdown(); } AppException _mapGrpcError(GrpcError error) { switch (error.code) { case StatusCode.unauthenticated: return AppException.unauthorized( message: error.message ?? '未经授权', ); caseStatusCode.notFound: return AppException_notFound( message: error.message ?? '未找到', ); case StatusCode.deadlineExceeded: return AppException.timeout(message: '请求超时'); case StatusCode.unavailable: return AppException.serverUnavailable( message: '服务不可用', ); default: return AppException.unknown( message: error.message ?? '未知错误', ); } } } 让我们来了解一下这个数据源中的一些关键决策。
通道:
_channel = ClientChannel( host, port: port, options: const ChannelOptions( credentials: ChannelCredentials.secure(), connectionTimeout: Duration(seconds: 10), ), );通道实际上是应用程序与服务器之间建立的HTTP/2连接。你只需创建一次这个通道,之后就可以在所有的请求中重复使用它。
ChannelCredentials.secure()用于启用TLS加密功能,而connectionTimeout则可以防止应用程序在服务器无法响应时无限期地等待。通道正是gRPC能够实现高性能的关键所在。通过一个持久的通道,所有请求都可以通过同一个HTTP/2连接来发送;如果为每个请求都创建一个新的通道,那么这种性能优势就会完全丧失,其运行效率甚至会低于REST接口。
通过元数据进行身份验证:
_client = BankingServiceClient( _channel, options: CallOptions( metadata: {'authorization': 'Bearer $authToken'}, timeout: const Duration(seconds: 30), ), );gRPC使用元数据(键值对)来实现类似HTTP头部信息的功能。将认证令牌作为元数据传递给
CallOptions,意味着通过这个客户端发出的所有请求都会自动包含“Authorization”元数据。你只需要编写一次这样的代码,它就会在所有地方发挥作用。错误映射:
AppException _mapGrpcError(GrpcError error) { switch (error.code) { case StatusCode.unauthenticated: return AppException.unauthorized(...); caseStatusCode_deadlineExceeded: return AppException.timeout(...); ... } }gRPC有一套自己的状态码,这些状态码与HTTP的状态码类似,但并不完全相同。在数据源层将这些状态码映射到应用程序中的异常类型上,就可以确保应用程序的其他部分永远不会直接处理那些特定于gRPC的错误。这样一来,你的代码结构就会更加清晰,且不会受到任何框架的限制。
仓库层
abstract class BankingRepository { Future<Result<UserProfile, AppException>>> getProfile(String userId); Stream<BalanceResponse>> watchBalance(String userId); Stream<Transaction>> streamTransactions(String userId); Future<Result<UploadResponse, AppException>>>> uploadDocument( List<Uint8List>> chunks, String documentType, ); } class BankingRepositoryImpl implements BankingRepository { final BankingRemoteDataSource _dataSource; BankingRepositoryImpl(this._dataSource); @override Future<Result<UserProfile, AppException>>>> getProfile(String userId) async { try { final profile = await _dataSource.getProfile(userId); return Result.success(profile); } on AppException catch (e) { return Result.failure(e); } } @override Stream<BalanceResponse>> watchBalance(String userId) { return _dataSource.watchBalance(userId); } @override Stream<Transaction>> streamTransactions(String userId) { return _dataSource.streamTransactions(userId); } @override Future<Result<UploadResponse, AppException>>>> uploadDocument( List<Uint8List>> chunks, String documentType, ) async { try { final response = await _dataSource.uploadDocument(chunks, documentType); return Result.success(response); } on AppException catch (e) { return Result.failure(e); } } }Riverpod提供者
part 'banking_providers.g.dart'; @riverpod StreambalanceStream(BalanceStreamRef ref, String userId) { final repository = ref.watch(bankingRepositoryProvider); return repository.watchBalance(userId); } @riverpod Stream transactionStream( TransactionStreamRef ref, String userId, ) { final repository = ref.watch(bankingRepositoryProvider); return repository.streamTransactions(userId); } 通过Riverpod,任何返回
Stream的提供者都会自动变成一个可供组件监听的AsyncValue。每当gRPC服务器发送新的数据时,相关的组件就会自动重新更新。用户界面
// lib/features/banking/presentation/pages/dashboard_page.dart class DashboardPage extends ConsumerWidget { final String userId; const DashboardPage({required this.userId, super.key}); @override Widget build(BuildContext context, WidgetRef ref) { final balanceAsync = ref.watch(balanceStreamProvider(userId)); final transactionsAsync = ref.watch(transactionStreamProvider(userId)); return Scaffold( appBar: AppBar(title: const Text('仪表盘')), body: Column( children: [ balanceAsync.when( data: (balance) => BalanceCard( amount: balance.balance, currency: balance(currency, ), loading: () => const BalanceShimmer(), error: (e, _) => ErrorCard(message: e.toString()), ), const SizedBox(height: 24), Expanded( child: transactionsAsync.when( data: (transaction) => TransactionTile( transaction: transaction, ), loading: () => const TransactionShimmer(), error: (e, _) => ErrorCard(message: e.toString()), ), ), ], ), ); } }余额信息和交易记录都是通过gRPC服务器流传输过来的。每当服务器有新数据发送时,这些信息都会实时更新。而Riverpod的
AsyncValue机制使得这两种数据的处理方式完全相同。这个框架确实涵盖了所有必要的功能模式。生产环境相关问题
使用拦截器进行身份验证
如果你需要更精细的身份验证控制机制,比如刷新过期的令牌或自动重试请求,你可以实现一个客户端拦截器:
class AuthInterceptor extends ClientInterceptor { final TokenService _tokenService; AuthInterceptor(this._tokenService); @override ResponseFutureinterceptUnary ( ClientMethodmethod, Q request, CallOptions options, ClientUnaryInvokerinvoker, ) { final token = _tokenService.currentToken; final authenticatedOptions = options.mergedWith( CallOptionsmetadata: {'authorization': 'Bearer $token'}, ); return invoker(method, request, authenticatedOptions); } @override ResponseStreaminterceptServerStreaming ( ClientMethodmethod, Q request, CallOptions options, ClientServerStreamingInvokerinvoker, ) { final token = _tokenService.currentToken; final authenticatedOptions = options.mergedWith( CallOptionsmetadata: {'authorization': 'Bearer $token'}, ); return invoker(method, request, authenticatedOptions); } }在创建客户端时,需要添加拦截器:
_client = BankingServiceClient( _channel, interceptors: [AuthInterceptor(tokenService)], );该拦截器会在每次调用时自动被触发。数据源代码永远不会直接处理认证逻辑。
错误处理:gRPC状态码
gRPC定义了一组标准的状态码,所有实现都必须遵循这些状态码:
| 状态码 | 含义 | >建议采取的措施 |
|---|---|---|
| OK (0) | 操作成功 | 可以直接使用响应结果 |
| CANCELLED (1) | 客户端主动取消了调用 | 忽略该错误或记录日志 |
| UNKNOWN (2) | 服务器出现未知错误 | 显示通用错误提示 |
| INVALID_ARGUMENT (3) | 请求数据无效 | 显示验证错误信息 |
| DEADLINE_EXCEEDED (4) | 调用超时 | 重试或显示超时提示 |
| NOT_FOUND (5) | 资源不存在 | 显示“未找到”界面 |
| ALREADY_EXISTS (6) | 资源重复 | 显示冲突提示信息 |
| PERMISSION_DENIED (7) | 权限不足 | 显示“访问被拒绝”提示 |
| UNAUTHENTICATED (16) | 凭证无效或已过期 | 跳转到登录页面 |
| RESOURCE_EXHAUSTED (8) | 达到请求速率限制 | 稍后重试 |
| UNAVAILABLE (14) | 服务器暂时不可用 | 显示“离线”提示 |
AppException _mapGrpcError(GrpcError error) {
switch (error.code) {
case StatusCode.unauthenticated:
return AppException.unauthorized(message: '会话已过期');
caseStatusCode.permissionDenied:
return AppException.forbidden(message: '访问被拒绝');
case StatusCode.notFound:
return AppException.notFound(message: error.message ?? '未找到');
case StatusCode.deadlineExceeded:
return AppException.timeout(message: '请求超时');
caseStatusCode.unavailable:
return AppException.serverUnavailable(message: '服务当前不可用');
case StatusCode.resourceExhausted:
return AppException.rateLimited(message: '请求次数过多');
case StatusCode.invalidArgument:
return AppException.validation(
message: error.message ?? '输入无效',
);
default:
return AppException.unknown(
message: error.message ?? '发生了错误',
);
}
}
截止时间与超时机制
每个gRPC调用都应该设定截止时间。如果没有截止时间,即使服务器响应缓慢,应用程序也可能会无限期地等待响应。
针对每次调用的截止时间设置:
Future<UserProfile>> getProfile(String userId) async {
return await _client.getProfile(
ProfileRequest(userId: userId),
options: CallOptions(timeout: const Duration(seconds: 10)),
);
}
所有调用的默认截止时间:
_client = BankingServiceClient(
_channel,
options: CallOptions(
timeout: const Duration(seconds: 30),
metadata: {'authorization': 'Bearer $authToken'},
),
);
当超过截止时间时,调用会抛出带有StatusCode.deadlineExceeded的GrpcError异常,你的错误处理机制会对此进行相应的处理。
日志拦截器
class LoggingInterceptor extends ClientInterceptor {
@override
ResponseFuture<R>> interceptUnary<Q, R>>(
ClientMethod<Q, R>> method,
Q request,
CallOptions options,
Client UnaryInvoker<Q, R>> invoker,
) {
final stopwatch = Stopwatch()..start();
debugPrint('[gRPC] --> ${method.path}');
final response = invoker(method, request, options);
response.then((_) {
stopwatch.stop();
debugPrint(
'[gRPC] <-- ${method.path} (${stopwatch.elapsedMilliseconds}ms)',
);
}).catchError((error) {
stopwatch.stop();
debugPrint(
'[gRPC] ERROR ${method.path}: $error (${stopwatch.elapsedMilliseconds}ms)',
);
});
return response;
}
}
gRPC、REST与WebSockets:何时使用哪种技术
何时使用REST
如果API被第三方开发者或外部合作伙伴使用,那么通过HTTP传输JSON数据是一种通用方案——任何语言的开发者都可以直接使用这种技术,而无需学习新的工具或技术。
REST采用简单的请求-响应模型,不需要进行流式数据处理,也只需要一个客户端平台即可。对于那些需要进行基本的CRUD操作的场景来说,REST更易于实现、更便于调试,也更容易进行测试。
公开文档的易读性也非常重要。使用OpenAPI/Swagger技术,REST接口可以生成结构清晰、便于浏览和测试的文档,开发者可以直接在浏览器中查看这些文档。
缓存也是非常重要的一项功能。REST请求的响应结果可以在CDN、反向代理以及浏览器缓存层中进行缓存;而gRPC请求则无法利用标准的HTTP缓存机制。
何时使用WebSockets
当你需要实现真正的双向实时通信,而且你的技术栈中还没有包含gRPC时,WebSockets就是最佳选择。聊天应用、多人游戏以及需要双方能够自由交流的合作工具,都是WebSockets的典型应用场景。
此外,浏览器必须直接支持WebSockets,而不能通过代理层来转发请求。现代所有的浏览器都原生支持WebSockets;而在浏览器中使用gRPC,则需要借助gRPC-Web框架以及代理层才能实现通信功能。
何时使用gRPC
当多个开发团队需要共享相同的服务接口时,gRPC是一个非常实用的选择。例如,如果Flutter、React、Go和Python这些不同语言编写的应用程序都需要调用同一个后端服务,那么使用.proto文件作为接口定义标准,就可以有效避免各个团队在接口实现上出现偏差。
对于那些需要传输大量数据的应用程序来说,gRPC也是更好的选择。因为protobuf支持二进制编码,而且生成的客户端代码也更加简洁高效。特别是当有大量的字段需要被传输,或者有很多应用程序都在使用这些数据时,使用gRPC的优势就更加明显了。
网络环境是多变的,数据包的大小也至关重要。使用2G或3G连接的用户,会直接从protobuf的紧凑二进制格式中受益。用protobuf表示的相同数据,其大小可能会比用JSON表示的数据小3到10倍,这样一来,对于那些数据流量有限的用户来说,加载速度就会更快,数据消耗也会更少。
如果需要进行实时流式传输,那么就需要一个能够适用于所有通信模式的统一框架。gRPC所提供的四种通信模式可以覆盖各种场景,而且无需在API之外再单独部署WebSocket服务器。
当需要高频进行服务之间的交互时,两种内部服务如果通过持久化的、多路复用的HTTP/2连接,并使用二进制protobuf格式进行数据传输,其性能将会显著优于通过HTTP/1.1使用JSON进行的通信。
混合架构
一个成熟的工程决策并不是在gRPC和REST之间做出选择,而是要明确两者各自适用的场景,并有意识地同时使用它们。
大多数规模较大的组织最终都会采用混合架构:公共的REST API用于满足外部用户对简洁性和JSON格式的需求;内部的gRPC网络则负责处理高频、高性能的服务间通信;而移动端的gRPC接口则能够通过高效的二进制连接,为Flutter客户端提供实时数据传输功能。
每一层都使用最适合自身需求的工具。这里不存在对某种协议的盲目推崇,纯粹是出于工程实践的考虑。
结论
远程过程调用的理念起源于这样一个简单的想法:网络通信不应该让开发人员花费额外的精力去处理相关细节。在另一台机器上调用一个函数,应该就像在自己的机器上调用该函数一样自然。
Protocol Buffers进一步解决了数据编码的问题。JSON格式虽然易于阅读,但文本量较大;而通过编译器强制执行的二进制编码方式,能够生成体积更小、解析速度更快的数据包,并且能确保这些数据包符合所有团队事先约定的格式规范。对于那些处于网络环境较差的环境中,或者需要每天处理海量请求的内部系统来说,这种效率无疑具有重大的商业价值。
gRPC将远程过程调用的语义、Protocol Buffers的编码方式以及HTTP/2传输协议结合在一起,形成了一个能够支持四种不同通信模式的框架:单向请求-响应模式、服务器端流式传输模式、客户端流式传输模式以及双向流式传输模式。所有这些功能都是通过同一个生成的客户端接口实现的,使用的是同一条持久化的连接,同时也受到相同的`.proto`协议规范的约束。
如果整个组织都共享一个Protocol Buffers代码库,那么gRPC就不仅仅是一种技术工具,而会真正成为一种工程实践标准。任何对协议规范的修改都会经过审查流程;那些可能破坏系统稳定性的变更也会被编译器及时检测出来。每个团队都会根据同一个源代码生成自己所需的强类型客户端程序,因此无论使用哪种编程语言,各个团队之间的开发进度都能保持同步。
在Flutter框架中,gRPC服务器端的流式传输功能可以与Dart语言中的`Stream`类型以及Riverpod提供的流处理机制完美结合。那些原本需要通过REST进行轮询查询,或者单独部署WebSocket服务器才能实现的实时余额更新和交易数据推送功能,现在都可以通过简单的流式订阅方式轻松实现。服务器负责发送数据,客户端负责接收数据,除此之外再不需要任何额外的配置。在选择使用 gRPC、REST 还是 WebSockets 时,关键不在于个人偏好,而在于让所选工具能够满足具体需求。公共 API 应该使用 REST 来实现;高频的内部服务通信则更适合使用 gRPC;那些被众多内部系统使用的、数据量较大的通信协议,应该选择 protobuf;而需要实时双向通信的场景,则应该采用 gRPC 流式传输技术。对于那些处于不稳定网络环境中的用户来说,提供数据量尽可能小的通信方案才是最合理的选择。
真正理解这些工具的作用原理、明白它们存在的意义,并且知道在什么情况下该使用哪种工具,这才是区分那些只会机械使用工具的工程师与那些能够从系统整体角度进行思考的工程师的关键所在。
祝编程愉快!
相关文章
如何利用Gemini构建人工智能功能:面向开发者的提示工程实用指南
大多数关于提示工程的教学教程都遵循相同的流程:安装SDK,输入API密钥,调用 generateContent 函数,然后打印输出结果。模型会生成一些看似合理的内容,之后教学教程也就结束了。 但当你真正尝试将这个系统投入实际使用时,才会发现其实真正的准备工作根本还没有开始。 “API返回的文本”与“让用户感到可信的实际功能”之间的差距,正是需要耗费大量精力去解决的地方。 这个差距中充满了各种棘手的问题:模型生成的内容听起来和其他聊天机器人没什么两样;它会编造用户从未说过的话;它返回的数据会被用Markdown格式包裹起来;系统会在凌晨2点出现故障;而对于那些只是想得到答案的用户来说,系统展示的
阅读全文
如何使用 shadcn/ui 在 React 中构建一个可重复使用的日期时间选择器
日期和时间选择器这类组件,在设计文件中看起来可能很简洁,但一旦开始实际开发,就会发现它们会消耗大量的资源。你需要一个日历、一个时间选择器,以及一个能够保证这两者同步的状态管理系统,通常还需要范围选择功能以及对应的多语言版本。 本指南将介绍一些现成的选择器组件,你可以直接将这些组件应用到你的React项目中:组合型日期和时间选择器、日期范围选择器以及时间选择器。 所有这些组件都可以作为 Shadcn日期和时间选择器 组件使用,你只需通过一条CLI命令即可安装它们,而无需从头开始开发。 这些组件都是基于Radix和Base UI的基础架构构建的,下面介绍的版本是使用Base UI实现的。此外,这些
阅读全文
如何使用MONAI在超声数据上训练肿瘤分割模型
大多数分割教程都是从选择一个模型开始,将图像输入该模型中,然后调整超参数直到相关指标得到改善。但这种方法忽略了通常最为关键的一步:理解数据本身。 在本教程中,我们首先会对数据集进行详细分析,随后会根据这些分析结果来决定MONAI分割流程中的每一个设计细节。 我们将涵盖以下内容: 本教程适合谁? 关于数据集 什么是MONAI,为什么使用它? 什么是Dice评分? 第1部分——建模前的数据分析 类别平衡对分割结果的影响 患者数量对数据划分的影响 第2部分——构建分割流程 单一配置对象 按患者分组的数据划分方式 由快照自动选择的转换操作 模型、损失函数与评估指标 结果解读 预测结果可视化 失败模式比
阅读全文
如何使用LangSmith来追踪和监控人工智能代理的行为
在本教程中,我将向您展示如何使用LangSmith来追踪和监控本地的AI代理。我们会构建一个简单的本地AI代理,然后为其启用LangSmith追踪功能,这样我们就能通过Web界面查看模型调用情况、工具使用情况以及请求处理延迟等信息。 我们将使用LangChain v1、Ollama、Qwen以及Python这些工具。除了用于实现观测功能的组件外,所有操作都在您的本地机器上完成,因此代理本身不会产生任何与模型API相关的费用。 目录 背景知识 什么是可观测性与监控? 什么是LangSmith? 开发动机与架构设计 步骤1:安装Ollama并下载模型 步骤2:安装Python相关依赖库 步骤3:启
阅读全文