在这里插入图片描述

📃个人主页: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 框架

对比项BRPCGRPC
性能更高(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库时引起冲突

image-20250912150832593

问题解决:直接卸载冲突版本解决问题

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 调用实现样例

服务端:

  1. 创建rpc服务子类继承pb中的EchoService服务类,并实现内部的业务接口逻辑创建rpc服务器类,
  2. 搭建服务器
  3. 向服务器类中添加 rpc子服务对象 – 告诉服务器收到什么请求用哪个接口处理
  4. 启动服务器

客户端:

  1. 创建网络通信信道
  2. 实例化pb中的 EchoService_Stub 类对象
  3. 发起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
};

职责

  1. 声明关注的服务(declared)
  2. 管理所有服务的 ServiceChannel(_services)
  3. 接收服务上下线通知(onServiceOnline / onServiceOffline)
  4. 为调用方提供“按服务名选择信道”的接口(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

【★,°:.☆( ̄▽ ̄)/$:.°★ 】那么本篇到此就结束啦,如果有不懂 和 发现问题的小伙伴可以在评论区说出来哦,同时我还会继续更新关于【】的内容,请持续关注我 !!

在这里插入图片描述

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