【c++中间件】brpc远程调用框架 && 二次封装
文章目录
Ⅰ. brpc的介绍
brpc 是用 c++ 语言编写的工业级 RPC 框架(远程调用框架),常用于搜索、存储、机器学习、广告、推荐等高性能系统。
你可以使用它:
- 搭建能在一个端口支持多协议的服务, 或访问各种服务
restful http/https,h2/gRPC。使用brpc的http实现比libcurl方便多了。从其他语言通过HTTP/h2+json访问基于protobuf的协议- redis 和 memcached,线程安全,比官方
client更方便 - rtmp/flv/hls,可用于搭建流媒体服务
hadoop_rpc(可能开源)- 支持rdma(即将开源)
- 支持thrift,线程安全,比官方
client更方便 - 各种百度内使用的协议: baidu_std, streaming_rpc,
hulu_pbrpc, sofa_pbrpc,nova_pbrpc,public_pbrpc,ubrpc和使用nshead的各种协议 - 基于工业级的RAFT算法实现搭建高可用分布式系统,已在braft开源
Server能同步或异步处理请求Client支持同步、异步、半同步,或使用组合channels简化复杂的分库或并发访问- 通过http界面调试服务, 使用cpu, heap, contention profilers
- 获得更好的延时和吞吐
- 把你组织中使用的协议快速地加入brpc,或定制各类组件, 包括命名服务 (
dns,zk,etcd), 负载均衡 (rr,random,consistent hashing)
对比其他的rpc框架

下图是同机单 client 单 server 在不同请求下的 QPS 指数:

