1. 为什么选择libhv和Protobuf来构建RPC服务器?

如果你正在用C++写网络服务,大概率遇到过这些头疼事:要自己处理TCP连接、管理线程池、解决粘包拆包、还得设计一套协议格式。光是这些底层细节就能耗掉你大半的开发时间,更别提还要保证高性能和高并发了。我自己以前也踩过不少坑,直到发现了libhv这个宝藏库,才真正体会到什么叫“把复杂留给自己,把简单留给用户”。

libhv是一个用纯C语言编写的高性能网络库,但它提供了完整的C++封装。我最喜欢它的evpp模块,你可以把它理解为一个C++风格的、更易用的libevent或者libuv。它帮你封装好了事件循环、TCP/UDP服务器、定时器这些基础设施,让你能像写业务逻辑一样轻松地构建网络应用。最关键的是,它的性能非常出色,我在实际项目里测过,单机轻松扛住上万并发连接不是问题。

而Protobuf,我相信做后端的朋友都不陌生。它是Google出品的一套序列化协议,简单来说就是能把你的C++结构体对象变成一串紧凑的二进制数据,方便在网络上传送或者存到文件里。相比JSON和XML,Protobuf生成的二进制数据体积小、解析速度快,天生就适合做RPC这种对性能要求高的场景。你只需要写一个.proto文件定义好数据格式,protoc编译器就能自动为你生成C++(或者其他语言)的读写代码,省去了手动组包解包的麻烦,也杜绝了字段拼写错误这种低级bug。

把这两者结合起来——用libhv处理高并发的网络I/O,用Protobuf处理高效、严谨的数据编解码——就成了构建现代RPC服务器的黄金组合。今天我要分享的,就是如何用大约300行C++代码,基于libhv的evpp模块,实现一个完整的、异步的Protorpc服务器。这个服务器是非阻塞的,能同时处理大量客户端请求,代码结构清晰,扩展起来也非常方便。无论你是想学习网络编程,还是急需一个轻量级的RPC框架用于自己的项目,我相信接下来的内容都会对你有帮助。

2. 环境准备与项目搭建

2.1 安装核心依赖:libhv和Protobuf

动手之前,我们得先把“厨房”收拾好。第一步就是安装libhv和Protobuf。这两个库的安装都很简单,我习惯从源码编译安装,这样最可控。

安装libhv: libhv的源码在GitHub上,直接克隆下来编译就行。它几乎没有外部依赖,编译速度很快。

git clone https://github.com/ithewei/libhv.git
cd libhv
# 使用CMake构建,这是现在更主流的方式
mkdir build && cd build
cmake ..
make -j4  # 根据你的CPU核心数调整,加快编译速度
sudo make install

安装完成后,libhv的头文件通常会放到/usr/local/include,库文件放到/usr/local/lib。你可以写个简单的测试程序,链接-lhv看看是否成功。

安装Protobuf: Protobuf的安装步骤稍微多几步,因为它需要生成一些配置脚本。

git clone https://github.com/protocolbuffers/protobuf.git
cd protobuf
# 注意,新版本可能需要先更新子模块
git submodule update --init --recursive
./autogen.sh  # 生成configure脚本
./configure
make -j4
sudo make install
sudo ldconfig  # 刷新系统的动态链接库缓存

安装完成后,检查一下protoc编译器是否可用:

which protoc
protoc --version

如果能看到版本号,说明安装成功了。protoc是我们后续将.proto文件生成C++代码的关键工具。

2.2 创建项目结构与定义协议

环境搞定后,我们创建一个干净的项目目录。我的习惯是这样的:

my_protorpc_project/
├── CMakeLists.txt          # 项目构建文件
├── proto/                  # 存放Protobuf协议定义
│   └── base.proto
├── include/                # 自己写的头文件
│   ├── handler.h
│   └── router.h
├── src/                   # 源代码
│   ├── main.cpp
│   ├── handler/
│   │   ├── calc.cpp
│   │   └── login.cpp
│   └── router.cpp
└── build/                 # 编译输出目录

接下来是重头戏:定义RPC的通信协议。我们在proto/base.proto文件中写下如下内容:

syntax = "proto3";
package protorpc;

message Error {
  int32 code = 1;
  string message = 2;
}

message Request {
  uint64 id = 1;          // 请求ID,用于匹配请求和响应
  string method = 2;      // 要调用的方法名,比如 "add", "login"
  repeated bytes params = 3; // 参数列表,每个参数本身是一个Protobuf bytes
}

