在这里插入图片描述

Ⅰ. brpc的介绍

官方文档

brpc 是用 c++ 语言编写的工业级 RPC 框架(远程调用框架),常用于搜索、存储、机器学习、广告、推荐等高性能系统

​ 你可以使用它:

  1. 搭建能在一个端口支持多协议的服务, 或访问各种服务
    • restful http/httpsh2/gRPC。使用 brpchttp 实现比libcurl方便多了。从其他语言通过 HTTP/h2+json 访问基于 protobuf 的协议
    • redismemcached,线程安全,比官方 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开源
  2. Server同步异步处理请求
  3. Client 支持同步异步半同步,或使用组合channels简化复杂的分库或并发访问
  4. 通过http界面调试服务, 使用cpu, heap, contention profilers
  5. 获得更好的延时和吞吐
  6. 把你组织中使用的协议快速地加入brpc,或定制各类组件, 包括命名服务 (dns, zk, etcd), 负载均衡 (rr, random, consistent hashing)

对比其他的rpc框架

在这里插入图片描述

​ 下图是同机单 clientserver 在不同请求下的 QPS 指数:
在这里插入图片描述

​ 还有其他环境的对比,这里就不展示,可以到文档中查看,下面是结论:

  • brpc:在吞吐,平均延时,长尾处理上都表现优秀。
  • UB:平均延时和长尾处理的表现都不错,吞吐的扩展性较差,提高线程数和 client 数几乎不能提升吞吐。
  • thrift:单机的平均延时和吞吐尚可,多机的平均延时明显高于 brpcUB。吞吐的扩展性较差,提高线程数和 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 类中有 ClosureRpcController 类!

​ 其中 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 实现的服务发现和注册功能结合起来!通过注册中心,我们能够得知谁提供了什么服务,进而通过创建对应服务的信道来发起服务调用

​ 所以下面我们要封装:

  1. 指定服务的信道管理类
    • 一个服务可能会有多个节点提供服务,每个节点都有自己的 channel,建立服务与信道的映射关系,并且关系是一对多,采用 RR 轮转策略进行获取。
  2. 总体的服务通道管理类
    • 将上述指定服务的信道管理类都给管理起来!提供 服务上线与下线处理接口,以及提供 进行服务声明的接口,因为在整个系统中,提供的服务有很多,但是当前可能并不一定会用到所有的服务,因此通过声明来告诉模块哪些服务是自己关心的,需要建立连接管理起来,没有添加声明的服务即使上线也不需要进行连接的建立。
#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);
}

在这里插入图片描述

在这里插入图片描述

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