还有其他环境的对比,这里就不展示,可以到文档中查看,下面是结论:
brpc:在吞吐,平均延时,长尾处理上都表现优秀。UB:平均延时和长尾处理的表现都不错,吞吐的扩展性较差,提高线程数和client数几乎不能提升吞吐。thrift:单机的平均延时和吞吐尚可,多机的平均延时明显高于brpc和UB。吞吐的扩展性较差,提高线程数和client数几乎不能提升吞吐。sofa-pbrpc:处理小包的吞吐尚可,大包的吞吐显著低于其他RPC,延时受长尾影响很大。hulu-pbrpc:单机表现和sofa-pbrpc类似,但多机的延时表现极差。gRPC:几乎在所有参与的测试中垫底,可能它的定位是给google cloud platform的用户提供一个多语言,对网络友好的实现,性能还不是要务。
Ⅱ. brpc的安装与使用
一、安装
先安装依赖:
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
二、类与接口介绍
日志输出类与接口
包含头文件: #include <butil/logging.h>
本质上其实我们用不着 brpc 的日志输出,因此在这里主要介绍如何关闭日志输出!
namespace logging {
// 该枚举类型用于指定日志的输出目的地,不同的枚举值代表不同的输出目标
enum LoggingDestination {
LOG_TO_NONE = 0 // 当将日志输出目标设置为该值时,不会产生任何日志输出
};
// 日志设置结构体,用于配置日志系统的相关参数
struct BUTIL_EXPORT LoggingSettings {
LoggingSettings();
LoggingDestination logging_dest;
};
// 初始化日志系统
bool InitLogging(const LoggingSettings& settings);
}
protobuf类与接口
protobuf 类中有 Closure 和 RpcController 类!
其中 Closure 类是一个抽象接口,定义了一个可以被调用执行的操作。它通常用于 异步操作完成时的回调机制(因为我们项目用的是同步操作,所以这里无需关注太多细节),当某个异步任务完成后,会调用 Closure 对象的 Run() 方法来执行相应的回调逻辑。
而 RpcController 类在 RPC 调用过程中起着 控制和状态报告的作用。
namespace google {
namespace protobuf {
// 通常用于异步操作完成时的回调机制(后面有个ClosureGuard来保护内存泄漏问题)
class PROTOBUF_EXPORT Closure
{
public:
Closure() {}
virtual ~Closure();
virtual void Run() = 0; // 当异步任务完成时,会调用此方法来执行回调逻辑。
};
// 此函数用于创建一个 Closure 对象,该对象会在其 Run 方法被调用时执行传入的函数
inline Closure* NewCallback(void (*function)());
// 提供了一些方法来检查 RPC 调用的状态,例如是否失败以及失败的错误信息
class PROTOBUF_EXPORT RpcController
{
public:
bool Failed(); // 检查 RPC 调用是否失败
std::string ErrorText(); // 获取 RPC 调用失败的错误信息
}
}
}
服务端类与接口
这里只介绍主要用到的成员与接口!
namespace brpc {
/**
* @brief 服务器启动选项结构体,用于配置服务器的行为。
*
* 该结构体包含了一系列用于配置 brpc::Server 的选项,可在启动服务器时指定。
*/
struct ServerOptions {
/**
* @brief 连接的空闲超时时间(秒)。
*
* 如果在指定的时间内没有数据传输,连接将被关闭。默认值为 -1,表示禁用此功能。
*/
int idle_timeout_sec;
/**
* @brief 服务器使用的线程数量。
*
* 指定服务器处理请求时使用的线程数量,默认值为 CPU 核心数。
*/
int num_threads;
// .... 可能还有其他未显示的选项
};
/**
* @brief 服务所有权枚举,用于指定服务器对服务对象的管理方式。
*
* 该枚举定义了在将服务添加到服务器时,服务器对服务对象的所有权和生命周期管理策略。
*/
enum ServiceOwnership {
/**
* @brief 服务器拥有服务对象的所有权。
*
* 当添加服务失败时,服务器将负责删除服务对象,管理其生命周期。
*/
SERVER_OWNS_SERVICE,
/**
* @brief 服务器不拥有服务对象的所有权。
*
* 即使添加服务失败,服务器也不会删除服务对象,由调用者负责管理其生命周期。
*/
SERVER_DOESNT_OWN_SERVICE
};
/**
* @brief 代表一个 RPC 服务器。
*
* 该类提供了管理服务器生命周期、添加服务以及启动和停止服务器的方法。
*/
class Server {
public:
/**
* @brief 向服务器添加一个服务。
*
* @param service 指向要添加的 google::protobuf::Service 对象的指针。
* @param ownership 服务所有权策略,指定服务器对服务对象的管理方式。
* @return 成功时返回 0,失败时返回-1。
*/
int AddService(google::protobuf::Service* service,
ServiceOwnership ownership);
/**
* @brief 启动服务器。
*
* @param port 服务器监听的端口号。
* @param opt 指向 ServerOptions 结构体的指针,用于配置服务器的行为。如果为 nullptr,则使用默认选项。
* @return 成功时返回 0,失败时返回非零错误码。
*/
int Start(int port, const ServerOptions* opt);
/**
* @brief 停止服务器。
*
* @param closewait_ms 关闭等待时间(毫秒),此参数已不再使用。
* @return 成功时返回 0,失败时返回非零错误码。
*/
int Stop(int closewait_ms/*not used anymore*/);
/**
* @brief 等待服务器停止。
*
* 阻塞当前线程,直到服务器完全停止。
* @return 成功时返回 0,失败时返回非零错误码。
*/
int Join();
/**
* @brief 运行服务器直到接收到退出信号。
*
* 该方法会使服务器进入运行状态,并休眠直到接收到 Ctrl+C 信号,或者调用 Stop 和 Join 方法停止服务器。
*/
void RunUntilAskedToQuit();
};
/**
* @brief 用于自动管理 google::protobuf::Closure 对象的生命周期。
*
* 该类在构造时接受一个 google::protobuf::Closure 指针,并在析构时自动调用其 Run 方法,确保回调操作被执行。
*/
class ClosureGuard {
public:
explicit ClosureGuard(google::protobuf::Closure* done);
/**
* @brief 析构函数,自动调用 Closure 对象的 Run 方法。
*
* 如果 _done 指针不为空,则调用其 Run 方法执行回调操作。
*/
~ClosureGuard() { if (_done) _done->Run(); }
private:
google::protobuf::Closure* _done;
};
/**
* @brief 表示 HTTP 请求或响应的头部信息。
*
* 该类提供了设置和获取 HTTP 头部字段、URI、方法和状态码的方法。
*/
class HttpHeader {
public:
void set_content_type(const std::string& type);
const std::string* GetHeader(const std::string& key);
void SetHeader(const std::string& key,
const std::string& value);
const URI& uri() const { return _uri; }
HttpMethod method() const { return _method; }
void set_method(const HttpMethod method);
int status_code();
void set_status_code(int status_code);
private:
URI _uri;
HttpMethod _method;
};
/**
* @brief 用于控制和报告 RPC 调用的状态。
*
* 该类继承自 google::protobuf::RpcController,并提供了额外的方法来设置超时时间、最大重试次数,
* 以及获取 HTTP 请求和响应的头部信息。
*/
class Controller : public google::protobuf::RpcController {
public:
/**
* @brief 设置 RPC 调用的超时时间(毫秒)。
*
* @param timeout_ms 要设置的超时时间(毫秒)。
*/
void set_timeout_ms(int64_t timeout_ms);
/**
* @brief 设置 RPC 调用的最大重试次数。
*
* @param max_retry 要设置的最大重试次数。
*/
void set_max_retry(int max_retry);
/**
* @brief 获取 RPC 调用的响应消息。
*
* @return 指向 google::protobuf::Message 对象的指针,该对象包含 RPC 调用的响应消息。
*/
google::protobuf::Message* response();
/**
* @brief 获取 HTTP 响应的头部信息。
*
* @return 一个引用,指向存储 HTTP 响应头部信息的 HttpHeader 对象。
*/
HttpHeader& http_response();
/**
* @brief 获取 HTTP 请求的头部信息。
*
* @return 一个引用,指向存储 HTTP 请求头部信息的 HttpHeader 对象。
*/
HttpHeader& http_request();
/**
* @brief 检查 RPC 调用是否失败。
*
* @return 如果 RPC 调用失败,返回 true;否则返回 false。
*/
bool Failed();
/**
* @brief 获取 RPC 调用失败的错误信息。
*
* @return 一个包含错误信息的字符串。
*/
std::string ErrorText();
/**
* @brief 定义 RPC 响应后的回调函数类型。
*
* 该类型的函数将在 RPC 调用完成并收到响应后被调用,可用于处理响应结果。
*/
using AfterRpcRespFnType = std::function<
void(Controller* cntl,
const google::protobuf::Message* req,
const google::protobuf::Message* res)>;
/**
* @brief 设置 RPC 响应后的回调函数。
*
* @param fn 一个 AfterRpcRespFnType 类型的函数对象,将在 RPC 调用完成并收到响应后被调用。
*/
void set_after_rpc_resp_fn(AfterRpcRespFnType&& fn);
};
}
其中上述 Server 中新增服务类的第一个参数 google::protobuf::Service* service,实际上是后面我们会在 .proto 文件中添加的 rpc 服务后生成的一个类 EchoService 的实现类比如说后面我们会实现的 EchoServiceImpl 类,然后这个类继承于 EchoService,实际上我们就要将这个 EchoServiceImpl 的指针作为参数传过去,表示 新增一个我们指定的服务!如下所示:
syntax = "proto3";
package example;
option cc_generic_services = true; // 需要继承RPC服务类就需要cc_generic_services设置开启才能生成
message EchoRequest {
string message = 1;
}
message EchoResponse {
string message = 1;
}
// 定义一个rpc服务
service EchoSerive {
rpc Echo(EchoRequest) returns (EchoResponse);
}