message Response {
  uint64 id = 1;          // 对应请求的ID
  oneof result_or_error {
    bytes result = 2;     // 成功时的结果
    Error error = 3;      // 失败时的错误信息
  }
}

这个协议设计得非常通用。Request里的method字段告诉服务器要调用哪个函数,params是一个数组,可以传递任意多个用Protobuf序列化好的二进制参数。Response使用oneof语法,意味着一个响应要么包含结果(result),要么包含错误(error),不会同时存在。这比用两个可选字段更清晰。

然后用protoc生成C++代码:

cd my_protorpc_project
protoc --cpp_out=./src ./proto/base.proto

这条命令会在src目录下生成base.pb.cc和base.pb.h两个文件。这就是我们之后序列化和反序列化要用到的所有代码,全是自动生成的,非常省心。

3. 理解libhv的evpp模块与异步事件模型

在开始写服务器主循环之前,我们得先搞明白libhv的evpp模块是怎么工作的。如果你用过libevent或Boost.Asio,理解起来会很快。它的核心思想是事件驱动和非阻塞I/O。

想象一下你去银行办业务。同步阻塞的方式就像只有一个柜台,你必须在那里干等着直到业务办完,期间什么都做不了。而异步非阻塞的方式,就像你取个号就可以去旁边坐着玩手机,柜台叫到你的号时(事件触发)再去处理。这样,一个柜台(一个线程)就能服务很多人(很多连接)。

在libhv的evpp里,这个“叫号系统”就是事件循环(EventLoop)。一个EventLoop对象管理着一个线程,不断地检查有没有新的事件发生:比如新的连接进来了、某个socket有数据可读了、或者定时器时间到了。我们的TcpServer就是构建在这个事件循环之上的。

创建一个基本的异步TCP服务器,用evpp模块简单到难以置信:

#include <hv/TcpServer.h>

using namespace hv;

int main() {
    TcpServer srv;
    srv.onConnection = [](const SocketChannelPtr& channel) {
        // 有新的连接建立或断开时会进入这个回调
        std::string peeraddr = channel->peeraddr();
        if (channel->isConnected()) {
            printf("%s connected! connfd=%d\n", peeraddr.c_str(), channel->fd());
        } else {
            printf("%s disconnected! connfd=%d\n", peeraddr.c_str(), channel->fd());
        }
    };

    srv.onMessage = [](const SocketChannelPtr& channel, Buffer* buf) {
        // 当该连接上有数据可读时,会进入这个回调
        // buf里就是收到的原始数据
        // 我们在这里处理请求逻辑
    };

    srv.setPort(12345);
    srv.setThreadNum(4); // 设置工作线程数,通常建议和CPU核心数一致
    srv.start();

    // 主线程可以在这里做其他事情,或者简单地等待
    while (1) hv_sleep(1);
    return 0;
}

看,不到30行,一个支持多线程并发的TCP服务器就跑起来了。setThreadNum(4)意味着启动了4个事件循环线程(也就是4个“柜台”),每个线程独立运行,共同处理连接。新连接会被均衡地分配到不同的线程上。onMessage回调是性能的关键,它运行在I/O线程中,必须快速处理,不能有阻塞操作(比如长时间的文件读写、同步的网络请求),否则会拖慢整个事件循环。

4. 设计并实现异步Protorpc服务器核心

4.1 粘包拆包:网络通信的第一道关卡

网络编程新手最容易栽在“粘包”问题上。TCP是流式协议,它只保证数据顺序,不保证边界。客户端连续发送两个请求“Hello”和“World”,服务器一次recv可能收到“HelloWorld”,也可能先收到“Hel”,再收到“loWorld”。所以,我们需要在应用层定义包的边界,这就是“拆包”。

常见的拆包方式有:定长包、分隔符、或者在包头声明长度。我们的Protorpc采用第三种,也是最通用高效的方式。我们设计一个简单的协议头:

+----------------+----------------+----------------+
|   Magic (2B)   |   Length (4B)  |     Body       |
+----------------+----------------+----------------+
  • Magic:2字节的魔数,比如0xABCE,用于快速校验这是一个合法的包起始。
  • Length:4字节的网络序整数,表示后面Body部分的长度。
  • Body:实际的数据,也就是我们Protobuf序列化后的Request或Response。

libhv的TcpServer贴心地提供了setUnpack方法,让我们可以自定义拆包规则,它会在底层帮我们做好缓冲和切割,保证onMessage回调里收到的Buffer* buf一定是一个完整的包。配置如下:

