ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

【C++三方组件】ZeroMQ:消息队列的另一种形态

【C++三方组件】ZeroMQ:消息队列的另一种形态 【C三方组件】ZeroMQ消息队列的另一种形态【摘要】ZeroMQ 是可以嵌入进程的消息通信库在 socket 之上提供消息边界、收发队列和多种通信模式。本文说明它与 broker 型消息系统的区别再通过请求应答、发布订阅、任务分发和多帧消息介绍常用 API。【版本基准】libzmq 4.3.5MPL-2.0cppzmq 4.11.0MITC17。完整代码见 zeromq_demo.cpp构建见 配套示例。1. WhatZeroMQ 是什么ZeroMQ 是一个消息通信库。程序创建 context 和特定类型的 socket通过bind、connect建立通信关系再以消息为单位收发数据。它可用于同进程线程之间也可通过 TCP 用于跨进程、跨机器通信。使用 ZeroMQ 不要求先部署一个独立 broker但应用也可以用它构建代理或中间节点。“不要求 broker”不等于所有架构都没有中间节点更不等于无需监控连接、队列和进程状态。它与 Kafka、RabbitMQ 等系统的直接区别是交付形态与责任范围ZeroMQ 首先提供通信机制内存消息队列本身不提供持久日志、消费确认或重启后的自动回放。可靠业务协议可以建立在其上但需要额外设计。socket 组合通信模式应用需要理解的约束REQ / REP请求与应答默认严格交替发送和接收PUB / SUB一对多发布订阅按前缀订阅允许丢消息PUSH / PULL任务分发与接收在可用下游间分配存在背压DEALER / ROUTER更灵活的异步请求和路由应用承担更多关联与协议管理PAIR / PAIR一对一通信适合专用连接不是通用分发器2. Why为什么不直接使用 TCP socketTCP 提供可靠有序的字节流但应用往往想发送的是“一条任务”“一次查询”或“一次状态更新”。两者之间还有一层工作需求自行实现的工作ZeroMQ 提供的能力保留消息边界编码长度、处理半包与合并读取消息与多帧接口发送期间对端暂时不可用管理连接状态和待发数据连接管理与队列行为受选项和模式影响一条数据发给多个订阅者维护订阅关系、逐个发送PUB/SUB任务分给多个 worker管理多个下游及分发状态PUSH/PULL本机线程与远程进程共用接口维护不同传输实现inproc://、tcp://等传输这减少了通信基础设施的代码量。应用仍要定义消息内容、版本兼容、处理成功的判定和故障恢复策略。send成功通常只表明消息按当前模式被接受不是“对端业务已经执行成功”的证明。若核心需求是持久化、消费确认、重放或跨重启恢复应优先比较现成消息中间件的能力。若需求是进程间任务通信、可丢失状态广播或自定义分布式组件ZeroMQ 更接近所需的构件。3. How接入与基本对象vcpkg 可以安装cppzmq和zeromq。cppzmq 是 C 绑定提供 RAII 的 context、socket 和 messagefind_package(cppzmq CONFIG REQUIRED) add_executable(app zeromq_demo.cpp) target_compile_features(app PRIVATE cxx_std_17) target_link_libraries(app PRIVATE cppzmq)#includezmq.hppzmq::context_t context;zmq::socket_tsocket(context,zmq::socket_type::req);socket.set(zmq::sockopt::linger,1000);socket.set(zmq::sockopt::sndtimeo,3000);socket.set(zmq::sockopt::rcvtimeo,3000);context 必须活得比使用它的 socket 更久。本文使用的普通 REQ、REP、PUB、SUB、PUSH、PULL socket 都不用于多线程并发共享worker 在线程内部创建、使用和销毁自己的 socket。配套示例检查send、recv返回值。超时会返回空结果其他错误可能抛出zmq::error_t不能忽略后继续假装消息已经处理。3.1 REQ/REP两个进程完成三次往返服务端创建 REP 并绑定端点客户端创建 REQ 并连接。共同循环的核心如下send、receive是源码中带超时检查的辅助函数for(inti1;i3;i){if(server){send(socket,reply: receive(socket));}else{send(socket,hello std::to_string(i));std::coutreceive(socket)\n;}}先在终端 A 启动服务端再在终端 B 启动客户端# 终端 A./build/zeromq_demo rep tcp://127.0.0.1:18082# 终端 B./build/zeromq_demo req tcp://127.0.0.1:18082服务端等候请求的超时设为 30 秒便于手动启动如果超时退出重新启动即可。客户端输出reply: hello 1 reply: hello 2 reply: hello 3REQ 默认要求发送后接收REP 默认要求接收后应答。它简化了一次只处理一个未完成请求的协议。若服务端执行失败需要通过应用消息表达结果或者定义超时后的重试与去重而不是把连接恢复当成业务已经恢复。配套的state模式故意连续发送两次用符号EFSM验证状态机约束。错误数字不应按操作系统猜测或硬编码。需要多个未完成请求时可以进一步学习 DEALER/ROUTER而不是随意打破 REQ 的顺序。3.2 发布订阅前缀匹配与订阅就绪SUB 设置感兴趣的前缀zmq::socket_tsubscriber(context,zmq::socket_type::sub);subscriber.set(zmq::sockopt::subscribe,news);subscriber.connect(tcp://127.0.0.1:18082);订阅按字节前缀匹配。因此news也会匹配以该前缀开头的其他主题需要严格主题划分时应在编码中加入明确分隔或其他约定。连接及订阅建立需要时间。为了让示例不依赖固定 sleep发布端使用XPUB它可以接收订阅通知等到收到news的订阅后再发消息。这里演示的是一个订阅者多个订阅者需要自己的就绪判定。conststd::string subscriptionreceive(publisher);if(subscription.empty()||subscription[0]!\1||subscription.substr(1)!news)throwstd::runtime_error(expected a news subscription);终端 A 启动发布端30 秒内在终端 B 启动订阅端# 终端 A./build/zeromq_demo pub tcp://127.0.0.1:18082# 终端 B./build/zeromq_demo sub tcp://127.0.0.1:18082发布端同时发送 news 和 weather订阅端只收到news1到news5。等待订阅就绪解决的是示例的启动协调不会把 PUB/SUB 变成带确认的可靠投递。3.3 多帧消息主题与正文分别发送上面的发布端把消息拆成两个帧send(publisher,news,zmq::send_flags::sndmore);send(publisher,std::to_string(i));接收端检查是否还有后续帧再取正文constautotopicreceive(subscriber);if(!subscriber.get(zmq::sockopt::rcvmore))throwstd::runtime_error(missing payload frame);constautopayloadreceive(subscriber);if(subscriber.get(zmq::sockopt::rcvmore))throwstd::runtime_error(unexpected extra frame);多帧消息适合“主题正文”“路由信封负载”等结构。对于此类 API一次recv取得的是一帧应用通过rcvmore识别整个多帧消息的边界。正文可以是 JSON、protobuf 或自定义二进制但要由通信双方约定。3.4 PUSH/PULL任务分发给多个 worker配套pipeline模式在同一个 context 中创建 PUSH以及两个各在线程内部创建的 PULL通过inproc://tasks通信。主线程等两位 worker 完成初始化再发送十条任务for(autosignal:ready)signal.get_future().get();for(inti0;i10;i)send(push,std::to_string(i));for(autoworker:workers)worker.join();worker 持续接收并计数演示中用 1 秒空闲超时结束。完整源码包含线程异常传递和 join正式任务系统通常应设计显式结束信号。./build/zeromq_demo pipeline一次验证输出为tasks10 workers5,5。分发依据是可用下游和队列状态不是库知道哪个 worker 的业务计算“最闲”不同启动时序与负载下不应把每个 worker 恰好处理五条作为协议保证。PUSH/PULL 的背压也不等于任务确认。worker 收到任务后崩溃应用仍需检测缺失结果、决定是否重新分配并考虑重复执行的影响。4. 队列、超时与退出ZMQ_SNDHWM、ZMQ_RCVHWM用于约束排队通常按对端队列和消息数量理解不能直接把值当成整个进程的内存字节上限。消息尺寸和连接数量同样影响内存。不同模式达到限制后的行为不同PUB 对慢订阅者可能丢弃消息PUSH 在没有可用下游时可能等待或者在非阻塞超时设置下返回失败。因此“加大 HWM”只是增加缓冲不会自动提高可靠性。linger决定关闭时如何处理待发送消息。0 表示丢弃待发消息正值表示最多等待相应毫秒数默认 -1 可能无限等待。本文统一给了有限等待时间退出仍须保证各线程关闭 socket再销毁 context。inproc的通信双方必须使用同一个 context。从 ZeroMQ 4.0 起inproc 已经允许先 connect 后 bind不应再把旧版本的顺序限制写成当前要求。5. 选型与参考主要需求需要重点比较进程间通信、状态发布、任务分发ZeroMQ 的模式与失败行为持久化、重放、消费确认、运营工具broker 型消息系统的完整能力已定义服务接口与跨语言方法调用gRPC 等 RPC 框架完全自定义传输与应用协议Asio 等底层网络库及自行维护成本是否允许丢消息只是其中一项。还要考虑业务确认、重试、部署方式和团队能承担的协议维护工作。ZeroMQ Guide基本模式。可靠性与请求应答模式。inproc 传输说明。cppzmqC 接口、选项与示例。
返回列表