VeighNa RpcService 模块深度指南:基于 ZeroMQ 的多进程分布式交易路由
【免费下载链接】vnpy基于Python的开源量化交易平台开发框架项目地址: https://gitcode.com/gh_mirrors/vn/vnpy
RpcService 是 VeighNa Trader 中用于将单个交易进程转化为 RPC 服务器的功能模块,对外提供交易路由、行情数据推送、持仓资金查询等服务。本文以 rpc_service.md 为骨架,结合仓库中 vnpy/rpc 的底层实现与 examples/client_server 的完整示例,讲解从服务端配置、客户端接入到源码级通讯机制的全过程,读完即可搭建一套"一条通道、多客户端并行交易"的分布式交易架构。
功能简介
RpcService 的核心作用是把一个 VeighNa Trader 进程升级为RPC 服务器(服务端),对外提供两类能力:
- 交易路由:客户端进程通过 RPC 远程调用服务端的下单、撤单等交易接口;
- 数据分发:服务端将收到的行情、委托、成交等事件数据主动推送给所有已连接的客户端。
服务端集中连接真实的交易接口(如 CTP),客户端无需再次配置账户密码,只需与服务端通讯,即可获得与本地直连几乎一致的使用体验。关于 RPC 更宏观的应用场景,见本文末尾【RPC 的应用场景】一节。
加载启动
通过 VeighNa Station 加载
启动并登录 VeighNa Station 后,点击【交易】按钮进入 VeighNa Trader,在配置对话框的【应用模块】栏勾选【RpcService】即可。
通过脚本加载
在启动脚本中按以下方式添加应用模块:
# 写在顶部 from vnpy_rpcservice import RpcServiceApp # 写在创建main_engine对象后 main_engine.add_app(RpcServiceApp)仓库自带的 examples/client_server/run_server.py 提供了完整的服务端启动脚本,除加载RpcServiceApp外,还演示了如何在无 GUI 的命令行模式下启动 RPC 引擎:
from vnpy_rpcservice import RpcServiceApp from vnpy_rpcservice.rpc_service.engine import RpcEngine, EVENT_RPC_LOG rpc_engine: RpcEngine = main_engine.add_app(RpcServiceApp) # 连接交易接口(以 CTP 为例) main_engine.connect(setting, "CTP") sleep(10) rep_address: str = "tcp://127.0.0.1:2014" pub_address: str = "tcp://127.0.0.1:4102" rpc_engine.start(rep_address, pub_address)从 examples/veighna_trader/run.py 可以看到,RpcServiceApp 也以注释形式作为可选应用模块预置在标准启动脚本中,取消注释即可启用。
启动模块
启动 RPC 服务模块之前,请先连接登录交易接口(连接方法见基本使用篇的"连接接口"部分)。正确连接后,VeighNa Trader 主界面【日志】栏会输出"合约信息查询成功",此时再启动 RPC 模块,可确保客户端接入后能立即查询到合约、持仓、资金等初始化信息。
连接交易接口成功后,通过菜单栏【功能】->【RPC 服务】,或点击左侧按钮栏的图标,即可进入 RPC 服务模块的 UI 界面。
配置与使用
配置 RPC 服务
RPC 服务基于ZeroMQ开发,对外提供两个通讯地址,职责各不相同:
请求响应地址(REP,Request-Reply 模式)
- 用于被动接收客户端发送过来的请求,执行对应任务后返回结果;
- 典型功能举例:
- 行情订阅;
- 委托下单;
- 委托撤单;
- 初始化信息查询(合约、持仓、资金等)。
事件广播地址(PUB,Publish-Subscribe 模式)
- 用于主动推送服务端收到的事件数据,到所有已连接的客户端;
- 典型功能举例:
- 行情推送;
- 委托推送;
- 成交推送。
两个地址均采用 ZeroMQ 的地址格式,由通讯协议(如tcp://)和通讯地址(如127.0.0.1:2014)两部分组成。RPC 服务支持的通讯协议如下:
| 协议 | 前缀 | 适用系统 | 通讯范围 | 说明 |
|---|---|---|---|---|
| TCP 协议 | tcp:// | Windows 和 Linux 均可 | 本机(127.0.0.1)或网络(局域网 IP) | 通用推荐,默认使用 |
| IPC 协议 | ipc:// | 仅 Linux(POSIX 本地端口通讯) | 仅限本机,后缀为任意字符串内容 | 低延时场景使用 |
一般推荐直接使用 TCP 协议及默认地址;对于使用 Ubuntu 系统、追求更低通讯延时的用户,可以改用 IPC 协议。
运行 RPC 服务
完成通讯地址配置后,点击【启动】按钮即可启动 RPC 服务,日志区域会输出"RPC 服务启动成功"。启动成功后,即可在另一个 VeighNa Trader 进程(客户端)中使用 RpcGateway 连接。如需停止服务,点击【停止】按钮,此时日志输出"RPC 服务已停止"。
连接客户端
VeighNa 提供了与 RpcService 配套使用的RpcGateway,作为客户端的标准接口来连接服务端并进行交易,对上层应用透明。从客户端的视角看,RpcGateway 是一个类似 CTP 的接口——因为服务端已经统一完成了外部交易账户的配置与连接,客户端只需与服务器端通讯,无需再次输入账户密码等信息。
在客户端加载 RpcGateway 接口后,进入 VeighNa Trader 主界面,点击菜单栏【系统】->【连接 RPC】,在弹出的窗口中点击【连接】即可使用。窗口中【主动请求地址】和【推送订阅地址】分别对应服务端配置的【请求响应地址】和【事件广播地址】,注意不要写反。
仓库的 examples/client_server/run_client.py 展示了客户端脚本的标准写法:加载RpcGateway作为交易接口,再叠加策略应用(如 CtaStrategyApp):
from vnpy_rpcservice import RpcGateway from vnpy_ctastrategy import CtaStrategyApp main_engine.add_gateway(RpcGateway) main_engine.add_app(CtaStrategyApp)RPC 简介
由于全局解释器锁 GIL 的存在,单一 Python 进程只能利用 CPU 单核的算力。远程过程调用(Remote Procedure Call Protocol, RPC)服务可以用于跨进程或者跨网络的服务功能调用,有效解决了上述问题:由一个特定进程连接交易接口充当服务器角色,在本地物理机或局域网内主动向其他独立的客户端进程推送事件,并处理客户端发来的相关请求。
源码剖析:RPC 底层通讯机制
理解了界面操作后,再深入仓库源码,可以看到 RPC 服务的完整实现位于 vnpy/rpc 目录,共三个文件:common.py、server.py、client.py。
通用配置:心跳参数
vnpy/rpc/common.py 定义了服务端与客户端共用的心跳常量:
HEARTBEAT_TOPIC = "heartbeat" # 心跳消息主题 HEARTBEAT_INTERVAL = 10 # 服务端心跳推送间隔(秒) HEARTBEAT_TOLERANCE = 30 # 客户端心跳容忍超时(秒)该文件还通过signal.signal(signal.SIGINT, signal.SIG_DFL)恢复 Ctrl-C 中断的默认行为,保证 RPC 接收线程可被正常打断退出。
服务端:RpcServer
vnpy/rpc/server.py 中的RpcServer是服务端核心类,其构造函数创建了两种 ZeroMQ Socket:
_socket_rep(类型zmq.REP):请求-应答模式服务端 Socket,对应【请求响应地址】;_socket_pub(类型zmq.PUB):发布-订阅模式服务端 Socket,对应【事件广播地址】。
核心方法包括:
start(rep_address, pub_address):对两个地址执行bind绑定,随后启动独立的工作线程(threading.Thread(target=self.run))并初始化下一次心跳推送时间。register(func):以函数名(func.__name__)为键,将可调用对象注册进_functions字典。RpcService 应用正是通过这一机制把行情查询、下单、撤单等引擎方法批量注册为可被客户端远程调用的函数。run():工作线程主循环。先对 REP Socket 轮询 1 秒,期间调用check_heartbeat();收到请求后通过recv_pyobj()反序列化出(函数名, 位置参数, 关键字参数),从_functions中取出对应函数执行,并将结果打包为[True, 返回值]或异常信息[False, traceback]后send_pyobj()回传。publish(topic, data):加锁后向 PUB Socket 推送[topic, data]格式的数据帧。check_heartbeat():每HEARTBEAT_INTERVAL(10 秒)向heartbeat主题推送一次当前时间戳,供客户端判断连接存活状态。
客户端:RpcClient
vnpy/rpc/client.py 中的RpcClient是客户端核心类,构造时创建_socket_req(zmq.REQ,请求-应答模式)与_socket_sub(zmq.SUB,发布-订阅模式),并设置TCP Keepalive参数(zmq.TCP_KEEPALIVE=1、空闲 60 秒),用于检测异常断连。
其最精妙的设计是__getattr__动态代理(配合@lru_cache(100)缓存):
def __getattr__(self, name: str) -> Any: def dorpc(*args: Any, **kwargs: Any) -> Any: timeout: int = kwargs.pop("timeout", 30000) # 默认超时 30 秒 req: list = [name, args, kwargs] with self._lock: self._socket_req.send_pyobj(req) n: int = self._socket_req.poll(timeout) if not n: raise RemoteException(f"Timeout of {timeout}ms reached for {req}") rep = self._socket_req.recv_pyobj() if rep[0]: return rep[1] else: raise RemoteException(rep[1]) return dorpc也就是说,只要服务端注册了某函数,客户端就能以同名属性直接调用,形如client.add(1, 3)即可触发一次远程调用;默认请求超时为 30000 毫秒,可通过关键字参数timeout覆盖。调用失败或超时时抛出RemoteException,其异常信息包含服务端traceback.format_exc()捕获的完整堆栈,便于排查服务端执行错误。
RpcClient的运行线程run()持续轮询 SUB Socket:一旦超过HEARTBEAT_TOLERANCE(30 秒)未收到服务端心跳,便调用on_disconnected()打印"RpcServer has no response..."告警;收到heartbeat主题时刷新_last_received_ping,收到业务主题时交给用户覆写的callback(topic, data)处理。订阅主题通过subscribe_topic(topic)完成,传入空字符串表示订阅全部主题。
最小可运行示例:simple_rpc
仓库在 examples/simple_rpc 中提供了不依赖 GUI 的最小演示,可用于快速理解 RPC 的本质:
- test_server.py:继承
RpcServer,在构造函数中register(self.add)注册远程函数,绑定tcp://*:2014(请求响应)与tcp://*:4102(事件广播),启动后每 2 秒向test主题发布一次服务器时间; - test_client.py:继承
RpcClient,覆写callback()打印收到的主题与数据,连接tcp://localhost:2014与tcp://localhost:4102,并循环执行tc.add(1, 3)验证远程调用。
注意这里的端口号2014/4102与 examples/client_server/run_server.py 中使用的默认地址tcp://127.0.0.1:2014/tcp://127.0.0.1:4102保持一致,是 VeighNa RPC 生态约定俗成的默认端口,日常使用可以沿用。
RPC 服务(RpcService)的应用场景
- 策略数量较多的个人用户:只需本地一条行情和交易通道,即可支持多个客户端进程同时交易,且每个客户端中的交易策略独立运行、互不影响;
- 中小型投资机构用户:可以在服务端加载各种交易接口以及 RiskManagerApp,实现一个轻量级的资管交易系统,多个交易员共享统一的交易通道,并实现基金产品级别的风险管理。
从 docs/community/info/introduction.md 对项目架构的描述看,RpcService 被定位为"允许将某一 VeighNa Trader 进程启动为服务端,作为统一的行情和交易路由通道,允许多客户端同时连接,实现多进程分布式系统",是 VeighNa 构建分布式量化交易架构的关键组件;配合 vnpy_rpcservice 中的 RpcGateway,客户端侧甚至可以实现"无本地账户、纯远程路由"的轻量接入模式。
小结
本文从功能定位、加载启动、地址配置、客户端接入四个层面完整介绍了 RpcService 的使用方法,并通过 vnpy/rpc 源码剖析了其基于 ZeroMQ REP/REQ 与 PUB/SUB 双通道的通讯模型、动态函数注册与远程代理、10 秒心跳 + 30 秒容忍的超时机制等实现细节。结合 examples/client_server 与 examples/simple_rpc 两个示例,读者可以快速在自己的环境里搭建起"一条通道、多客户端并行交易"的分布式交易架构。
【免费下载链接】vnpy基于Python的开源量化交易平台开发框架项目地址: https://gitcode.com/gh_mirrors/vn/vnpy
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考