// 在服务器初始化时设置
unpack_setting_t protorpc_unpack_setting;
memset(&protorpc_unpack_setting, 0, sizeof(unpack_setting_t));
protorpc_unpack_setting.mode = UNPACK_BY_LENGTH_FIELD; // 按长度字段拆包
protorpc_unpack_setting.package_max_length = DEFAULT_PACKAGE_MAX_LENGTH; // 最大包长,防攻击
protorpc_unpack_setting.body_offset = PROTORPC_HEAD_LENGTH; // Body在包中的起始位置
protorpc_unpack_setting.length_field_offset = PROTORPC_HEAD_LENGTH_FIELD_OFFSET; // 长度字段的偏移(2,跳过Magic)
protorpc_unpack_setting.length_field_bytes = PROTORPC_HEAD_LENGTH_FIELD_BYTES; // 长度字段占几个字节(4)
protorpc_unpack_setting.length_field_coding = ENCODE_BY_BIG_ENDIAN; // 网络字节序(大端)
srv.setUnpack(&protorpc_unpack_setting);

配置好之后,我们就不用再操心粘包问题了,可以专注于业务逻辑。

4.2 请求路由与处理器设计

服务器收到一个完整的请求包后,需要根据method字段找到对应的处理函数。这就是路由。一个清晰的路由机制能让代码易于维护和扩展。我通常用一个结构体数组来定义路由表:

// router.h
typedef void (*ProtorpcHandler)(const protorpc::Request& req, protorpc::Response* res);

struct ProtorpcRouter {
    const char* method;
    ProtorpcHandler handler;
};

// 在某个全局或静态区域定义路由表
extern ProtorpcRouter g_router[];
extern const size_t g_router_num;
// router.cpp
#include "handler/calc.h"
#include "handler/login.h"

ProtorpcRouter g_router[] = {
    {"add", calc_add},
    {"sub", calc_sub},
    {"mul", calc_mul},
    {"div", calc_div},
    {"login", login},
    // 新的方法在这里添加...
};

const size_t g_router_num = sizeof(g_router) / sizeof(g_router[0]);

处理器的签名是统一的:接收一个常量请求引用,和一个用于填充的响应指针。例如,一个加法处理器:

// handler/calc.cpp
#include "base.pb.h"

void calc_add(const protorpc::Request& req, protorpc::Response* res) {
    // 1. 反序列化参数。req.params() 是一个 repeated bytes
    // 2. 执行计算
    // 3. 将结果序列化成bytes,放入 res->mutable_result()
    // 4. 如果出错,则设置 res->mutable_error()
}

这种设计的好处是,路由逻辑与业务逻辑完全解耦。要新增一个API,只需要在g_router数组里加一项,然后实现对应的处理器函数即可。服务器的主循环完全不用动。

4.3 组装完整的服务器类

现在我们把所有零件组装起来。我们创建一个ProtoRpcServer类,继承自hv::TcpServer,在构造函数里完成所有初始化。

class ProtoRpcServer : public hv::TcpServer {
public:
    ProtoRpcServer() : TcpServer() {
        // 1. 设置连接回调
        onConnection = [](const SocketChannelPtr& channel) {
            // ... 日志记录连接/断开事件
        };

        // 2. 设置消息回调,指向我们的静态处理函数
        onMessage = &ProtoRpcServer::handleMessage;

        // 3. 配置拆包规则
        unpack_setting_t unpack_setting;
        // ... 填充上述拆包配置
        setUnpack(&unpack_setting);
    }

    int listen(int port) {
        return createsocket(port); // 创建监听socket
    }

private:
    // 静态成员函数,处理所有消息
    static void handleMessage(const SocketChannelPtr& channel, Buffer* buf) {
        // 这里是核心处理流程,下一步详细展开
    }
};

注意handleMessage被声明为static。因为onMessage回调需要的是一个普通函数指针或std::function,而普通的成员函数有隐式的this指针。将其设为静态,我们可以在函数内部通过其他方式访问全局或静态的路由表。

5. 核心流程剖析:从接收到响应的300行代码

现在,我们深入最核心的handleMessage函数。这不到100行的代码,完成了Protorpc服务器的全部工作。我会逐段解释,并分享一些我踩过坑的细节。

第一步:拆包与校验

