【框架工具#4】Brpc 安装和使用

📃个人主页:island1314
⛺️ 欢迎关注:👍点赞 👂🏽留言 😍收藏 💞 💞 💞
- 生活总是不会一帆风顺,前进的道路也不会永远一马平川,如何面对挫折影响人生走向 – 《人民日报》
🔥 目录
brpc 安装和使用
1. 基本概述
概念:百度开源的高性能、跨平台的 RPC框架,主要目的就是简化分布式系统中服务间通信。brpc提供了一个简单高效的方式来构建和管理微服务架构中的服务调用。支持同步与异步两种RPC调用
使用流程
服务端
- 定义服务接口(如使用 Protobuf)
- 实现服务端的具体逻辑(通过继承生成的服务接口类并实现方法)
- 启动 brpc 服务并监听请求
客户端
- 创建
brpc::Channel对象并初始化连接 - 使用
example::EchoService_Stub生成的客户端接口进行 RPC 调用 - 处理同步或异步的请求与响应
特点
- 高性能:brpc 使用了多种优化技术,如零拷贝(zero-copy)和高效的序列化方式,能够处理高并发的请求;支持多线程和线程池机制,能够在多核机器上有效分配计算资源
- 异步支持:支持异步 RPC 调用,允许客户端在不阻塞的情况下发起请求,并在后续处理完成时得到回调,极大提升了并发性能
- 跨平台:brpc 是一个跨平台的框架,支持 Linux、macOS、Windows 等多个操作系统,且支持多种编程语言(C++ 和 Python)
- 高可用和负载均衡:brpc 提供了连接池、自动重试、超时控制等特性来保证高可用性
- 服务发现与管理:brpc 与服务发现系统(如 ZooKeeper、Consul)集成,支持自动注册与发现服务。服务管理通过 brpc::Server 类来实现
- 接口定义与协议:brpc 支持通过 Protobuf(Protocol Buffers)来定义服务接口和数据模型,使得服务接口的定义清晰且与编程语言无关
类比理解
如果将RPC服务想象成从一个城市图书馆获取服务的逻辑,那么梳理各个部分的职责
图书馆服务(EchoService_stud):这个类通过继承重写出来的方法,就是图书馆可以提供的服务,例如图书馆可以提供上传服务(时查询是否有这本书,然后提供给用户文件内容),捐赠书籍(接收用户提供的书籍,然后将着呢书放到书架)
- 所以这个服务就可以理解成在图书馆设定的各个岗位,这些岗位的工作人员每一个提供不同的服务
- 如
EchoService_stud就是这个“岗位说明书”——它声明了能提供哪些服务(如echo()),然后去实现它
图书馆(服务器):一般项目中为了代码繁杂,一般对该处代码的构建使用建造者模式。
- 那么就可以理解成通过建造一个图书馆,也就是有这个提供服务服务器,然后再有后面发生的事情
- 整个图书馆是一个物理空间,有多个岗位(服务实现),由管理员(RPC Server)统一管理、调度。
读者(客户端):读者不知道图书馆内部结构,只关心“我要借一本书”(创建请求) → 找对应的工作人员(发起请求) → 等待工作人员回应(等待回应) → 得到结果(处理响应)。这完全就是 Client 调用service.echo("hello")的过程
服务注册(地图标注/做宣传):新开一家图书馆,但没人知道在哪 → 服务没注册,客户端就找不到。而注册中心就是 “城市电子地图”,让所有客户端能动态发现“哪里有 EchoService”
深思
① “服务注册”不只是“贴个标签”
- 地图上标注了图书馆,但读者怎么知道哪个管理员能处理我的请求?
- → 对应:服务注册时,不仅要注册地址,还要注册“服务名”和“版本”(如
EchoService:v1)
- → 对应:服务注册时,不仅要注册地址,还要注册“服务名”和“版本”(如
- 如果图书馆太忙,来了10个读者,只有一个管理员?
- → 对应:负载均衡 —— 注册中心可能注册了多个“分馆”(多个服务实例),客户端随机选一个。
② “读者”不是每次都亲自跑过去
- 如果每次借书都要跑一趟图书馆,太慢了!
- → 对应:客户端缓存服务地址、使用连接池,就像你办了借书卡,不用每次都登记。
- 如果图书馆关门了(服务宕机)?
- → 对应:心跳检测 + 服务剔除,注册中心会定期检查“图书馆是否还在营业”。
③ “建造者模式”不只是为了“代码整洁”
实际上,RPC Server 的构建通常是这样的:
server := grpc.NewServer()
pb.RegisterEchoServiceServer(server, &EchoServiceImpl{})
lis, _ := net.Listen("tcp", ":50051")
server.Serve(lis)
- 这就是一个“建造者”:先创建 Server,再注册服务,再监听端口 —— 步骤清晰,职责分离。
④ “服务实现类”其实是“契约执行者”
EchoService_stud是接口(Contract),它规定:“只要你实现了这个方法,我就认你是合格的借书员”。- 客户端不关心你内部怎么查书、怎么扫描、怎么归位 —— 只要你按约定返回结果就行。
- → 这就是“接口隔离”和“面向契约编程”,是解耦的核心!
对比其他 rpc 框架
| 对比项 | BRPC | GRPC |
|---|---|---|
| 性能 | 更高(C++ 优化极致) | 高,但略低 |
| 协议支持 | 多协议(HTTP、Redis、bRPC 等) | 主要 gRPC |
| 服务发现 | 内置 etcd/ZK/DNS | 需外部集成 |
| 负载均衡 | 多种策略内置 | 较少 |
| 监控 | bvar + Prometheus | 需 Prometheus |
| 社区 | 国内活跃 | 国际主流,且下载 submodule 需 vpn |
| 易用性 | 中等(C++ 接口) | 高(多语言支持) |
应用场景
- 微服务通信:服务间高性能调用
- 分布式存储:如分布式 KV、文件系统
- 推荐系统:高并发特征查询
- 广告系统:实时竞价(RTB)
- 内部中间件:自研消息队列、配置中心
2. 安装配置
安装依赖
sudo apt-get install -y git g++ make libssl-dev libprotobuf-dev libprotoc-dev protobuf-compiler libleveldb-dev
安装 brpc
git clone https://github.com/apache/brpc.git
cd brpc/
mkdir build && cd build
cmake -DCMAKE_INSTALL_PREFIX=/usr .. && cmake --build . -j6
make && sudo make install
PB 多版本冲突问题
链接Protobuf库时引起冲突

问题解决:直接卸载冲突版本解决问题
sudo apt remove protobuf-compiler libprotobuf-dev
// 更新动态库
sudo ldconfig
3. 类和接口介绍
3.1 接口学习
该处接口基于客户端和服务端代码学习
brpc::Channel管理着与服务器的通信,负责连接和数据交换brpc::Controller管理RPC的调用状态,记录错误信息和控制调用流程example::EchoService_Stub是客户端与服务端通信的接口,包含RPC请求方法example::EchoRequest和 example::EchoResponse 分别定义了请求和响应消息的格式google::protobuf::Closure和 NewCallback 用于在异步调用中处理回调brpc::Server是 brpc 框架中最重要的类之一,负责管理服务的注册、启动和监听请求。通过 AddService 和 Start 等方法来注册服务和启动服务器brpc::ServiceOwnership控制服务对象的生命周期管理,选择由 brpc 服务器或开发者管理brpc::ClosureGuard是一个辅助工具,确保在异步 RPC 调用结束时自动调用回调函数brpc::ServerOptions提供了配置 brpc 服务器的一些参数,如空闲超时、线程数等
3.2 类说明
① brpc :: Channel类:主要负责客户端和服务端的通信,管理底层的TCP连接、RPC请求的发送和响应的接收;Channel 本身就是“通道”,此处表示的是客户端和服务端之间的连接通道
Init(const std::string& server, const ChannelOptions* options):初始化通道并连接到指定的服务器地址CallMethod:发起RPC请求,通常与 Controller 配合使用
brpc::Channel channel;
brpc::ChannelOptions options;
options.connect_timeout_ms = 1000; // 设置连接超时1秒
options.timeout_ms = 5000; // 设置请求超时5秒
options.max_retry = 3; // 设置最大重试次数为3
options.protocol = "baidu_std"; // 设置通信协议为baidu_std
int ret = channel.Init("127.0.0.1:8080", &options); // 初始化信道并连接到服务器
if (ret != 0) {
std::cout << "初始化信道失败!" << std::endl;
}
② brpc :: Controller类:Controller本身是“控制器”的含义,该类就是控制和管理RPC的调用状态。用于管理RPC调用的状态,包括请求的发送、响应的接收以及错误信息的处理。每个RPC请求都需要一个 Controller 对象,通常在每次请求时创建并传递给服务存根(Stub)进行调用
- Failed():返回是否发生了错误。如果RPC调用失败,返回 true
- ErrorText():获取错误信息(如果 Failed() 返回 true)
- Reset():重置 Controller 的状态,适用于同一个 Controller 重复使用的情况
brpc::Controller cntl;
example::EchoResponse response;
cntl.Reset(); // 如果需要重用,重置状态
stub.Echo(&cntl, &req, &response, nullptr);
if (cntl.Failed()) {
std::cout << "Rpc调用失败:" << cntl.ErrorText() << std::endl;
} else {
std::cout << "收到响应: " << response.message() << std::endl;
}
③ example :: EchoService_stub 类:stub本身是代理的意思,在RPC框架中表示客户端的代码,其提供了与远程服务进行交互的接口
由 Protocol Buffers 自动生成的客户端存根(Stub)类。它提供了与远程 EchoService 服务进行通信的接口。在该类中,定义了 Echo 方法来向服务端发送请求,并接收响应
Echo(Controller* cntl, const EchoRequest* request, EchoResponse* response, google::protobuf::Closure* done):执行RPC请求并接收响应,支持同步和异步调用。done 参数用于传入回调函数
// 同步调用
example::EchoService_Stub stub(&channel);
example::EchoRequest req;
req.set_message("Hello, Echo Server");
example::EchoResponse rsp;
brpc::Controller cntl;
stub.Echo(&cntl, &req, &rsp, nullptr);
if (cntl.Failed()) {
std::cerr << "RPC调用失败:" << cntl.ErrorText() << std::endl;
} else {
std::cout << "收到响应: " << rsp.message() << std::endl;
}
④ example :: EchoRequest 类:Protocol Buffers 自动生成的消息类,用于定义RPC请求的结构。在本例中,它包含一个 message 字段,用于传递请求消息
set_message(const std::string& msg):设置请求消息message():获取请求消息
example::EchoRequest req;
req.set_message("你好,RPC服务!"); // 设置请求消息内容
⑤ example :: EchoResponse 类:Protocol Buffers 自动生成的消息类,用于定义RPC响应的结构
et_message(const std::string& msg):设置响应消息message():获取响应消息
example::EchoResponse rsp;
rsp.set_message("Hello, client!"); // 设置响应消息内容
⑥ google :: protobuf :: Closure 和 NewCallback:Closure在编程中是“闭包”,对回调函数的一种封装,异步调用完成后调用其来处理结果
NewCallback表示创建回调,也就是用于创建Closure对象,将回调函数和其参数进行绑定
google::protobuf::Closure是 Protocol Buffers 用于封装回调函数的类。NewCallback 是一个静态工厂函数,用于创建回调函数对象,并将回调与传入的参数绑定google::protobuf::NewCallback(callback, args...):创建一个回调对象,callback 是回调函数,args… 是传递给回调的参数
auto closure = google::protobuf::NewCallback(callback, cntl, rsp);
stub.Echo(&cntl, &req, &rsp, closure); // 传递回调函数closure,进行异步调用
⑦ brpc :: Server:brpc 框架中的核心类,用于管理和启动服务端。它负责注册服务、启动服务器并监听客户端请求,以及处理 RPC 请求
参数说明
-
AddService(Service* service, ServiceOwnership ownership)-
service:添加到服务器中的服务对象,通常是继承自 brpc::Service 的自定义服务类
-
ownership:指定 brpc 是否管理服务对象的生命周期(SERVER_OWNS_SERVICE 或 SERVER_DOESNT_OWN_SERVICE)
-
-
Start(int port, const ServerOptions* options = nullptr)
- port:服务器监听的端口
- options:配置选项,包含如线程数、超时时间等
-
RunUntilAskedToQuit():无参数。进入服务器的主事件循环,直到接收到退出请求
class Server {
public:
// 向服务器添加服务
int AddService(Service* service, ServiceOwnership ownership);
// 启动服务器,开始监听指定端口
int Start(int port, const ServerOptions* options = nullptr);
// 进入主事件循环,等待请求并处理
void RunUntilAskedToQuit();
};
// 使用
brpc::Server server;
// 向服务器添加服务
server.AddService(&echo_service, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
// 启动服务器监听端口8080
server.Start(8080, nullptr);
// 进入事件循环,直到收到退出信号
server.RunUntilAskedToQuit();
⑧ brpc :: ServiceOwnership:一个枚举类型,指示服务器是否管理服务对象的生命周期
- SERVER_OWNS_SERVICE:brpc 服务器负责服务对象的生命周期,服务器关闭时自动销毁服务对象
- SERVER_DOESNT_OWN_SERVICE:服务对象的生命周期由用户控制,服务器不会管理服务对象
enum class ServiceOwnership {
SERVER_OWNS_SERVICE, // 服务器负责管理服务对象生命周期
SERVER_DOESNT_OWN_SERVICE // 服务器不负责管理服务对象生命周期
};
⑨ brpc :: ClosureGuard: 一个 RAII 类,确保回调函数(通过 Closure 对象传递的)会在作用域结束时自动执行;确保 done->Run() 会在方法结束时被调用,通常用于处理异步 RPC 回调
class ClosureGuard {
public:
explicit ClosureGuard(google::protobuf::Closure* closure);
~ClosureGuard();
};
⑩ brpc :: ServerOptions:服务器选项 配置 brpc 服务器的一些参数,例如线程数、超时等
class ServerOptions {
public:
int idle_timeout_sec; // 连接空闲超时时间,单位为秒
int num_threads; // 处理请求的线程数
};
4. Rpc 调用实现样例
服务端:
- 创建rpc服务子类继承pb中的EchoService服务类,并实现内部的业务接口逻辑创建rpc服务器类,
- 搭建服务器
- 向服务器类中添加 rpc子服务对象 – 告诉服务器收到什么请求用哪个接口处理
- 启动服务器
客户端:
- 创建网络通信信道
- 实例化pb中的
EchoService_Stub类对象 - 发起rpc请求,获取响应进行处理
创建 proto 文件 —— main.proto
syntax="proto3";
package example;
option cc_generic_services = true;
// 定义 Echo 方法请求参数结构
message EchoRequest{
string message = 1;
}
// 定义 Echo 方法响应参数结构
message EchoResponse{
string message = 1;
}
// 定义 RPC 远端方法
service EchoService{
rpc Echo(EchoRequest) returns (EchoResponse);
}
- 然后编译:
protoc --cpp_out=./ main.proto
Makefile
all : server client
server : server.cc main.pb.cc
g++ -std=c++17 $^ -o $@ -lbrpc -lgflags -lssl -lcrypto -lprotobuf -lleveldb
client : client.cc main.pb.cc
g++ -std=c++17 $^ -o $@ -lbrpc -lgflags -lssl -lcrypto -lprotobuf -lleveldb
.PHONY: clean
clean:
rm -rf server client
server.cc
-
服务定义和实现:通过继承 Protobuf 生成的 EchoService 接口并实现其中的 Echo 方法,提供具体的业务逻辑。Echo 方法接收到客户端的消息后,简单地在消息后面加上“–这是响应!!”并将其返回
-
brpc 服务器启动
- 创建并配置 brpc::Server
- 将 EchoServiceImpl 服务添加到服务器中
- 启动服务器并监听指定端口(8080)
#include <brpc/server.h>
#include <butil/logging.h>
#include "main.pb.h"
// 1. 继承于EchoService创建一个子类,并实现rpc调用的业务功能
class EchoServiceImpl : public example::EchoService {
public:
EchoServiceImpl(){}
~EchoServiceImpl(){}
void Echo(google::protobuf::RpcController* controller,
const ::example::EchoRequest* request,
::example::EchoResponse* response,
::google::protobuf::Closure* done)
{
brpc::ClosureGuard rpc_guard(done);
std::cout << "收到消息:" << request->message() << "\n";
std::string str = request->message() + "--这是响应!!";
response->set_message(str);
// done->Run();
}
};
int main(int argc, char *argv[]){
// 关闭 brpc 的默认日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
// 2. 构造服务器对象
brpc::Server server;
// 3. 向服务器对象中,新增EchoService服务
EchoServiceImpl echo_service;
int ret = server.AddService(&echo_service, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if (ret == -1) {
std::cout << "添加Rpc服务失败!\n";
return -1;
}
// 4. 启动服务器
brpc::ServerOptions options;
options.idle_timeout_sec = -1; //连接空闲超时时间-超时后连接被关闭
options.num_threads = 1; // io线程数量
ret = server.Start(8080, &options);
if (ret == -1) {
std::cout << "启动服务器失败!\n";
return -1;
}
server.RunUntilAskedToQuit();//修改等待运行结束
return 0;
}
client.cc
- 通过
EchoService_Stub向服务器发送一个请求,并在请求完成后通过回调函数处理响应 - 实现同步调用和异步调用两种方式
#include <brpc/channel.h>
#include <thread>
#include "main.pb.h"
void callback(brpc::Controller* cntl, ::example::EchoResponse* response) {
std::unique_ptr<brpc::Controller> cntl_guard(cntl);
std::unique_ptr<example::EchoResponse> resp_guard(response);
if (cntl->Failed() == true) {
std::cout << "Rpc调用失败: " << cntl->ErrorText() << "\n";
return;
}
std::cout << "收到响应: " << response->message() << "\n";
}
int main(int argc, char *argv[]){
// 1. 构造 Channel 信道 连接服务器
brpc::ChannelOptions options;
options.connect_timeout_ms = -1; // 连接等待超时时间,-1表示一直等待
options.timeout_ms = -1; // rpc请求等待超时时间,-1表示一直等待
options.max_retry = 3; // 请求重试次数
options.protocol = "baidu_std"; // 序列化协议,默认使用baidu_std
brpc::Channel channel;
int ret = channel.Init("127.0.0.1:8080", &options);
if (ret == -1) {
std::cout << "初始化信道失败!\n";
return -1;
}
// 2. 构造 EchoService_stub 对象 进行 rpc 调用
example::EchoService_Stub stub(&channel);
// 3. rpc 调用
example::EchoRequest req;
req.set_message("你好 Island");
brpc::Controller *cntl = new brpc::Controller();
example::EchoResponse *rsp = new example::EchoResponse();
// 同步处理流程
stub.Echo(cntl, &req, rsp, nullptr);
if (cntl->Failed() == true) {
std::cout << "Rpc调用失败:" << cntl->ErrorText() << std::endl;
return -1;
}
std::cout << "收到响应: " << rsp->message() << std::endl;
delete cntl;
delete rsp;
// 异步调用处理流程
// auto clusure = google::protobuf::NewCallback(callback, cntl, rsp);
// stub.Echo(cntl, &req, rsp, clusure);
// std::cout << "异步调用结束!\n";
// std::this_thread::sleep_for(std::chrono::seconds(3));
return 0;
}
运行结果 如下:
lighthouse@VM-8-10-ubuntu:code$ ./server
收到消息:你好 Island
lighthouse@VM-8-10-ubuntu:code$ ./client
收到响应: 你好 Island--这是响应!!
如果把 client.cc 同步调用那代码改成异步调用,此时结果如下:
lighthouse@VM-8-10-ubuntu:code$ ./server
收到消息:你好 Island
lighthouse@VM-8-10-ubuntu:code$ ./client
异步调用结束!
收到响应: 你好 Island--这是响应!!
5. 二次封装
5.1 基本思想
原因分析:brpc 的核心功能是RPC调用,但是在分布式系统中,单纯的RPC调用无法满足需求,需要考虑清楚三个问题
-
提供服务的主体:服务发现,例如通过注册中心(Etcd)动态获取服务列表(后续会将两者借结合)
-
高效管理服务通信的方式:负载均衡、连接复用等
-
扩展和维护服务端实现:动态增减服务节点,避免大量代码被修改
综上所述,brpc的二次封装主要目的是信道管理,不是对RPC方法的封装
核心思想
① 指定服务的信道管理类
- 一个服务可能会有多个节点提供服务,每个节点都对应一个独立的
Channel - 封装的目标就是将这些节点的信道管理起来,建立服务于信道之间的映射关系(一般情况下应该是多对一的映射关系,也就是一个服务名对应多个节点的Channel)
- 负载均衡实现:通过轮询策略,在多个节点中选择一个Channel发起请求
② 全局信道管管理类
- 对多个服务的信道进行 统一管理
- 提供一个全局的入口,便于通过服务名称动态获取对应的信道
- 动态加载或者移除服务信道,支持线上服务的动态拓展和下线
③ 封装实现逻辑分析
-
服务与信道的管理
- 将每个服务节点的信道管理起来,每个节点有独立的Channel
- 通过服务ID找到对应的信道组
-
发起RPC调用时的信道选择
- 根据服务名称获取信道组;在信道组中根据负载均衡策略选择一个信道
- 使用选择的信道发起RPC调用
-
动态拓展
- 当新的服务节点上线时,动态创建并添加信道
- 当服务节点下线时,动态移除信道,保证资源回收
5.2 具体实现
下面都是 Channel.hpp 里面的类, 只是分成了两份代码
① ServiceChannel 类是一个 基于 brpc 的服务通信信道管理器,用于管理 某个服务 的多个节点(host)的 brpc::Channel,并提供 轮询(RR)负载均衡。代码如下:
// 1. 封装单个服务的信道管理类
class ServiceChannel{
public:
using ptr = std::shared_ptr<ServiceChannel>; // 类目别名 简化使用
using ChannelPtr = std::shared_ptr<brpc::Channel>;
ServiceChannel(const std::string& name)
:_service_name(name), _index(0)
{}
// 1. 服务上线新节点, 调用 append 新增信道
void append(const std::string &host){
// 1. 创建信道
auto channel = std::make_shared<brpc::Channel>();
// 2. 配置 ChannelOptions
brpc::ChannelOptions options;
options.connect_timeout_ms = -1; // 连接等待超时时间
options.timeout_ms = -1; // rpc请求等待超时时间
options.max_retry = 3; // 请求重试次数
options.protocol = "baidu_std"; // 序列化协议,默认使用baidu_std
// 3. 信道初始化
if(channel->Init(host.c_str(), &options) == -1){
LOG_ERROR("初始化{}-{}信道失败", _service_name, host);
return;
}
// 4. 加锁 & 添加数据
std::unique_lock<std::mutex> lock(_mutex);
_hosts.insert(std::make_pair(host, channel));
_channels.push_back(channel);
}
// 2. 服务下线新节点, 调用 remove 释放信道
void remove(const std::string &host){
// 1. 查找释放存在该信道
std::unique_lock<std::mutex> lock(_mutex);
auto it = _hosts.find(host);
if(it == _hosts.end()){
LOG_WARN("{}-{}节点删除信道时,没有找到信道信息!", _service_name, host);
return;
}
// 2. 遍历信道列表删除信道
for(auto i = _channels.begin(); i != _channels.end(); ++i){
if(*i == it->second){
_channels.erase(i);
break;
}
}
_hosts.erase(it); // 记得清除对应映射
}
// 3. 获取 Channel 用于发起对应服务的 Rpc 调用 -- RR 轮转策略
ChannelPtr choose(){
std::unique_lock<std::mutex> lock(_mutex);
if(_channels.size() == 0){
LOG_ERROR("当前没有能够提供 {} 服务的节点!", _service_name);
return ChannelPtr();
}
int idx = _index++ % _channels.size();
return _channels[idx];
}
private:
std::mutex _mutex;
int32_t _index; // 轮询计数器
std::string _service_name; // 服务名称
std::vector<ChannelPtr> _channels; // 信道列表(用于轮询)
std::unordered_map<std::string, ChannelPtr> _hosts; // 主机地址与信道映射关系
};
功能
- 管理一个服务(如
user_service)的所有可用节点(host:port) - 每个节点对应一个
brpc::Channel(brpc 的客户端通信通道) - 支持动态添加(
append)和删除(remove)节点 - 提供 轮询(Round-Robin) 策略选择
Channel发起 RPC 调用
② ServiceManager 类是一个 服务管理中枢,它与前面的 ServiceChannel 配合,实现了 服务发现 + 动态信道管理 + 负载均衡选择 的完整客户端逻辑
// 2. 总体的服务信道管理类
class ServiceManager {
public:
using ptr = std::shared_ptr<ServiceManager>;
ServiceManager() {}
// 1. 为指定服务选择一个可用的 Channel(用于发起 RPC)
ServiceChannel::ChannelPtr choose(const std::string &service_name) {
std::unique_lock<std::mutex> lock(_mutex);
// 1. 查找是否存在
auto sit = _services.find(service_name);
if(sit == _services.end()){
LOG_ERROR("当前没有能够提供 {} 服务的节点!", service_name);
return ServiceChannel::ChannelPtr();
}
// 2. RR 策略选择
return sit->second->choose();
}
// 2. 声明我关注某个服务,注意 不关心的就不用管理
void declared(const std::string &service_name) {
std::unique_lock<std::mutex> lock(_mutex);
_follow_services.insert(service_name);
}
// 3. 处理“服务上线”事件,添加新节点
void onServiceOnline(const std::string &service_instance, const std::string &host) {
// 1. 解析服务名
std::string service_name = getServiceName(service_instance);
ServiceChannel::ptr service;
{
std::unique_lock<std::mutex> lock(_mutex);
// 2. 检查是否关注该服务
auto fit = _follow_services.find(service_name);
if (fit == _follow_services.end()) {
LOG_DEBUG("{}-{} 服务上线了,但是当前并不关心!", service_name, host);
return;
}
// 3. 获取/创建 对应的 ServiceChannel
auto sit = _services.find(service_name);
if (sit == _services.end()) {
// 不存在则创建
service = std::make_shared<ServiceChannel>(service_name);
_services.insert(std::make_pair(service_name, service));
}else {
service = sit->second;
}
}
if (!service) {
LOG_ERROR("新增 {} 服务管理节点失败!", service_name);
return ;
}
// 4. 添加信道
service->append(host);
LOG_DEBUG("{}-{} 服务上线新节点,进行添加管理!", service_name, host);
}
// 4. 处理“服务下线”事件,删除节点
void onServiceOffline(const std::string &service_instance, const std::string &host) {
// 1. 解析服务名
std::string service_name = getServiceName(service_instance);
ServiceChannel::ptr service;
{
std::unique_lock<std::mutex> lock(_mutex);
// 2. 检查是否关注
auto fit = _follow_services.find(service_name);
if (fit == _follow_services.end()) {
LOG_DEBUG("{}-{} 服务下线了,但是当前并不关心!", service_name, host);
return;
}
// 3. 查找对应 服务信道
auto sit = _services.find(service_name);
if(sit == _services.end()){
LOG_WARN("删除{}服务节点时,没有找到管理对象", service_name);
return;
}
service = sit->second;
}
// 4. 调用 remove 删除信道
service->remove(host);
LOG_DEBUG("{}-{} 服务下线节点,进行删除管理!", service_name, host);
}
private:
// 从服务实例路径提取服务名
// 如: /service/user/127.0.0.1:8080 ---> /service/user
std::string getServiceName(const std::string &service_instance) {
auto pos = service_instance.find_last_of('/');
if (pos == std::string::npos) return service_instance;
// 提取字符串
return service_instance.substr(0, pos);
}
private:
std::mutex _mutex;
std::unordered_set<std::string> _follow_services; // 关心哪些服务
std::unordered_map<std::string, ServiceChannel::ptr> _services; // 服务名 -> ServiceChannel
};
职责
- 声明关注的服务(
declared) - 管理所有服务的
ServiceChannel(_services) - 接收服务上下线通知(
onServiceOnline/onServiceOffline) - 为调用方提供“按服务名选择信道”的接口(
choose)
注意:在上面服务下线的 代码中,是调用 service->remove(host),而不是直接从 _services 中删除整个服务?
为“服务下线” ≠ “服务不存在”
它只是 某个节点(host)下线了,而不是整个服务消失。所以: 要从
ServiceChannel中 移除这个 host 的信道,而不是从_services中删除整个ServiceChannel
举个生活中的例子(想象一个外卖平台)
- 服务名:
骑手调度服务 - 节点:
骑手A、骑手B、骑手C - 当
骑手A下线(休息):不代表“骑手调度服务”没了,只是少了一个可用节点,但是其他骑手(B、C)还能接单 - 因此应该:从 “可用骑手列表” 中移除
骑手A,但保留 “骑手调度服务” 的管理对象
5.3 与 ectd 的联调测试
Makefile
all: discovery registry
discovery: discovery.cc main.pb.cc
g++ -o $@ $^ -std=c++17 -letcd-cpp-api -lcpprest -lgflags -lfmt -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
registry: registry.cc main.pb.cc
g++ -o $@ $^ -std=c++17 -letcd-cpp-api -lcpprest -lgflags -lfmt -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
.PHONY: clean
clean:
rm -f discovery registry
registry.cc
DEFINE_bool(run_mode, false, "程序运行模式: false-调试; true 发布");
DEFINE_string(log_file, "", "发布模式下指定日志的输出文件");
DEFINE_int32(log_level, 0, "发布模式下指定日志的输出等级");
DEFINE_string(etcd_host, "http://127.0.0.1:2379", "服务注册中心地址");
DEFINE_string(base_service, "/service", "服务监控根目录");
DEFINE_string(instance_name, "/echo/instance", "当前实例名称");
DEFINE_string(access_host, "127.0.0.1:7070", "当前实例的外部访问地址");
DEFINE_int32(listen_port, 7070, "Rpc服务器监听端口");
class EchoServiceImpl : public example::EchoService {
public:
EchoServiceImpl(){}
~EchoServiceImpl(){}
void Echo(google::protobuf::RpcController* controller,
const ::example::EchoRequest* request,
::example::EchoResponse* response,
::google::protobuf::Closure* done)
{
brpc::ClosureGuard rpc_guard(done);
std::cout << "收到消息:" << request->message() << "\n";
std::string str = request->message() + "--这是响应!!";
response->set_message(str);
}
};
int main(int argc, char* argv[]){
google::ParseCommandLineFlags(&argc, &argv, true);
init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 服务器改造
// 关闭 brpc 的默认日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
// 2. 构造服务器对象
brpc::Server server;
// 3. 向服务器对象中,新增EchoService服务
EchoServiceImpl echo_service;
int ret = server.AddService(&echo_service, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if (ret == -1) {
std::cout << "添加Rpc服务失败!\n";
return -1;
}
// 4. 启动服务器
brpc::ServerOptions options;
options.idle_timeout_sec = -1; // 连接空闲超时时间-超时后连接被关闭
options.num_threads = 1; // io 线程数量
ret = server.Start(FLAGS_listen_port, &options);
if (ret == -1) {
std::cout << "启动服务器失败!\n";
return -1;
}
// 5. 注册服务
Registry::ptr rclient = std::make_shared<Registry>(FLAGS_etcd_host);
// 注册: /service/echo/instance 发现: /service/echo
rclient->registry(FLAGS_base_service + FLAGS_instance_name, FLAGS_access_host);
server.RunUntilAskedToQuit(); // 休眠修改等待运行结束
return 0;
}
discovery.cc
DEFINE_bool(run_mode, false, "程序运行模式: false-调试; true 发布");
DEFINE_string(log_file, "", "发布模式下指定日志的输出文件");
DEFINE_int32(log_level, 0, "发布模式下指定日志的输出等级");
DEFINE_string(etcd_host, "http://127.0.0.1:2379", "服务注册中心地址");
DEFINE_string(base_service, "/service", "服务监控根目录");
DEFINE_string(call_service, "/service/echo", "服务监控根目录");
int main(int argc, char* argv[]){
google::ParseCommandLineFlags(&argc, &argv, true);
init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 1. 构造 rpc 信道管理对象
auto sm = std::make_shared<ServiceManager>();
sm->declared(FLAGS_call_service);
auto put_cb = std::bind(&ServiceManager::onServiceOnline, sm.get(), std::placeholders::_1, std::placeholders::_2);
auto del_cb = std::bind(&ServiceManager::onServiceOffline, sm.get(), std::placeholders::_1, std::placeholders::_2);
// 2. 构造服务发现对象
Discovery::ptr rclient = std::make_shared<Discovery>(FLAGS_etcd_host, FLAGS_base_service, put_cb, del_cb);
while(1){
// 3. 提供 rpc 信道管理对象 获取提供 Echo 服务的信道
auto channel = sm->choose(FLAGS_call_service);
if(!channel){
std::this_thread::sleep_for(std::chrono::seconds(1));
// continue;
return -1;
}
// 4. 发起 EchoRpc 调用
example::EchoService_Stub stub(channel.get());
example::EchoRequest req;
req.set_message("你好 Island");
brpc::Controller *cntl = new brpc::Controller();
example::EchoResponse *rsp = new example::EchoResponse();
// 同步处理流程
stub.Echo(cntl, &req, rsp, nullptr);
if (cntl->Failed() == true) {
std::cout << "Rpc调用失败: " << cntl->ErrorText() << std::endl;
delete cntl;
delete rsp;
std::this_thread::sleep_for(std::chrono::seconds(1));
continue;
}
std::cout << "收到响应: " << rsp->message() << std::endl;
std::this_thread::sleep_for(std::chrono::seconds(1));
}
return 0;
}
结果输出如下:
# 先服务注册
lighthouse@VM-8-10-ubuntu:test$ ./registry
收到消息:你好 Island
^C
# 再服务发现
lighthouse@VM-8-10-ubuntu:test$ ./discovery
[default-logger][17:58:04][3401509][debug ][../../../common/channel.hpp:125] /service/echo-127.0.0.1:7070 服务上线新节点,进行添加管理!
收到响应: 你好 Island--这是响应!!
收到响应: 你好 Island--这是响应!!
收到响应: 你好 Island--这是响应!!
I0907 17:58:07.802816 3401562 4294969856 /home/lighthouse/code/project/chat-client/code/server/example/brpc/brpc/src/brpc/socket.cpp:2589 CheckHealth] Checking Socket{id=0 addr=127.0.0.1:7070} (0x55fc163c6170)
Rpc调用失败: [E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R1][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R2][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R3][E112]Not connected to 127.0.0.1:7070 yet, server_id=0
Rpc调用失败: [E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R1][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R2][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R3][E112]Not connected to 127.0.0.1:7070 yet, server_id=0
Rpc调用失败: [E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R1][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R2][E112]Not connected to 127.0.0.1:7070 yet, server_id=0 [R3][E112]Not connected to 127.0.0.1:7070 yet, server_id=0
[default-logger][17:58:10][3401564][debug ][../../../common/channel.hpp:150] /service/echo-127.0.0.1:7070 服务下线节点,进行删除管理!
[default-logger][17:58:10][3401564][debug ][../../../common/etcd.hpp:72] 删除服务: /service/echo/instance-127.0.0.1:7070
[default-logger][17:58:10][3401509][error ][../../../common/channel.hpp:59] 当前没有能够提供 /service/echo 服务的节点!
[warn] watcher does't exit normally
lighthouse@VM-8-10-ubuntu:test$
对 discovery.cc 的部分代码作修改,如下:
if(!channel){
std::this_thread::sleep_for(std::chrono::seconds(1));
continue;
// return -1;
}
// ....
std::cout << "收到响应: " << rsp->message() << std::endl;
// std::this_thread::sleep_for(std::chrono::seconds(1));
break;
此时结果输出
# 先服务发现
lighthouse@VM-8-10-ubuntu:test$ ./discovery
[default-logger][17:53:05][3399406][error ][../../../common/channel.hpp:85] 当前没有能够提供 /service/echo 服务的节点!
[default-logger][17:53:06][3399406][error ][../../../common/channel.hpp:85] 当前没有能够提供 /service/echo 服务的节点!
[default-logger][17:53:07][3399406][error ][../../../common/channel.hpp:85] 当前没有能够提供 /service/echo 服务的节点!
[default-logger][17:53:08][3399451][debug ][../../../common/channel.hpp:125] /service/echo-127.0.0.1:7070 服务上线新节点,进行添加管理!
[default-logger][17:53:08][3399451][debug ][../../../common/etcd.hpp:69] 新增服务:/service/echo/instance-127.0.0.1:7070
收到响应: 你好 Island--这是响应!!
[warn] watcher does't exit normally
# 再服务注册
lighthouse@VM-8-10-ubuntu:test$ ./registry
收到消息:你好 Island
【★,°:.☆( ̄▽ ̄)/$:.°★ 】那么本篇到此就结束啦,如果有不懂 和 发现问题的小伙伴可以在评论区说出来哦,同时我还会继续更新关于【】的内容,请持续关注我 !!

更多推荐
所有评论(0)