1. 项目概述为什么用C手搓Redis发布订阅最近在社区里看到不少朋友在讨论中间件和网络编程Redis的发布订阅模式作为一个经典的消息通信模型经常被提及。很多人会用Python、Go或者Java的客户端库去调用但如果你是一名C/C后台开发者或者正在深入学习系统编程仅仅会调用API可能还不够过瘾。我们得搞清楚这玩意儿到底是怎么从零到一跑起来的。所以我决定动手用C纯手工实现一个简化版的Redis发布订阅模式这不仅仅是为了“造轮子”更是为了彻底吃透网络编程、I/O多路复用、并发模型这几个后台开发的硬核知识点。这个项目能做什么简单说就是构建一个迷你消息中间件。它允许客户端连接到我们的服务器订阅Subscribe自己感兴趣的主题Channel然后当有其他客户端向这个主题发布Publish消息时所有订阅了该主题的客户端都能实时收到这条消息。这就像是一个聊天室或者系统内部的事件通知总线。通过亲手实现你会对TCP长连接管理、非阻塞I/O、事件驱动架构有肌肉记忆般的理解这些是构建高性能服务端的基石。适合谁来搞如果你已经熟悉C基础语法和面向对象对Socket编程有点概念但想深入或者被“高并发”、“Reactor模式”这些词搞得云里雾里那么这个项目就是一个绝佳的练手场。我们不依赖任何第三方网络库比如libevent、asio就从最原始的Socket API开始一步步搭起来。过程中踩的坑、做的抉择都是宝贵的经验。2. 核心架构设计与技术选型2.1 为什么选择Reactor事件驱动模型当我们面对成百上千个需要同时保持连接的客户端时传统的“一个连接一个线程”的阻塞式模型会迅速耗尽系统资源。线程的创建、销毁、上下文切换都是巨大的开销。因此高并发网络服务的首选模型是事件驱动。在事件驱动模型中Reactor模式是经典实现。它的核心思想是用一个主线程或少量线程死循环监听所有网络连接上的事件比如可读、可写当某个事件发生时再分发Dispatch给对应的处理函数去执行。这样我们用少量的线程就能管理海量的连接。在我们的Redis发布订阅服务器中Reactor模式再合适不过。主线程只需要监听新的客户端连接请求listenfd上的可读事件。已连接客户端发来的数据clientfd上的可读事件。需要向客户端发送数据的时机clientfd上的可写事件通常我们采用水平触发有数据就直接写不一定依赖写事件。我选择使用Linux系统原生的epoll作为我们的I/O多路复用器。相比于早期的select和pollepoll在管理大量文件描述符时具有显著性能优势它是基于事件回调的不会像select那样需要每次遍历整个描述符集合。注意在Windows环境下对应的机制是IOCP完成端口这是一个Proactor模型与Reactor的思维略有不同。为了聚焦核心逻辑和跨平台简化本项目以Linuxepoll为例进行讲解。如果你需要在Windows上运行可以考虑使用libevent或asio这样的跨平台库来抽象底层差异。2.2 数据结构设计如何高效管理频道与订阅关系这是发布订阅模式的核心。我们需要一个数据结构能够快速完成两种操作给定一个频道名找到所有订阅了它的客户端用于消息广播。给定一个客户端找到它订阅的所有频道用于客户端断开连接时清理订阅关系。最直观的方法是使用两个std::unordered_mapstd::unordered_mapstd::string, std::unordered_setint channel_subscribers;Key: 频道名std::stringValue: 订阅了此频道的所有客户端的文件描述符集合std::unordered_setint用途当向频道“news”发布消息时O(1)时间复杂度找到所有订阅者。std::unordered_mapint, std::unordered_setstd::string client_subscriptions;Key: 客户端的文件描述符intValue: 该客户端订阅的所有频道名集合std::unordered_setstd::string用途当客户端断开连接时O(1)时间复杂度找到其所有订阅的频道并从channel_subscribers中将其移除避免内存泄漏和无效广播。这种双向索引的设计以空间换时间确保了订阅和发布操作的高效性。当然在极端高频的场景下可能需要考虑对集合加锁或使用更高效的无锁结构但对我们这个学习项目而言这已经是一个非常清晰且有效的起点了。2.3 通信协议设计模仿Redis RESPRedis客户端与服务端通信使用一种名为RESP (REdis Serialization Protocol)的简单协议。它易于解析人类可读并且能区分数据类型。我们不完全实现RESP的所有类型但借鉴其数组格式来设计我们的命令。例如一个订阅命令可能看起来像这样以RESP数组格式*2\r\n$9\r\nSUBSCRIBE\r\n$4\r\nnews\r\n解析后等价于[SUBSCRIBE, news]一个发布命令可能像这样*3\r\n$7\r\nPUBLISH\r\n$4\r\nnews\r\n$11\r\nHello World\r\n解析后等价于[PUBLISH, news, Hello World]为什么选择文本协议而不是二进制协议对于这个项目文本协议如RESP更易于调试和实现。你可以直接用telnet或nc命令连接服务器进行测试非常直观。在实际工业级系统中根据性能要求可能会选择更紧凑的二进制协议如Protobuf、FlatBuffers。我们的服务器解析逻辑就是读取客户端发送的字节流按照\r\n分割解析出命令和参数然后执行相应的业务逻辑订阅、发布、退订等。3. 核心模块实现与代码拆解3.1 网络层Epoll事件循环的搭建这是服务器的引擎。我们创建一个Epoll类来封装相关操作。// File: epoll_wrapper.h #ifndef EPOLL_WRAPPER_H #define EPOLL_WRAPPER_H #include sys/epoll.h #include vector #include functional class Epoll { public: using EventCallback std::functionvoid(int); // fd - void Epoll(); ~Epoll(); bool addFd(int fd, uint32_t events); bool modFd(int fd, uint32_t events); bool delFd(int fd); int wait(int timeoutMs -1); // -1 表示阻塞等待 const struct epoll_event* getEvents() const { return events_.data(); } private: int epollFd_; static const int MAX_EVENTS 1024; std::vectorstruct epoll_event events_; }; #endif // EPOLL_WRAPPER_H// File: epoll_wrapper.cpp #include epoll_wrapper.h #include unistd.h #include cstring #include iostream Epoll::Epoll() { epollFd_ epoll_create1(0); if (epollFd_ -1) { perror(epoll_create1 failed); exit(EXIT_FAILURE); } events_.resize(MAX_EVENTS); } Epoll::~Epoll() { if (epollFd_ 0) { close(epollFd_); } } bool Epoll::addFd(int fd, uint32_t events) { struct epoll_event ev; memset(ev, 0, sizeof(ev)); ev.events events; ev.data.fd fd; if (epoll_ctl(epollFd_, EPOLL_CTL_ADD, fd, ev) -1) { std::cerr Failed to add fd fd to epoll std::endl; return false; } return true; } // modFd 和 delFd 实现类似...主事件循环在main.cpp或Server类中// 伪代码展示主循环逻辑 Epoll epoller; // 1. 创建监听socket (listenFd)设置为非阻塞并添加到epoll监听读事件 int listenFd create_and_bind_server_socket(port); set_nonblocking(listenFd); epoller.addFd(listenFd, EPOLLIN | EPOLLET); // 边缘触发(ET)模式效率更高 while (!stop) { int numEvents epoller.wait(100); // 等待100毫秒 for (int i 0; i numEvents; i) { int fd epoller.getEvents()[i].data.fd; uint32_t events epoller.getEvents()[i].events; if (fd listenFd) { // 处理新连接 handle_new_connection(listenFd, epoller); } else if (events EPOLLIN) { // 客户端有数据可读 handle_client_message(fd, epoller); } else if (events EPOLLERR || events EPOLLHUP) { // 客户端出错或挂断 handle_client_disconnect(fd, epoller); } // 注意写事件(EPOLLOUT)的处理通常在我们有数据要发送时再监听避免busy-loop } }3.2 协议解析器拆解RESP格式命令我们需要一个ProtocolParser类来将收到的字节流解析成结构化的命令。// File: protocol_parser.h #ifndef PROTOCOL_PARSER_H #define PROTOCOL_PARSER_H #include string #include vector enum class CommandType { UNKNOWN, SUBSCRIBE, UNSUBSCRIBE, PUBLISH, PING, QUIT }; struct ParsedCommand { CommandType type; std::vectorstd::string args; // 命令参数如对于PUBLISHargs[0]是频道args[1]是消息 }; class ProtocolParser { public: // 状态机可能一次recv的数据不足以构成完整命令需要缓冲 bool feed(const char* data, size_t len); bool has_complete_command() const; ParsedCommand parse_next_command(); void reset(); private: std::string buffer_; // 更复杂的实现可能需要处理嵌套数组这里简化处理简单RESP数组 bool parse_simple_resp_array(const std::string str, size_t pos, ParsedCommand cmd); }; #endif // PROTOCOL_PARSER_H解析的核心在于处理*数组、$批量字符串等前缀并按照\r\n进行分割。这是一个典型的状态机解析过程。对于学习而言实现一个能解析简单RESP数组的版本就足够了。3.3 业务逻辑核心订阅与发布的管理这是我们的PubSubServer类它持有之前提到的两个unordered_map并处理具体的命令。// File: pubsub_server.h #ifndef PUBSUB_SERVER_H #define PUBSUB_SERVER_H #include unordered_map #include unordered_set #include string #include mutex // 考虑线程安全时需要 class PubSubServer { public: // 处理来自客户端的命令 std::string process_command(int clientFd, const ParsedCommand cmd); // 客户端断开连接时清理 void client_disconnected(int clientFd); private: // 关键数据结构 std::unordered_mapstd::string, std::unordered_setint channel_subscribers_; std::unordered_mapint, std::unordered_setstd::string client_subscriptions_; // 注意在多线程环境下操作这两个map需要加锁如std::shared_mutex // 本项目是单Reactor线程暂时不需要。 // 内部处理方法 std::string handle_subscribe(int clientFd, const std::vectorstd::string args); std::string handle_unsubscribe(int clientFd, const std::vectorstd::string args); std::string handle_publish(int clientFd, const std::vectorstd::string args); }; #endif // PUBSUB_SERVER_Hhandle_publish函数的实现是精华std::string PubSubServer::handle_publish(int publisherFd, const std::vectorstd::string args) { if (args.size() 2) return -ERR wrong number of arguments for publish command\r\n; const std::string channel args[0]; const std::string message args[1]; auto it channel_subscribers_.find(channel); if (it channel_subscribers_.end()) { // 没有订阅者根据Redis协议返回接收到消息的客户端数量0 return :0\r\n; } // 构建要广播的消息格式。Redis的发布订阅消息是推送式的格式固定。 // 例如message channel actual_message std::string broadcast_msg *3\r\n$7\r\nmessage\r\n$ std::to_string(channel.length()) \r\n channel \r\n$ std::to_string(message.length()) \r\n message \r\n; // 遍历订阅者集合发送消息 for (int subFd : it-second) { if (subFd ! publisherFd) { // 通常不发给发布者自己除非他也订阅了 // 这里需要将消息写入每个客户端的发送缓冲区 // 在实际代码中我们可能有一个 ClientSession 对象管理每个连接的缓冲区 // 例如get_client_session(subFd).send_buffer.append(broadcast_msg); // 然后修改epoll监听该fd的写事件(EPOLLOUT)在下次事件循环中写出 send_to_client(subFd, broadcast_msg); } } // 返回接收到消息的客户端数量 return : std::to_string(it-second.size()) \r\n; }实操心得这里有一个重要的设计点。send_to_client不能直接调用write或send因为TCP socket的发送缓冲区可能已满特别是网络慢或客户端处理慢时直接调用会阻塞线程破坏我们的事件循环。正确的做法是将待发送数据追加到该客户端对应的应用层发送缓冲区然后通过epoll监听该socket的可写事件EPOLLOUT。当可写事件触发时再从缓冲区尝试发送数据。如果一次没发完继续等待下次可写事件。这就是所谓的“写缓冲区”管理是异步非阻塞编程的必备技巧。3.4 客户端会话管理我们需要一个ClientSession类来封装每个连接的状态。// File: client_session.h #ifndef CLIENT_SESSION_H #define CLIENT_SESSION_H #include string #include protocol_parser.h class ClientSession { public: ClientSession(int fd) : fd_(fd), parser_() {} int fd() const { return fd_; } // 接收数据并喂给解析器 void append_recv_data(const char* data, size_t len) { recv_buffer_.append(data, len); parser_.feed(data, len); } // 从接收缓冲区中提取并解析完整命令 bool has_complete_command() { return parser_.has_complete_command(); } ParsedCommand get_next_command() { return parser_.parse_next_command(); } // 发送缓冲区管理 void append_send_data(const std::string data) { send_buffer_ data; } bool has_pending_data() const { return !send_buffer_.empty(); } // 尝试发送数据返回是否全部发送完毕 bool try_send(); const std::string get_send_buffer() const { return send_buffer_; } void clear_sent_data(size_t len) { if (len send_buffer_.size()) { send_buffer_.clear(); } else { send_buffer_.erase(0, len); } } private: int fd_; std::string recv_buffer_; // 应用层接收缓冲区 std::string send_buffer_; // 应用层发送缓冲区 ProtocolParser parser_; }; // try_send 的实现需要考虑非阻塞write和EAGAIN/EINTR错误 bool ClientSession::try_send() { if (send_buffer_.empty()) return true; ssize_t n ::send(fd_, send_buffer_.data(), send_buffer_.size(), MSG_NOSIGNAL); if (n 0) { clear_sent_data(n); return send_buffer_.empty(); // 如果清空了返回true } else if (n -1) { if (errno EAGAIN || errno EWOULDBLOCK) { // 发送缓冲区满了等下次EPOLLOUT事件 return false; } else { // 真正的错误应关闭连接 perror(send error); return true; // 标记为需要关闭 } } // n 0 通常意味着连接关闭 return true; // 标记为需要关闭 }这样在主事件循环中handle_client_message就变成了从socketrecv数据到ClientSession的接收缓冲区。调用session.append_recv_data。循环检查session.has_complete_command()有则取出命令交给PubSubServer::process_command处理。将服务器返回的响应如“:3\r\n”或广播消息通过session.append_send_data加入发送缓冲区。如果发送缓冲区非空则通过epoll修改该fd监听EPOLLOUT事件。在EPOLLOUT事件触发时调用session.try_send()。4. 编译、运行与基础测试4.1 环境准备与编译你需要一个Linux环境或WSL2和GCC/Clang编译器。确保安装了必要的开发工具。项目目录结构建议如下cpp_redis_pubsub/ ├── CMakeLists.txt ├── src/ │ ├── main.cpp │ ├── epoll_wrapper.h / .cpp │ ├── protocol_parser.h / .cpp │ ├── pubsub_server.h / .cpp │ └── client_session.h / .cpp └── build/一个简单的CMakeLists.txtcmake_minimum_required(VERSION 3.10) project(cpp_redis_pubsub) set(CMAKE_CXX_STANDARD 11) set(CMAKE_CXX_STANDARD_REQUIRED ON) add_executable(pubsub_server src/main.cpp src/epoll_wrapper.cpp src/protocol_parser.cpp src/pubsub_server.cpp src/client_session.cpp ) target_include_directories(pubsub_server PRIVATE src)编译命令mkdir build cd build cmake .. make4.2 启动服务器与基础功能测试编译后在build目录下会生成pubsub_server可执行文件。你可以指定端口运行./pubsub_server 6379现在我们可以用telnet或netcat(nc)来模拟客户端进行测试。打开两个终端。终端1订阅者:nc localhost 6379 SUBSCRIBE news你应该会收到服务器返回的订阅确认消息模仿Redis协议。终端2发布者:nc localhost 6379 PUBLISH news Hello from publisher!此时在终端1应该会立即收到一条推送消息格式类似于*3 $7 message $4 news $23 Hello from publisher!注意实际传输是\r\n分隔的终端显示可能换行这就完成了最基本的发布订阅流程。你可以测试多个订阅者、多个频道、退订UNSUBSCRIBE等功能。4.3 使用Redis-cli进行兼容性测试进阶因为我们模仿了RESP协议理论上可以与官方的redis-cli进行简单交互。但请注意我们的实现是极简版只支持SUBSCRIBEPUBLISH等少数命令redis-cli的很多交互模式如订阅模式下的特殊提示我们可能不支持。可以尝试redis-cli -p 6379 127.0.0.1:6379 PUBLISH mychannel test (integer) 1 # 如果有一个订阅者会返回1这可以验证我们协议格式的基本正确性。5. 性能优化、问题排查与扩展思考5.1 性能瓶颈分析与优化方向一个单Reactor线程的服务器其性能瓶颈通常在于CPU协议解析字符串分割、类型转换、数据结构查找哈希表、内存拷贝缓冲区操作。I/Oepoll_wait返回的大量事件处理、send/recv系统调用。优化思路缓冲区设计使用连续内存如std::vectorchar或链表式缓冲区减少拷贝。可以考虑使用readv/writev进行分散/聚集I/O。协议解析优化避免在解析过程中频繁创建std::string可以尝试零拷贝解析直接操作接收缓冲区的原始数据。数据结构当频道数巨大百万级时std::unordered_map的哈希冲突和内存开销可能成为问题。可以考虑使用性能更好的哈希表如absl::flat_hash_map或分片哈希。Reactor线程模型单Reactor处理所有I/O和业务在业务逻辑复杂时可能阻塞事件循环。可以升级为“单Reactor 线程池”模型主Reactor线程只负责I/O事件accept read write将解析好的命令对象投递到一个工作线程池中执行具体的PubSubServer::process_command逻辑然后将结果返回给主线程由主线程负责写回socket。这能有效利用多核防止业务逻辑阻塞网络I/O。5.2 常见问题与调试技巧客户端连接不上检查服务器是否成功绑定端口netstat -tlnp | grep 端口号。检查防火墙是否阻止了连接sudo ufw status。检查服务器代码中的bind或listen调用是否失败检查errno并打印。订阅后收不到发布的消息检查发布命令的参数格式是否正确频道名是否完全匹配大小写敏感调试在handle_publish函数中打印channel_subscribers_[channel]集合的大小和内容。检查广播消息的构建格式是否正确可以用printf或写入日志文件查看生成的RESP字符串。检查send_to_client逻辑是否正确是否因为发送缓冲区满而将消息丢弃检查try_send函数的返回值和对EAGAIN的处理。内存泄漏关键点确保在handle_client_disconnect中不仅要从epoll中删除fd关闭socket还要从client_subscriptions_和channel_subscribers_中彻底清理该客户端的订阅信息。工具使用valgrind --leak-checkfull ./pubsub_server进行内存检查。服务器在高并发下崩溃或卡死检查是否正确处理了EAGAIN/EWOULDBLOCK非阻塞socket的read返回-1且errno为EAGAIN时表示数据还没准备好这不是错误应该继续循环。检查epoll是否使用了边缘触发ET模式如果用了ET必须循环read直到返回EAGAIN否则会丢失事件。检查是否有死循环或阻塞操作在主事件循环中比如在process_command中执行了耗时的数据库查询。5.3 项目扩展方向这个基础版本可以作为一个起点向多个方向深化支持认证AUTH在ClientSession中增加一个authenticated状态位。在连接建立后要求客户端首先发送AUTH password命令验证通过后才能执行其他命令。支持模式订阅PSUBSCRIBE允许使用通配符如news.*订阅多个频道。这需要将数据结构从精确匹配的哈希表升级为支持模式匹配的结构如前缀树Trie或将模式编译成正则表达式。持久化与持久订阅将频道和订阅关系持久化到磁盘如RDB快照或AOF日志服务器重启后能恢复。还可以实现持久订阅即使客户端断开重连也能收到离线期间的消息需要消息队列。集群化单个服务器有性能上限。可以设计一个代理层根据频道名的哈希值将订阅和发布请求路由到不同的后端服务器节点上。这涉及到节点间通信、数据同步等更复杂的分布式系统问题。完善Redis协议实现更多的Redis命令如PING,QUIT,UNSUBSCRIBE甚至尝试实现一些数据结构的命令让它更像一个真正的迷你Redis。通过这个项目你亲手触摸到了高性能网络服务的核心脉络。从epoll的边沿触发到应用层缓冲区的管理从协议设计到并发数据结构的选择每一个细节都影响着服务的稳定与高效。下次当你再用redis-cli发布一条消息时你看到的将不再是一个黑盒而是一个由事件循环、哈希表和TCP流组成的、清晰可见的精密机器。这才是动手实现的意义所在。