static void handleMessage(const SocketChannelPtr& channel, Buffer* buf) {
    // 1. 拆包
    protorpc_message msg;
    memset(&msg, 0, sizeof(msg));
    int packlen = protorpc_unpack(&msg, buf->data(), buf->size());
    if (packlen < 0) {
        printf("protorpc_unpack failed!\n");
        return; // 拆包失败,可能是数据损坏,直接丢弃
    }
    assert(packlen == buf->size()); // 在调试模式下确认

    // 2. 校验包头(例如检查魔数)
    if (protorpc_head_check(&msg.head) != 0) {
        printf("protorpc_head_check failed!\n");
        return;
    }

protorpc_unpack是一个辅助函数,根据我们之前设置的规则,从buf中解析出包头和指向Body的指针。protorpc_head_check可以检查魔数是否正确,这是一个快速失败的手段,能尽早过滤掉非法的或损坏的数据包。

第二步:反序列化请求

    // 3. 反序列化Protobuf Request
    protorpc::Request req;
    protorpc::Response res;
    if (req.ParseFromArray(msg.body, msg.head.length)) {
        // 成功反序列化
        printf("> %s\n", req.DebugString().c_str()); // 打印请求日志,调试用
        res.set_id(req.id()); // 响应ID必须与请求ID一致

ParseFromArray是Protobuf生成的接口,它尝试从二进制数据msg.body中解析出Request对象。如果失败,说明客户端发送的数据不符合Request的消息格式,我们应当返回一个“Bad Request”错误。

第三步:路由与处理

        // 4. 路由查找
        const char* method = req.method().c_str();
        bool found = false;
        for (size_t i = 0; i < g_router_num; ++i) {
            if (strcmp(method, g_router[i].method) == 0) {
                found = true;
                // 找到处理器,调用它
                g_router[i].handler(req, &res);
                break;
            }
        }
        if (!found) {
            // 没有找到对应方法,返回404 Not Found
            not_found(req, &res);
        }
    } else {
        // 反序列化失败,返回400 Bad Request
        bad_request(req, &res);
    }

这里就是一个简单的线性查找。如果你的方法非常多(比如上百个),可以考虑用std::unordered_map<std::string, ProtorpcHandler>来提升查找效率。not_found和bad_request是辅助函数,负责填充Response中的error字段。

第四步:序列化响应并发送

    // 5. 序列化响应并封包
    protorpc_message_init(&msg); // 复用msg结构体,但清空body指针
    msg.head.length = res.ByteSize(); // 获取序列化后的长度
    int needed_len = protorpc_package_length(&msg.head); // 计算整个包(头+体)的长度

    // 动态分配内存来存放整个包
    unsigned char* writebuf = (unsigned char*)malloc(needed_len);
    if (writebuf == nullptr) {
        // 处理内存分配失败,在实际项目中这里应该记录日志并关闭连接
        return;
    }

    // 先打包,填充头部(Magic和Length)
    packlen = protorpc_pack(&msg, writebuf, needed_len);
    if (packlen > 0) {
        // 将Protobuf响应序列化到包体的位置
        res.SerializeToArray(writebuf + PROTORPC_HEAD_LENGTH, msg.head.length);
        printf("< %s\n", res.DebugString().c_str()); // 打印响应日志

        // 6. 发送响应
        channel->write(writebuf, packlen);
    }
    free(writebuf); // 释放临时缓冲区
}

这里有几个关键点:

  1. res.ByteSize():获取序列化后的大小,用于设置包头的长度字段。
  2. protorpc_pack:这个函数会根据我们定义的包头格式,将msg.head里的信息(主要是长度)编码到writebuf的前几个字节。
  3. SerializeToArray:将Response对象序列化到包体的内存位置。
  4. channel->write:这是libhv提供的异步发送接口。它不会阻塞,数据会被放入该连接的发送缓冲区,由libhv在后台通过事件循环发送出去。这是高性能的关键,handleMessage函数可以迅速返回去处理下一个请求。

整个流程就是一个经典的“接收-处理-发送”管道,清晰且高效。所有的网络I/O(channel->write)都是非阻塞的,所有的CPU密集型工作(Protobuf编解码、业务逻辑)都在当前线程快速完成,这正是异步服务器的精髓。

6. 编写客户端与进行完整测试

服务器写好了,我们还需要一个客户端来验证它。客户端的结构与服务器类似,但更简单,因为它通常不需要处理并发路由。

客户端核心步骤:

  1. 建立连接:使用hv::TcpClient连接到服务器地址。
  2. 构造请求:填充Request的id、method和params字段。params是repeated bytes,你需要先把每个参数(比如整数、字符串)用Protobuf序列化成bytes。
  3. 封包与发送:将Request序列化,加上我们自定义的包头,通过client->send发送。
  4. 接收与解包:在onMessage回调中,用和服务器一样的拆包逻辑,解析出Response。
  5. 处理响应:根据响应中的id匹配请求,判断是result还是error,进行相应处理。

一个简单的同步请求客户端可能长这样(伪代码):

