如何通过ZMQ发送大字符串?大尺寸数据收发异常求助
ZMQ PUB/SUB模式发送大消息无响应问题解决
问题场景
使用ZMQ的PUB/SUB模式发送大消息时,当消息尺寸为32MB(1024321024),客户端recv无响应;但尺寸为**32KB(1024*32)**时收发正常。
客户端代码
#include "zmq.hpp" int main(int argc, char ** argv) { int port = 9443; zmq::context_t con(1); zmq::socket_t sock(con, ZMQ_SUB); std::string address = "tcp://address:" + std::to_string(port); // address为目标主机IP sock.setsockopt(ZMQ_SUBSCRIBE, 0, 0); sock.connect(address.c_str()); size_t size = 1024 * 32 * 1024; char* buff = new char[size]; while (1) { size_t s = sock.recv(buff, size); printf("recv size %zu\n", s); } }
服务端代码
#include "zmq.hpp" int main(int argc, char ** argv) { int port = 9443; zmq::context_t con(1); zmq::socket_t sock(con, ZMQ_PUB); std::string address = "tcp://*:" + std::to_string(port); sock.bind(address.c_str()); size_t size = 1024 * 32 *1024; // 问题点:大尺寸消息发送后客户端无响应 char* buff = new char[size]; std::string sss; std::cin >> sss; std::cout << sock.send(buff, size) << std::endl; }
解决思路与代码修改
问题根源在于ZMQ默认的缓冲区和高水位线限制无法容纳大消息,需要调整以下关键配置:
1. 调整套接字高水位线(HWM)
ZMQ的ZMQ_SNDHWM(发送高水位线)和ZMQ_RCVHWM(接收高水位线)默认限制了队列中待处理的消息数量/大小,大消息容易触发阻塞。将其设置为0表示取消限制:
服务端修改:
// 在bind之后添加 int hwm = 0; sock.setsockopt(ZMQ_SNDHWM, &hwm, sizeof(hwm));
客户端修改:
// 在connect之后添加 int hwm = 0; sock.setsockopt(ZMQ_RCVHWM, &hwm, sizeof(hwm));
2. 增大TCP缓冲区大小
系统默认的TCP发送/接收缓冲区可能不足以承载32MB的大消息,需要通过ZMQ_SNDBUF和ZMQ_RCVBUF设置足够大的缓冲区(建议设置为消息大小的1.5倍或更大):
服务端修改:
// 在bind之后添加 int snd_buf = 1024 * 32 * 1024 * 2; // 64MB缓冲区 sock.setsockopt(ZMQ_SNDBUF, &snd_buf, sizeof(snd_buf));
客户端修改:
// 在connect之后添加 int rcv_buf = 1024 * 32 * 1024 * 2; // 64MB缓冲区 sock.setsockopt(ZMQ_RCVBUF, &rcv_buf, sizeof(rcv_buf));
3. 确保消息完整发送
ZMQ的send方法在非阻塞模式下可能无法一次性发送完大消息,建议使用阻塞模式发送(默认标志为0),同时检查返回值确认发送状态:
服务端send修改:
// 阻塞模式发送,确保消息完整发出 bool send_ok = sock.send(buff, size, 0); std::cout << "Send result: " << (send_ok ? "success" : "failed") << std::endl;
4. 系统级TCP缓冲区调整(可选)
如果上述配置仍无效,需调整系统的TCP缓冲区上限:
- Linux系统:编辑
/etc/sysctl.conf,添加或修改:
执行net.core.wmem_max = 67108864 # 64MB net.core.rmem_max = 67108864sysctl -p使配置生效。
完整修改后的代码示例
服务端
#include "zmq.hpp" int main(int argc, char ** argv) { int port = 9443; zmq::context_t con(1); zmq::socket_t sock(con, ZMQ_PUB); std::string address = "tcp://*:" + std::to_string(port); sock.bind(address.c_str()); // 取消发送高水位线限制 int hwm = 0; sock.setsockopt(ZMQ_SNDHWM, &hwm, sizeof(hwm)); // 设置64MB TCP发送缓冲区 int snd_buf = 1024 * 32 * 1024 * 2; sock.setsockopt(ZMQ_SNDBUF, &snd_buf, sizeof(snd_buf)); size_t size = 1024 * 32 *1024; char* buff = new char[size]; std::string sss; std::cin >> sss; // 阻塞发送大消息 bool send_ok = sock.send(buff, size, 0); std::cout << "Send result: " << (send_ok ? "success" : "failed") << std::endl; }
客户端
#include "zmq.hpp" int main(int argc, char ** argv) { int port = 9443; zmq::context_t con(1); zmq::socket_t sock(con, ZMQ_SUB); std::string address = "tcp://address:" + std::to_string(port); // address为目标主机IP sock.setsockopt(ZMQ_SUBSCRIBE, 0, 0); sock.connect(address.c_str()); // 取消接收高水位线限制 int hwm = 0; sock.setsockopt(ZMQ_RCVHWM, &hwm, sizeof(hwm)); // 设置64MB TCP接收缓冲区 int rcv_buf = 1024 * 32 * 1024 * 2; sock.setsockopt(ZMQ_RCVBUF, &rcv_buf, sizeof(rcv_buf)); size_t size = 1024 * 32 * 1024; char* buff = new char[size]; while (1) { size_t s = sock.recv(buff, size); printf("recv size %zu\n", s); } }
内容的提问来源于stack exchange,提问作者nothingisme
相关产品推荐
相关产品推荐