此外我们需要在 【服务端】代码中重写上述 EchoService 类中的 Echo 接口,完成对应的服务,这个会在下面实现的时候讲!
客户端类与接口
namespace brpc {
/**
* @brief 用于配置 brpc 通道的选项结构体
*
* 该结构体包含了一系列参数,用于控制通过 brpc 通道发起的 RPC 调用的行为,
* 例如连接超时时间、请求超时时间、重试次数以及序列化协议类型等。
*/
struct ChannelOptions {
/**
* @brief 请求连接超时时间,单位为毫秒。
*
* 当通过通道发起请求时,若在该时间内未能成功建立与服务器的连接,
* 则认为连接超时,请求将失败。默认值为 200 毫秒。
*/
int32_t connect_timeout_ms;
/**
* @brief RPC 请求超时时间,单位为毫秒。
*
* 从请求发送开始计时,若在该时间内未收到服务器的响应,
* 则认为请求超时,请求将失败。默认值为 500 毫秒。
*/
int32_t timeout_ms;
/**
* @brief 最大重试次数。
*
* 当请求失败时,通道会尝试重新发送请求,最多重试该指定的次数。
* 默认值为 3 次。
*/
int max_retry;
/**
* @brief 序列化协议类型。
*
* 指定在进行 RPC 通信时所使用的序列化协议,例如可以设置为 "baidu_std" 等。
* 不同的协议在数据格式、性能等方面可能存在差异。
*/
AdaptiveProtocolType protocol;
//....其他未显示的选项
};
/**
* @brief 代表一个 brpc 通道,用于与服务器进行 RPC 通信。
*
* 该类继承自 ChannelBase,提供了初始化通道的方法,以便后续通过该通道发起 RPC 请求。
*/
class Channel : public ChannelBase {
public:
/**
* @brief 初始化通道。
*
* 该方法用于配置并启动通道,使其可以与指定的服务器进行通信。
*
* @param server_addr_and_port 服务器的地址和端口信息,格式通常为 "ip:port"。
* @param options 指向 ChannelOptions 结构体的指针,用于指定通道的配置选项。
* 如果为 nullptr,则使用默认的配置选项。
* @return 初始化成功时返回 0;若初始化过程中出现错误,则返回非零的错误码。
*/
int Init(const char* server_addr_and_port,
const ChannelOptions* options);
};
}

