尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

用 ZeroMQ 承载 Thrift RPC:TZmq 传输层实现原理与实战指南

用 ZeroMQ 承载 Thrift RPC:TZmq 传输层实现原理与实战指南 用 ZeroMQ 承载 Thrift RPCTZmq 传输层实现原理与实战指南【免费下载链接】thriftApache Thrift项目地址: https://gitcode.com/GitHub_Trending/thr/thrift本指南围绕 Apache Thrift 仓库中 contrib/zeromq 目录下的 ZeroMQ 适配代码展开讲解如何将 Thrift RPC 跑在消息队列 ZeroMQ 之上。文中包含完整可运行的客户端/服务端示例Python 与 C、两种消息模式的 socket 选择规则以及普通方法与 oneway 方法共存的 TZmqMultiServer 解决方案帮助读者理解 Thrift 流式传输与 ZeroMQ 消息式传输之间的阻抗失配并掌握一套可直接落地的混合传输方案。背景Thrift 面向流ZeroMQ 面向消息Thrift 的传输层TTransport天然是为 TCP 这类**流式stream-based接口设计的数据按字节流连续读写连接保持长活消息边界由协议层如二进制协议自行划分。而 ZeroMQ 是消息式message-based**通信库一次send/recv就是一个完整消息消息边界由内核维护。两者之间的阻抗失配大部分被适配层隐藏了——开发者依然用标准的TTransport接口读写数据底层则由TZmqClient/TZmqServer负责把字节流包装成一条条 ZeroMQ 消息。但有一个问题无法被隐藏必须在应用层面处理oneway 方法与普通方法的应答语义不同见 README.md。普通方法需要返回结果属于请求-应答模式oneway 方法oneway void发出即忘不需要应答。ZeroMQ 要求在创建 socket 时就声明消息模式REQ/REP、PUSH/PULL、PUB/SUB 等不可能在同一连接上逐条消息地决定是否回包。因此这套实现把责任交给客户端普通方法必须走ZMQ_REQsocketoneway 方法必须走ZMQ_DOWNSTREAM实际代码中使用的是ZMQ_PUSH语义等同下游发送socket。相应地同时暴露两类方法的服务必须监听两个端口、部署两个服务端而TZmqMultiServer负责把两个服务端放进同一个线程统一调度。实验接口storage.thrift整套代码围绕一个极简的存储服务演示接口展开storage.thriftservice Storage { oneway void incr(1: i32 amount); i32 get(); }接口刻意同时包含两种方法类型就是为了演示前面提到的模式分裂问题incr是oneway方法客户端只发请求、不等应答get是普通方法客户端需要拿到当前存储值。两者由同一个 handler 提供状态C 与 Python 的实现保持一致内部维护一个int32_t value_incr累加、get返回当前值见 test-server.cpp 与 test-server.py。构建与运行环境要求安装 Thrift 与 ZeroMQPython 侧还需要 pyzmq文档记录的测试环境为ZeroMQ 2.0.7与pyzmq afabbb5b9bd3README.md。代码基于较老的 zmq API如zmq.hpp、ctx.socket(...)、PULL/PUSH 消息模式在更新版本上编译运行时可能需要进行 API 适配这是阅读与复用时需要注意的前提。编译在contrib/zeromq目录下直接执行makemake会依次完成三件事见 Makefile用thrift --gen cpp从storage.thrift生成 C 代码Storage.h、Storage.cpp、storage_types.h、storage_types.cpp用thrift --gen py生成 Python 代码storage/包编译四个可执行程序test-client、test-server、test-sender、test-receiver链接-lzmq -lthrift。如果 Thrift 或 ZeroMQ 安装在非标准路径需要在 make 命令行显式指定两个变量README.mdmake THRIFT/path/to/thrift PKG_CONFIG_PATH/path/to/pkgconfigTHRIFTThrift 代码生成器thrift 编译器的位置PKG_CONFIG_PATH需要包含 Thrift 与 ZeroMQ 两者的 pkgconfig 文件所在目录。运行方式测试服务端不需要任何参数测试客户端则有两种模式README.md不带参数走 REQ 模式连接tcp://127.0.0.1:9090调用get()并打印当前存储值带一个整数参数走 PUSH 模式连接tcp://127.0.0.1:9091调用incr(amount)累加该数值。先启动服务端再依次调用客户端即可验证# 终端 1启动服务端 ./test-server # 终端 2读取初始值输出 0 ./test-client # 终端 3累加 5 ./test-client 5 # 终端 2再次读取输出 5 ./test-clientPython 侧运行方式相同test-client.py、test-server.py可直接用 Python 解释器执行。注意 oneway 调用发出后立即返回客户端在incr后time.sleep(0.05)/usleep(50000)短暂等待让消息送达服务端再查询结果test-client.py、test-client.cpp。客户端实现TZmqClient 传输层Python 实现TZmqClient.py 定义了TZmqClient它同时继承TTransportBase与CReadableTransport因此可以配合TBinaryProtocolAcceleratedCython 加速版二进制协议使用class TZmqClient(TTransportBase, CReadableTransport): def __init__(self, ctx, endpoint, sock_type): self._sock ctx.socket(sock_type) self._endpoint endpoint self._wbuf StringIO() self._rbuf StringIO() def open(self): self._sock.connect(self._endpoint) ...工作流程与标准 Thrift 传输层一致TZmqClient.pywrite/flush协议层把请求序列化写入_wbufflush()时把整段字节取出作为一条ZeroMQ 消息send出去read先读_rbuf缓冲区空了就recv()一条应答消息放入_rbuf再继续读从而把消息抽象成流对上层透明。构造函数接收三个参数ctxzmq.Context 实例、endpoint如tcp://127.0.0.1:9090、sock_type如zmq.REQ/zmq.PUSH。C 实现C 侧 TZmqClient.h 把TZmqClient定义为TTransport的子类内部持有zmq::socket_t、两个TMemoryBuffer写缓冲与读缓冲以及一个zmq::message_tclass TZmqClient : public TTransport { public: TZmqClient(zmq::context_t ctx, const std::string endpoint, int type); void open() { if (zmq_type_ ZMQ_PUB) { sock_.bind(endpoint_.c_str()); // PUB 模式下客户端做 bind } else { sock_.connect(endpoint_.c_str()); // 其余模式做 connect } } ... };一个值得注意的细节C 客户端对ZMQ_PUB类型做了特殊处理此时客户端自身执行bind而不是connect这为发布/订阅模式的扩展预留了接口。read_virt/write_virt/writeEnd的具体实现位于 TZmqClient.cpp整体与 Python 版本对称写缓冲攒满后整体作为一条消息发出读到空缓冲时从 socket 接收新消息再填充读缓冲。使用示例Python 客户端装配方式test-client.pyctx zmq.Context() transport TZmqClient.TZmqClient(ctx, endpoint, socktype) protocol thrift.protocol.TBinaryProtocol.TBinaryProtocolAccelerated(transport) client storage.Storage.Client(protocol) transport.open()C 客户端装配方式test-client.cppzmq::context_t ctx(1); shared_ptrTZmqClient transport(new TZmqClient(ctx, endpoint, socktype)); shared_ptrTBinaryProtocol protocol(new TBinaryProtocol(transport)); StorageClient client(protocol); transport-open();两种语言都把TZmqClient当作普通 TTransport 注入协议与生成的 client上层完全感知不到底层是 ZeroMQ。服务端实现TZmqServer 与 TZmqMultiServer单 socket 服务端TZmqServer.py 中TZmqServer继承自标准TServerclass TZmqServer(thrift.server.TServer.TServer): def __init__(self, processor, ctx, endpoint, sock_type): thrift.server.TServer.TServer.__init__(self, processor, None) self.zmq_type sock_type self.socket ctx.socket(sock_type) self.socket.bind(endpoint) def serveOne(self): msg self.socket.recv() itrans thrift.transport.TTransport.TMemoryBuffer(msg) otrans thrift.transport.TTransport.TMemoryBuffer() iprot self.inputProtocolFactory.getProtocol(itrans) oprot self.outputProtocolFactory.getProtocol(otrans) try: self.processor.process(iprot, oprot) except Exception: logging.exception(Exception while processing request) if self.zmq_type zmq.REP: msg otrans.getvalue() self.socket.send(msg)处理模型TZmqServer.py非常直观从 socketrecv()一条请求消息包进TMemoryBuffer作为输入传输层创建输出TMemoryBuffer交给processor.process(iprot, oprot)执行实际业务逻辑仅当 socket 类型为zmq.REP时才把输出缓冲区内容作为应答send回去——这正是回复与否由 socket 类型决定的体现PULL类型的 oneway socket 处理完即结束不回任何消息处理过程中出现异常会记入日志并继续走完流程保证请求有始有终。C 版本 TZmqServer.cpp 的逻辑与之一一对应serveOne接收消息构造输入输出TMemoryBuffer调用processor_-process(...)ZMQ_REP时把输出缓冲拷贝进zmq::message_t回发。双端口混合服务TZmqMultiServer前面提到同时含普通方法与 oneway 方法的服务必须开两个端口。TZmqMultiServer就是为在同一个线程里把两个服务端跑在一起而生的README.md。Python 实现基于zmq.PollerTZmqServer.pyclass TZmqMultiServer(object): def __init__(self): self.servers [] def serveOne(self, timeout-1): self._serveActive(self._setupPoll(), timeout) def serveForever(self): poll_info self._setupPoll() while True: self._serveActive(poll_info, -1) def _setupPoll(self): server_map {} poller zmq.Poller() for server in self.servers: server_map[server.socket] server poller.register(server.socket, zmq.POLLIN) return (server_map, poller)原理是把每个子TZmqServer的 socket 注册进同一个 pollerserveForever循环里poll()出就绪的 socket再分发给对应的serveOne()。C 版本 TZmqServer.cpp 使用zmq::poll与zmq::pollitem_t数组实现同样的多路复用并在轮询后逐个检查revents ZMQ_POLLIN决定是否派发。组装两个服务端的完整服务端代码Python 版test-server.pyhandler StorageHandler() processor storage.Storage.Processor(handler) ctx zmq.Context() reqrep_server TZmqServer.TZmqServer(processor, ctx, tcp://0.0.0.0:9090, zmq.REP) oneway_server TZmqServer.TZmqServer(processor, ctx, tcp://0.0.0.0:9091, zmq.PULL) multiserver TZmqServer.TZmqMultiServer() multiserver.servers.append(reqrep_server) multiserver.servers.append(oneway_server) multiserver.serveForever()C 版本结构相同test-server.cpp仅 API 风格不同servers().push_back(...)。消息模式选择规则客户端责任模型综合客户端与服务端代码这套适配层的消息模式映射关系可以总结如下方法类型服务端 socket客户端 socket端点示例普通方法如getZMQ_REPzmq.REPZMQ_REQzmq.REQtcp://127.0.0.1:9090oneway 方法如incrZMQ_PULLzmq.PULLZMQ_PUSHzmq.PUSHtcp://127.0.0.1:9091客户端的模式切换逻辑见 test-client.py 与 test-client.cpp不带参数时默认 REQ/9090 读值带整数参数时切换为 PUSH/9091 发增量。README 中提到的ZMQ_DOWNSTREAM是早期 ZeroMQ 对下游push/pull 模式socket 的称呼当前代码实际落地的类型是ZMQ_PUSH/ZMQ_PULL。这套客户端责任模型是设计上的权衡既然 ZeroMQ 不允许在消息级别动态决定是否回包就把该方法是否需要应答这一语义在客户端编码进 socket 类型里服务端则纯粹按 socket 类型决定是否回发REP 回、PULL 不回两侧各司其职。目录内其他文件与扩展方向test-sender.cpp 与 test-receiver.cpp一对独立的发送/接收程序用于验证纯 PUSH/PULL 数据通路可作为 oneway 场景的最小验证工具csharpC# 语言版本示例。需要说明的是C# 已不再是 Apache Thrift 的官方支持目标官方推荐以 netstd.NET Standard 实现见 lib/netstd替代README.md 明确这段代码仅保留供学习参考除非有人将其移植到 netstd否则不再维护TZmqClient.cpp / TZmqServer.hC 侧传输层与 server 的完整声明与实现其中TZmqServer支持serveOne/serveForever两种驱动方式。已知局限与适用边界README 作者明确自评这套代码还称不上生产级README.md主要局限有二未接入 Thrift 的完整钩子机制标准 Thrift server 的事件钩子、部分生命周期管理能力在 TZmq 适配层中并未完整实现存在不必要的内存拷贝请求/应答在 ZeroMQ 消息与TMemoryBuffer之间反复搬运性能并非最优。此外正如消息模式选择规则一节所述应答与否依赖客户端自觉选择正确的 socket 类型若客户端用错模式例如对普通方法使用 PUSH服务端将不会回包调用方会因等不到应答而阻塞。因此该方案适合教学演示、原型验证或对消息边界有强诉求的内部场景追求高吞吐的生产环境建议评估 Thrift 官方的 TCP/HTTP 传输或对这套代码做针对性的零拷贝与钩子接入改造后再使用。【免费下载链接】thriftApache Thrift项目地址: https://gitcode.com/GitHub_Trending/thr/thrift创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表