hv::TcpClient client;
client.onMessage = [&](const SocketChannelPtr& channel, Buffer* buf) {
    // 拆包,反序列化Response
    // 根据res.id()找到对应的请求future,设置值
};
client.start(); // 连接并启动事件循环

// 构造请求
protorpc::Request req;
req.set_id(generate_id());
req.set_method("add");
// ... 填充params

// 封包并发送
send_packed_request(req);

// 等待响应(这里可以用条件变量或future实现简单的同步)
auto resp = wait_for_response(req.id());

在实际项目中,客户端需要管理一个请求ID到回调的映射,以实现异步调用。为了测试,我们可以直接用libhv示例中提供的protorpc_client。

完整测试流程:

  1. 启动服务器:./protorpc_server 1234
  2. 测试正常调用:./protorpc_client 127.0.0.1 1234 add 1 2。观察服务器和客户端输出,应该能看到请求和响应的日志,客户端打印出结果3。
  3. 测试错误处理:
    • 调用不存在的方法:./protorpc_client 127.0.0.1 1234 xyz 1 2。服务器应返回404 Not Found错误。
    • 触发业务逻辑错误:比如除法除零./protorpc_client ... div 1 0。你的calc_div处理器应该能捕获这个错误并返回400 Bad Request。
  4. 压力测试:你可以写一个简单的脚本,用多个进程或线程并发调用客户端,观察服务器的CPU和内存占用是否平稳。libhv的事件循环模型在处理大量空闲连接时非常节省资源。

通过这几轮测试,你就能确信这个300行代码构建的RPC服务器不仅功能完整,而且健壮可靠。

7. 性能调优与生产环境考量

虽然我们的基础版本已经可以工作,但要用于生产环境,还需要考虑更多。这里分享几个我实践中总结的优化点。

连接管理与资源释放:目前的onConnection回调只打了日志。在生产环境中,你需要在这里管理连接资源。比如,用一个std::unordered_map<int, SocketChannelPtr>来记录所有活跃连接,在连接断开时清理相关资源(如用户会话)。注意,这个映射表可能需要线程安全,因为不同连接可能在不同的I/O线程中断开。

线程池与工作队列:我们的handleMessage在当前I/O线程中执行。如果某个处理器函数非常耗时(比如涉及复杂的数据库查询或计算),它会阻塞当前事件循环,影响其他连接的处理。一个经典的优化模式是引入线程池。

// 在服务器类中增加一个线程池
#include <hv/ThreadPool.h>

class ProtoRpcServer : public hv::TcpServer {
    hv::ThreadPool worker_pool_;
public:
    ProtoRpcServer() : worker_pool_(4, 1024) { // 4个线程,任务队列容量1024
        // ... 其他初始化
        onMessage = [this](const SocketChannelPtr& channel, Buffer* buf) {
            // 将耗时的处理任务提交到线程池
            worker_pool_.commit([channel, buf_copy = *buf]() {
                // 注意:这里需要拷贝buf的数据,因为原buf可能在回调结束后被复用
                // 在这个线程中执行反序列化、路由、业务逻辑等耗时操作
                protorpc::Response res = processRequest(buf_copy);
                // 将响应发送回I/O线程(需要通过channel的loop提交任务)
                channel->loop()->queueInLoop([channel, res]() {
                    sendResponse(channel, res);
                });
            });
        };
    }
};

这样,I/O线程只负责快速的网络数据收发,耗时任务被卸载到后台工作线程,极大提升了并发能力。

日志与监控:把printf换成真正的日志库,比如spdlog或glog。记录请求的耗时、错误率、QPS等指标,这对于排查线上问题至关重要。你可以在handleMessage的开始和结束记录时间戳。

超时与心跳:网络是不稳定的。你需要为每个请求设置超时,防止因为某个慢请求耗尽资源。此外,TCP连接可能因为中间网络设备而静默断开。实现一个简单的心跳机制(比如一个空的ping/pong方法),定期检查连接健康度是很有必要的。

安全性:目前的服务器没有任何安全措施。在生产环境中,至少需要考虑:

  1. 认证:在login这样的方法中实现 token 验证,并在其他方法调用前检查 token。
  2. 限流:防止单个客户端恶意请求耗尽服务器资源。
  3. 数据校验:对反序列化后的参数进行严格的类型和范围检查。

把这些点都考虑到并实现,你的Protorpc服务器就从一个玩具变成了一个真正能扛住生产流量、易于维护的RPC服务框架。这其中的每一步,都是我在实际项目中踩过坑、交过学费才积累下来的经验。

Logo

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

更多推荐