在 EchoService 中有个子类 EchoService_Stub,其实是客户端使用的一个类,通过调用 Echo 接口实现请求和响应的处理,我们可以设置同步和异步调用方式,这个下面会讲!所以要注意与 EchoService 区分开!
三、同步调用
同步调用是指 客户端会阻塞收到 server 端的响应或发生错误。
下面我们以 Echo(输出 hello world)方法为例来讲解基础的同步 RPC 请求是如何实现的!
.proto 文件:
syntax = "proto3";
package example;
option cc_generic_services = true; // 需要继承RPC服务类就需要cc_generic_services设置开启才能生成
message EchoRequest {
string message = 1;
}
message EchoResponse {
string message = 1;
}
// 定义RPC远端方法
service EchoService {
rpc Echo(EchoRequest) returns (EchoResponse);
}
在 rpc 框架中,EchoService 是一个服务定义,它仅仅描述了服务的接口,即有哪些 RPC 方法以及这些方法的输入输出消息类型。而在服务端,我们需要实现这些接口,定义每个 RPC 方法具体要执行的操作,这样当客户端调用这些方法时,服务端才能正确处理请求并返回响应。通过创建 EchoServiceImpl 类并继承 EchoService,我们可以将服务的接口定义和具体实现分离开来,使得代码结构更加清晰,易于维护和扩展。
重写 Echo() 函数,主要就是处理请求和发送响应,然后调用 Run() 接口(我们会用类似守卫锁来让它自动调用)。
server.cc 文件:
#include <brpc/server.h>
#include <butil/logging.h>
#include "main.pb.h"
// 别忘了要声明命名空间example
class EchoServiceImpl : public example::EchoService {
public:
EchoServiceImpl() {}
~EchoServiceImpl() {}
// 重写Echo接口,实现具体的业务逻辑,以响应客户端发起的 Echo RPC 调用
void Echo(::google::protobuf::RpcController* controller,
const ::example::EchoRequest* request,
::example::EchoResponse* response,
::google::protobuf::Closure* done)
{
// 1. 使用RAII方式自动释放done对象,内部会调用run()接口
brpc::ClosureGuard guard(done);
// 2. 进行请求输出和响应
std::cout << "收到消息:" << request->message() << std::endl;
response->set_message("这是一条用于测试的响应信息!");
}
};
int main(int argc, char* argv[])
{
// 1. 关闭rpc自带的日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
// 2. 创建服务器对象
brpc::Server server;
// 3. 向服务器对象中新增EchoService服务(局部变量,不需要去服务器去清除服务对象)
EchoServiceImpl impl;
int ret = server.AddService(&impl, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if(ret != 0) {
std::cout << "添加rpc服务失败!\n";
return -1;
}
// 4. 启动服务器
brpc::ServerOptions options;
options.num_threads = 1;
ret = server.Start(8080, &options);
if(ret != 0) {
std::cout << "启动rpc服务器失败!\n";
return -1;
}
server.RunUntilAskedToQuit(); // 休眠等待运行结束
return 0;
}
而客户端只需要调用 example::EchoService_Stub 中的 Echo() 接口即可,然后根据需要,如果要用异步调用的话,则需要再写一个回调函数传给 Echo() 接口来处理,这个看后面的异步调用方式!
client.cc 文件:
#include <brpc/channel.h>
#include <thread>
#include "main.pb.h"
int main(int argc, char* argv[])
{
// 1. 构造Channel信道,连接服务器
brpc::Channel channel;
int ret = channel.Init("127.0.0.1:8080", nullptr);
if(ret != 0) {
std::cout << "初始化信道失败!\n";
return -1;
}
// 2. 构造EchoService_Stub对象,用于进行rpc调用
example::EchoService_Stub stub(&channel);
// 3. 进行Rpc调用(同步调用)
example::EchoRequest req;
req.set_message("你好,利刃!");
brpc::Controller *ctl = new brpc::Controller();
example::EchoResponse *resp = new example::EchoResponse();
stub.Echo(ctl, &req, resp, nullptr); // 同步方式调用,不需要Closure对象
if(ctl->Failed() == true) {
std::cout << "Rpc调用失败:" << ctl->ErrorText() << std::endl;
return -1;
}
std::cout << "收到响应: " << resp->message() << std::endl;
// 别忘了释放资源
delete ctl;
delete resp;
std::this_thread::sleep_for(std::chrono::seconds(3));
return 0;
}
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
四、异步调用
异步调用是指客户端注册一个响应处理回调函数, 当调用一个 RPC 接口时立即返回,不会阻塞等待响应, 当 server 端返回响应时会调用传入的回调函数处理响应。
#include <brpc/channel.h>
#include <thread>
#include "main.pb.h"
void callback(brpc::Controller* ctl, example::EchoResponse* resp)
{
std::unique_ptr<brpc::Controller> ctl_guard(ctl);
std::unique_ptr<example::EchoResponse> resp_guard(resp);
if (ctl->Failed() == true) {
std::cout << "Rpc调用失败:" << ctl->ErrorText() << std::endl;
return;
}
std::cout << "收到响应: " << resp->message() << std::endl;
}
int main(int argc, char* argv[])
{
// 1. 构造Channel信道,连接服务器
brpc::Channel channel;
int ret = channel.Init("127.0.0.1:8080", nullptr);
if(ret != 0) {
std::cout << "初始化信道失败!\n";
return -1;
}
// 2. 构造EchoService_Stub对象,用于进行rpc调用
example::EchoService_Stub stub(&channel);
// 3. 进行Rpc调用(异步调用)
example::EchoRequest req;
req.set_message("你好,利刃!");
brpc::Controller *ctl = new brpc::Controller();
example::EchoResponse *resp = new example::EchoResponse();
// 设置回调函数,然后返回Closure对象指针
auto closure = google::protobuf::NewCallback(callback, ctl, resp);
stub.Echo(ctl, &req, resp, closure);
std::this_thread::sleep_for(std::chrono::seconds(3));
return 0;
}
Ⅲ. brpc的信道管理封装
rpc 调用这里的封装,因为 不同的服务调用使用的是不同的 Stub,这个封装起来的意义不大,因此我们 只需要封装通信所需的 Channel 管理,这样当需要进行什么样的服务调用的时候,只需要通过服务名称获取对应的 channel,然后实例化 Stub 进行调用即可。
注意 rpc 调用要和我们前面学到的 etcd 实现的服务发现和注册功能结合起来!通过注册中心,我们能够得知谁提供了什么服务,进而通过创建对应服务的信道来发起服务调用!
所以下面我们要封装:
- 指定服务的信道管理类
- 一个服务可能会有多个节点提供服务,每个节点都有自己的
channel,建立服务与信道的映射关系,并且关系是一对多,采用RR轮转策略进行获取。
- 一个服务可能会有多个节点提供服务,每个节点都有自己的
- 总体的服务通道管理类
- 将上述指定服务的信道管理类都给管理起来!提供 服务上线与下线处理接口,以及提供 进行服务声明的接口,因为在整个系统中,提供的服务有很多,但是当前可能并不一定会用到所有的服务,因此通过声明来告诉模块哪些服务是自己关心的,需要建立连接管理起来,没有添加声明的服务即使上线也不需要进行连接的建立。
#pragma once
#include <brpc/channel.h>
#include <string>
#include <vector>
#include <unordered_map>
#include <mutex>
#include "logger.hpp"
// 单个服务的信道管理类
class ChannelManager {
public:
using channel_ptr = std::shared_ptr<brpc::Channel>;
using ptr = std::shared_ptr<ChannelManager>;
public:
ChannelManager(const std::string& service_name)
: _service_name(service_name)
, _index(0)
{}
// 服务上线了一个节点,则调用append新增信道
void append(const std::string& host)
{
// 构造初始化Channel信道
channel_ptr channel = std::make_shared<brpc::Channel>();
brpc::ChannelOptions options;
options.connect_timeout_ms = -1;
options.timeout_ms = -1;
options.max_retry = 3;
options.protocol = "baidu_std";
int ret = channel->Init(host.c_str(), &options);
if(ret != 0) {
LOG_ERROR("初始化 {}-{} 信道失败!", _service_name, host);
return;
}
// 先加锁再添加信息
std::unique_lock<std::mutex> lock(_mtx);
_channels.push_back(channel);
_hosts[host] = channel;
}
// 服务下线了一个节点,则调用remove释放信道
void remove(const std::string& host)
{
std::unique_lock<std::mutex> lock(_mtx);
auto it = _hosts.find(host);
if(it == _hosts.end())
{
LOG_WARN("删除 {}-{} 信道失败,没有找到该信道!", _service_name, host);
return;
}
for(auto ait = _channels.begin(); ait != _channels.end(); ++ait)
if(*ait == it->second)
{
_channels.erase(ait);
break;
}
_hosts.erase(it);
}
// 通过RR轮转策略,获取一个Channel用于发起对应服务的rpc调用
channel_ptr get()
{
std::unique_lock<std::mutex> lock(_mtx);
if(_channels.size() == 0)
{
LOG_ERROR("当前无信道可用,已为你创建新信道");
return channel_ptr();
}
int32_t idx = _index++ % _channels.size();
return _channels[idx];
}
private:
std::mutex _mtx;
int32_t _index; // 当前轮转下标计数器
std::string _service_name; // 服务名
std::vector<channel_ptr> _channels; // 存放channel的集合
std::unordered_map<std::string, channel_ptr> _hosts; // 存放主机号与channel的映射关系
};
// 总体服务的信道管理类
class ServiceManager {
public:
using channel_ptr = std::shared_ptr<brpc::Channel>;
using ptr = std::shared_ptr<ServiceManager>;
public:
// 获取对应服务的一个channel对象,用于rpc调用
ChannelManager::channel_ptr getChannel(const std::string& service_name)
{
std::unique_lock<std::mutex> lock(_mtx);
auto it = _services.find(service_name);
if(it == _services.end())
{
LOG_ERROR("没有提供 {} 服务的节点", service_name);
return ChannelManager::channel_ptr();
}
return it->second->get();
}
// 声明关注哪些服务的上下线调用,不关注的服务不需要处理
void declared(const std::string& service_name)
{
std::unique_lock<std::mutex> lock(_mtx);
_follow_services.insert(service_name);
}
// 服务上线时调用的回调接口(即etcd.hpp中Discovery类的put_cb对象):为服务添加主机地址
void online(const std::string& instance_name, const std::string& host)
{
std::string service_name = instance_to_service(instance_name);
ChannelManager::ptr service;
{
std::unique_lock<std::mutex> lock(_mtx);
auto fit = _follow_services.find(service_name);
if(fit == _follow_services.end())
{
LOG_DEBUG("{}-{} 服务上线了,但是当前并不关心!\n", service_name, host);
return;
}
auto it = _services.find(service_name);
if(it == _services.end())
{
// 说明是新添加的服务节点,此时创建并且插入即可
service = std::make_shared<ChannelManager>(service_name);
_services[service_name] = service;
}
else
service = it->second;
}
if (!service)
{
LOG_ERROR("新增 {} 服务管理节点失败!", service_name);
return ;
}
service->append(host);
LOG_DEBUG("{}-{} 服务上线新节点,进行添加管理!", service_name, host);
}
// 服务下线时调用的回调接口(即etcd.hpp中Discovery类的del_cb对象):删除服务节点
void offline(const std::string& instance_name, const std::string& host)
{
std::string service_name = instance_to_service(instance_name);
ChannelManager::ptr service;
{
std::unique_lock<std::mutex> lock(_mtx);
auto fit = _follow_services.find(service_name);
if(fit == _follow_services.end())
{
LOG_DEBUG("{}-{} 服务下线了,但是当前并不关心!", service_name, host);
return;
}
auto it = _services.find(service_name);
if(it == _services.end())
{
LOG_WARN("删除服务节点失败:没找到{}节点", service_name);
return;
}
service = it->second;
}
service->remove(host);
LOG_DEBUG("{}-{} 服务下线节点,进行删除管理!", service_name, host);
}
private:
// 将实例名转化为服务名(实际上就是去掉最后一个'/'后面的内容
std::string instance_to_service(const std::string& instance)
{
auto pos = instance.find_last_of('/');
if (pos == std::string::npos) return instance;
return instance.substr(0, pos);
}
private:
std::mutex _mtx;
std::unordered_set<std::string> _follow_services; // 存放关注的服务集合
std::unordered_map<std::string, ChannelManager::ptr> _services; // 存放服务名和ChannelManager映射的集合
};
与 etcd 的联调
下面我们结合注册服务以及注册发现,也就是 etcd 来实现服务上下线的处理!
这部分实际上就是我们前面 brpc 接口的使用样例,以及 etcd 封装的测试样例结合起来的样例!
discovery.cc 文件:
// discovery.cc
#include <gflags/gflags.h>
#include "../header/etcd.hpp"
#include "../header/channel.hpp"
#include "main.pb.h"
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(basedir, "/service", "服务监控的主目录");
DEFINE_string(call_service, "/service/echo", "测试代码中要获取的提供echo服务节点");
int main(int argc, char* argv[])
{
// 初始化gflags
google::ParseCommandLineFlags(&argc, &argv, true);
init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 1、构造rpc信道管理对象,然后关注“/service/echo”这个服务
ServiceManager::ptr sm = std::make_shared<ServiceManager>();
sm->declared(FLAGS_call_service);
// 2、绑定服务上下线的回调处理函数,然后传入回调函数构造服务发现对象Discovery
auto put_cb = std::bind(&ServiceManager::online, sm.get(), std::placeholders::_1, std::placeholders::_2);
auto del_cb = std::bind(&ServiceManager::offline, sm.get(), std::placeholders::_1, std::placeholders::_2);
Discovery::ptr dis = std::make_shared<Discovery>(FLAGS_etcd_host, FLAGS_basedir, put_cb, del_cb);
while(true)
{
// 3、通过rpc信道管理对象,获取提供Echo服务的信道对象
auto channel = sm.get()->getChannel(FLAGS_call_service);
if(!channel)
{
std::this_thread::sleep_for(std::chrono::seconds(1));
return -1;
}
// 4、构造EchoService_Stub对象,用于进行rpc调用(同步调用)
example::EchoService_Stub stub(channel.get());
example::EchoRequest req;
req.set_message("你好,利刃!");
brpc::Controller *ctl = new brpc::Controller();
example::EchoResponse *resp = new example::EchoResponse();
stub.Echo(ctl, &req, resp, nullptr); // 同步方式调用,不需要Closure对象
if(ctl->Failed() == true) {
std::cout << "Rpc调用失败:" << ctl->ErrorText() << std::endl;
delete ctl;
delete resp;
std::this_thread::sleep_for(std::chrono::seconds(1));
continue;
}
std::cout << "收到响应: " << resp->message() << std::endl;
std::this_thread::sleep_for(std::chrono::seconds(1));
}
return 0;
}
registery.cc 文件:
// registery.cc
#include <brpc/server.h>
#include <butil/logging.h>
#include <gflags/gflags.h>
#include "main.pb.h"
#include "../header/etcd.hpp"
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(basedir, "/service", "服务监控的主目录");
DEFINE_string(instance, "/echo/instance", "当前实例名称");
DEFINE_string(host, "127.0.0.1:8080", "服务的主机地址");
DEFINE_int32(listen_port, 8080, "rpc服务器的监听端口");
// 别忘了要声明命名空间example
class EchoServiceImpl : public example::EchoService {
public:
EchoServiceImpl() {}
~EchoServiceImpl() {}
// 重写Echo接口,实现具体的业务逻辑,以响应客户端发起的 Echo RPC 调用
void Echo(::google::protobuf::RpcController* controller,
const ::example::EchoRequest* request,
::example::EchoResponse* response,
::google::protobuf::Closure* done)
{
// 1. 使用RAII方式自动释放done对象,内部会调用run()接口
brpc::ClosureGuard guard(done);
// 2. 进行请求输出和响应
std::cout << "收到消息:" << request->message() << std::endl;
response->set_message(request->message() + "这是一条用于测试的响应信息!");
}
};
int main(int argc, char* argv[])
{
// 初始化gflags
google::ParseCommandLineFlags(&argc, &argv, true);
init_logger(FLAGS_run_mode, FLAGS_log_file, FLAGS_log_level);
// 1. 关闭rpc自带的日志输出
logging::LoggingSettings settings;
settings.logging_dest = logging::LoggingDestination::LOG_TO_NONE;
logging::InitLogging(settings);
// 2. 创建服务器对象
brpc::Server server;
// 3. 向服务器对象中新增EchoService服务(局部变量,不需要去服务器去清除服务对象)
EchoServiceImpl impl;
int ret = server.AddService(&impl, brpc::ServiceOwnership::SERVER_DOESNT_OWN_SERVICE);
if(ret != 0) {
std::cout << "添加rpc服务失败!\n";
return -1;
}
// 4. 启动服务器
brpc::ServerOptions options;
options.num_threads = 1;
ret = server.Start(FLAGS_listen_port, &options);
if(ret != 0) {
std::cout << "启动rpc服务器失败!\n";
return -1;
}
// 5、注册服务(创建注册类对象并用智能指针保护起来)
Registry::ptr reg = std::make_shared<Registry>(FLAGS_etcd_host);
reg->regiter(FLAGS_basedir + FLAGS_instance, FLAGS_host);
server.RunUntilAskedToQuit(); // 休眠等待运行结束
return 0;
}
makefile 文件:
all : registry discovery
registry : registry.cc main.pb.cc
g++ -std=c++17 -o $@ $^ -lspdlog -lfmt -lgflags -letcd-cpp-api -lcpprest -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
discovery : discovery.cc main.pb.cc
g++ -std=c++17 -o $@ $^ -lspdlog -lfmt -lgflags -letcd-cpp-api -lcpprest -lbrpc -lssl -lcrypto -lprotobuf -lleveldb
.PHONY:clean
clean:
rm -f discovery registry
然后 protobuf 文件就用我们在前面样例中使用的即可:
syntax = "proto3";
package example;
option cc_generic_services = true; // 需要继承RPC服务类就需要cc_generic_services设置开启才能生成
message EchoRequest {
string message = 1;
}
message EchoResponse {
string message = 1;
}
// 定义RPC远端方法
service EchoService {
rpc Echo(EchoRequest) returns (EchoResponse);
}


更多推荐


所有评论(0)