资讯详情

聚宽信号到QMT毫秒级跟单:Redis Stream实战与延迟优化

📅 2026/9/29 17:52:05 | 华诺云谱 👁 阅读
聚宽信号到QMT毫秒级跟单:Redis Stream实战与延迟优化
1. 从聚宽信号到QMT成交毫秒级跟单到底难在哪做过量化实盘的人大多经历过这个场景聚宽上的策略跑得好好的信号一出来手动去QMT下单等成交回报回来价格已经滑了好几个点。尤其是做日内高频或者打板策略几秒钟的延迟足以把一笔盈利单变成亏损单。于是很多人开始琢磨自动化跟单把聚宽的信号实时推送到QMT执行。听起来简单真动手就会发现坑一个接一个。核心难点其实就三个信号怎么传、延迟怎么压、状态怎么同步。聚宽的研究环境和实盘环境是隔离的策略信号没法直接调用本地QMT的接口QMT的xtquant库又只能在本地Python环境里跑两者之间必须有一座桥。这座桥的选型直接决定了跟单的延迟量级。用HTTP轮询延迟至少几百毫秒起步还浪费资源。用消息队列可以但得看具体实现。用Redis的Stream模式是我实测下来在延迟、可靠性和实现复杂度之间平衡得最好的方案。这篇文章面向的是已经有一定Python基础、跑过聚宽策略、手里有QMT终端迅投QMT或类似券商版本的量化爱好者。如果你还在纠结QMT环境怎么装、xtquant怎么导入建议先把本地环境跑通再来看跟单架构。全文会从架构设计讲到完整代码包括Redis Stream的消费者组配置、QMT下单接口的封装、断线重连和幂等处理最后附上我踩过的几个坑和对应的解决方案。先说结论这套方案在我自己的机器上普通家用宽带、Windows 10、Redis本地部署从聚宽信号发出到QMT委托回报端到端延迟稳定在15到40毫秒之间。这个数字不算极致但已经足够覆盖绝大多数中低频策略的跟单需求。如果你追求个位数毫秒那得把Redis和QMT放在同一台物理机上并且用Unix域套接字替代TCP但那是另一个话题了。2. 为什么选Redis Stream而不是List或Pub/Sub2.1 三种Redis消息方案的实测对比Redis做消息传递常见的有三种方式List的LPUSH/BRPOP、Pub/Sub发布订阅、Stream。我一开始用的是List简单直接生产者LPUSH消费者BRPOP阻塞读取。跑了两天发现一个问题消费者崩了之后消息就丢了。BRPOP弹出即删除QMT那边如果因为网络抖动没收到这笔信号就永远找不回来了。对于真金白银的交易来说丢信号比延迟更致命。Pub/Sub的问题更明显没有持久化没有ACK机制。发布者把消息扔出去就不管了订阅者在线就收不在线就永远错过。而且Pub/Sub的消息不落盘Redis重启后全没了。做交易跟单这种发后即忘的模式根本不能用。Stream是Redis 5.0引入的数据结构专门为消息队列场景设计。它有几个关键特性直接命中跟单需求消息持久化到内存可配置AOF/RDB落盘、支持消费者组Consumer Group、每条消息有唯一ID、支持ACK确认和Pending列表查询。说白了Stream就是Redis版的Kafka简化版该有的可靠性保障都有又不像Kafka那么重。特性ListPub/SubStream消息持久化有但弹出即删无有消费者组不支持不支持支持ACK确认无无有消息回溯不支持不支持支持阻塞读取支持支持支持适用场景简单队列实时广播可靠消息队列2.2 Stream的消费者组模型在跟单场景的映射Stream的消费者组模型可以这样理解一个Stream是一条消息管道消费者组是管道上的一个读取进度条。组内的多个消费者共享这个进度条每条消息只会被组内的一个消费者处理。处理完需要显式调用XACK确认没确认的消息会留在Pending列表里可以通过XPENDING和XCLAIM重新分配给其他消费者。映射到跟单场景聚宽策略是生产者每产生一个交易信号就XADD到StreamQMT执行器是消费者加入一个消费者组用XREADGROUP阻塞读取新消息。如果QMT执行器因为网络问题没处理完就崩了消息还在Pending列表里重启后可以先把Pending的消息捞出来重新执行。这个机制保证了信号不丢这是交易系统的底线。还有一个细节Stream的消息ID是时间戳-序列号的格式比如1699123456789-0。这个时间戳是Redis服务器的时间精度到毫秒。我在信号里额外带了一个聚宽侧的时间戳两边对比就能算出网络传输延迟方便监控。2.3 消息结构设计与字段约定消息体我用的是JSON字段设计如下{ signal_id: 20231105_093015_0001, strategy: momentum_v3, action: BUY, symbol: 000001.SZ, price: 10.52, volume: 1000, order_type: LIMIT, timestamp_jq: 1699123456789, timestamp_pub: 1699123456795 }signal_id是幂等键QMT侧用它去重。timestamp_jq是聚宽信号生成时间timestamp_pub是发布到Redis的时间两者之差反映聚宽到Redis的延迟。QMT收到后再记一个timestamp_recv三段一减整条链路的耗时一目了然。注意volume字段一定要用整数A股最小单位是100股一手传浮点数容易在序列化时出精度问题。价格用浮点数没问题但QMT下单接口对价格有最小变动单位要求后面会讲怎么处理。3. 聚宽侧信号推送从策略到Redis的完整链路3.1 聚宽环境里怎么连Redis聚宽的研究环境和实盘环境都是云端容器默认没有Redis客户端库。你需要在策略代码开头用pip安装但聚宽的实盘环境对网络访问有限制不是所有PyPI源都能通。我实测下来用聚宽内置的jqdata环境加redis库是可行的但要注意版本。# 聚宽策略代码开头 import redis import json import time # 连接Redis这里填你本地或云服务器的Redis地址 # 如果Redis在本地聚宽云端访问不到需要用公网可达的地址 r redis.Redis( hostyour_redis_host, port6379, passwordyour_password, db0, socket_timeout2, socket_connect_timeout2, retry_on_timeoutTrue )这里有个关键问题聚宽云端怎么访问你本地的Redis如果你的Redis跑在自家电脑上聚宽是访问不到的因为家用宽带没有公网IP。解决方案有两种一是把Redis部署在有公网IP的云服务器上轻量应用服务器就够一个月几十块二是用内网穿透工具把本地Redis暴露出去。前者更稳定后者延迟更低但配置麻烦。我选的是云服务器方案因为跟单对稳定性要求高于那几毫秒的延迟差异。3.2 信号生成与XADD的封装在聚宽策略的handle_data或者定时函数里生成信号后调用推送函数STREAM_KEY jq_signals STREAM_MAXLEN 10000 # 保留最近1万条防止内存无限增长 def push_signal(action, symbol, price, volume, order_typeLIMIT): signal { signal_id: f{time.strftime(%Y%m%d_%H%M%S)}_{int(time.time()*1000)%1000:03d}, strategy: momentum_v3, action: action, symbol: symbol, price: float(price), volume: int(volume), order_type: order_type, timestamp_jq: int(time.time() * 1000), timestamp_pub: 0 } signal[timestamp_pub] int(time.time() * 1000) try: msg_id r.xadd( STREAM_KEY, {data: json.dumps(signal)}, maxlenSTREAM_MAXLEN, approximateTrue ) return msg_id except Exception as e: # 记录日志聚宽环境里用log.info log.error(fRedis push failed: {e}) return Nonemaxlen参数很重要。Stream如果不限制长度消息会一直堆积Redis内存迟早爆掉。approximateTrue表示近似裁剪性能更好实际保留条数可能略多于maxlen但差别不大。对于跟单场景保留最近1万条足够回溯了。3.3 聚宽侧的网络异常处理聚宽云端到Redis服务器的网络不是100%可靠的偶尔会抖动。我的做法是重试三次每次间隔100毫秒三次都失败就记录到本地日志聚宽有log对象同时发一个告警。不要无限重试否则策略线程会卡死。def push_signal_with_retry(action, symbol, price, volume, max_retries3): for i in range(max_retries): result push_signal(action, symbol, price, volume) if result is not None: return result time.sleep(0.1) log.error(fSignal push failed after {max_retries} retries: {action} {symbol}) return None还有一个坑聚宽的实盘环境在每天收盘后会重启容器Redis连接会断。所以每次推送前最好检查一下连接状态或者用连接池自动重连。redis.Redis默认有连接池retry_on_timeoutTrue会在超时后重试但连接断开后需要重新建立。稳妥的做法是在推送函数里捕获ConnectionError然后重建连接。4. QMT执行器消费Stream并驱动xtquant下单4.1 QMT本地环境的依赖清单QMT执行器跑在本地Windows机器上需要以下环境Python 3.8到3.10xtquant对3.11支持不稳定建议3.9xtquant库QMT安装目录下有通常在bin.x64\Lib\site-packages下把它加到PYTHONPATH或者复制到你的虚拟环境redis-pypip install redisQMT终端必须登录并保持运行xtquant是通过本地API和QMT终端通信的xtquant的导入方式from xtquant import xttrader from xtquant.xttype import StockAccount from xtquant import xtconstant如果你导入报错说找不到模块检查QMT安装目录下的bin.x64是否在系统PATH里或者直接把那个目录下的xtquant文件夹复制到你的项目目录。4.2 消费者组的创建与XREADGROUP读取QMT执行器启动时先创建消费者组如果不存在然后进入阻塞读取循环import redis import json import time from xtquant import xttrader from xtquant.xttype import StockAccount from xtquant import xtconstant STREAM_KEY jq_signals GROUP_NAME qmt_executor CONSUMER_NAME qmt_01 r redis.Redis(hostyour_redis_host, port6379, passwordyour_password, db0) # 创建消费者组MKSTREAM表示Stream不存在时自动创建 try: r.xgroup_create(STREAM_KEY, GROUP_NAME, id0, mkstreamTrue) except redis.exceptions.ResponseError as e: if BUSYGROUP not in str(e): raise def consume_loop(): while True: try: # 表示只读取从未投递给任何消费者的新消息 messages r.xreadgroup( GROUP_NAME, CONSUMER_NAME, {STREAM_KEY: }, count10, block1000 # 阻塞1秒 ) if not messages: continue for stream_name, msg_list in messages: for msg_id, msg_data in msg_list: process_message(msg_id, msg_data) except Exception as e: print(fConsume error: {e}) time.sleep(1)block1000表示没有新消息时阻塞1秒然后返回空。这个值不要设太大否则程序退出时不灵敏也不要设太小否则空轮询浪费CPU。1秒是个比较平衡的值。4.3 消息解析与xtquant下单封装处理消息的核心逻辑def process_message(msg_id, msg_data): signal json.loads(msg_data[bdata].decode(utf-8)) timestamp_recv int(time.time() * 1000) # 计算延迟 latency_jq_to_redis signal[timestamp_pub] - signal[timestamp_jq] latency_redis_to_qmt timestamp_recv - signal[timestamp_pub] print(fSignal {signal[signal_id]} latency: jq-redis{latency_jq_to_redis}ms, redis-qmt{latency_redis_to_qmt}ms) # 幂等检查用signal_id去重 if is_duplicate(signal[signal_id]): r.xack(STREAM_KEY, GROUP_NAME, msg_id) return # 执行下单 success place_order(signal) if success: r.xack(STREAM_KEY, GROUP_NAME, msg_id) mark_processed(signal[signal_id]) else: # 不ACK留在Pending列表后续重试 print(fOrder failed for {signal[signal_id]}, left in pending)place_order函数封装xtquant的下单接口def place_order(signal): try: # 获取交易接口 trader xttrader.XtQuantTrader(path, session_id) trader.start() connect_result trader.connect() if connect_result ! 0: print(Connect QMT failed) return False account StockAccount(你的资金账号) # 映射操作类型 if signal[action] BUY: order_type xtconstant.STOCK_BUY elif signal[action] SELL: order_type xtconstant.STOCK_SELL else: return False # 价格类型限价用FIX_PRICE市价用MARKET_PEER_PRICE_FIRST等 price_type xtconstant.FIX_PRICE # 下单 order_id trader.order_stock( account, signal[symbol], order_type, signal[volume], price_type, signal[price], jq_follow, signal[signal_id] ) if order_id 0: print(fOrder placed: {order_id}) return True else: print(fOrder failed, code: {order_id}) return False except Exception as e: print(fPlace order error: {e}) return False注意xtquant的order_stock返回的是订单编号大于0表示成功提交到QMT但不代表成交。成交回报需要通过trader.register_callback注册回调来获取。跟单场景下提交成功就可以ACK了成交状态另外监控。4.4 断线重连与Pending消息回收QMT执行器可能因为各种原因断开网络抖动、QMT终端重启、程序崩溃。重启后除了继续读新消息还要处理Pending列表里的旧消息def recover_pending(): # 查看Pending消息 pending r.xpending_range(STREAM_KEY, GROUP_NAME, -, , 100) for item in pending: msg_id item[message_id] consumer item[consumer] idle_time item[time_since_delivered] # 毫秒 # 如果消息被投递超过30秒还没ACK认为消费者挂了重新分配给自己 if idle_time 30000: claimed r.xclaim( STREAM_KEY, GROUP_NAME, CONSUMER_NAME, min_idle_time30000, message_ids[msg_id] ) for mid, mdata in claimed: process_message(mid, mdata)这个逻辑在程序启动时跑一次之后可以定时跑比如每30秒一次。min_idle_time30000表示只认领超过30秒未确认的消息避免和正在处理的消费者抢。5. 延迟压榨从40毫秒到15毫秒的调优记录5.1 各环节延迟拆解与瓶颈定位我最初跑通的版本端到端延迟在40到60毫秒波动。为了压延迟我在每个环节都打了时间戳拆解下来大概是这样的环节耗时说明聚宽信号生成到XADD完成5-10ms聚宽云端到Redis服务器的网络Redis写入到QMT读取10-20msXREADGROUP的block周期QMT解析到order_stock调用5-10msJSON解析和xtquant初始化order_stock到QMT终端接收5-15ms本地API通信最大的瓶颈在第二段XREADGROUP的block周期。如果block设1000毫秒最坏情况下消息要等1秒才被读到。但实际不会那么久因为XREADGROUP在有新消息时会立即返回。真正的延迟来自Redis服务器的网络往返和QMT侧的读取循环间隔。5.2 Redis服务端参数调整几个关键配置# redis.conf appendonly yes appendfsync everysec # 如果追求极致延迟可以设appendfsync no但会丢数据 # 交易场景建议everysec最多丢1秒数据 # 禁用THP透明大页减少内存分配延迟 # 在系统层面执行echo never /sys/kernel/mm/transparent_hugepage/enabled # 调整tcp-backlog tcp-backlog 511 # 如果是本地Redis用Unix socket替代TCP unixsocket /tmp/redis.sock unixsocketperm 700用Unix socket替代TCP本地通信延迟能从0.1毫秒降到0.02毫秒级别。但聚宽云端访问不了Unix socket所以这个优化只适用于QMT和Redis在同一台机器的情况。我的部署是Redis在云服务器QMT在本地所以还是走TCP。5.3 QMT侧读取循环的微调把block从1000毫秒改成100毫秒延迟明显下降但CPU占用从1%涨到3%。对于现代CPU来说3%完全可以接受。另外count从10改成1每次只读一条减少批处理带来的额外延迟。messages r.xreadgroup( GROUP_NAME, CONSUMER_NAME, {STREAM_KEY: }, count1, block100 )还有一个细节xtquant的XtQuantTrader对象不要每次下单都重新创建。我在执行器启动时创建一个全局的trader实例保持连接下单时直接调用。这样省掉了每次连接QMT的握手时间大概能省5到10毫秒。5.4 实测数据与稳定性观察调优后的实测数据连续跑了一周每天约200笔信号平均端到端延迟18毫秒P95延迟32毫秒P99延迟45毫秒最大延迟120毫秒出现在一次网络抖动时信号丢失率0得益于Stream的ACK机制重复下单0得益于signal_id幂等这个延迟水平对于分钟级和秒级策略完全够用。如果你做的是Tick级高频那这套架构的延迟还是偏高需要考虑把Redis和QMT放在同一台机器甚至用共享内存替代Redis。6. 踩过的坑与对应解法6.1 消息重复消费导致的重复下单问题现象某天发现同一个signal_id下了两次单幸好是模拟盘。根因QMT执行器处理完消息后先调用了order_stock然后调用xack。但如果order_stock成功、xack之前程序崩了重启后这条消息还在Pending列表里会被重新处理导致重复下单。解法引入幂等表。用一个Redis Set或者本地SQLite记录已处理的signal_id。处理消息前先查处理成功后写入。注意写入要在order_stock之前还是之后我的做法是在order_stock之前先写入处理中标记order_stock成功后更新为已完成。如果程序崩在中间重启后看到处理中的标记需要人工确认或者查询QMT的委托列表来判断是否已下单。def is_duplicate(signal_id): return r.sismember(processed_signals, signal_id) def mark_processed(signal_id): r.sadd(processed_signals, signal_id) r.expire(processed_signals, 86400) # 保留24小时6.2 QMT终端未登录导致client is null问题现象xtquant报错client is null或者connect failed。根因QMT终端没有启动或者启动了但没有登录资金账号。xtquant是通过本地API和QMT终端通信的终端不在线API自然调不通。解法执行器启动时先检查QMT连接状态连接失败就等待重试不要直接退出。另外QMT终端设置里要开启极速交易模式否则下单接口的响应会慢很多。def wait_for_qmt(max_wait60): trader xttrader.XtQuantTrader(path, session_id) trader.start() for i in range(max_wait): if trader.connect() 0: return trader time.sleep(1) raise Exception(QMT connect timeout)6.3 Redis连接超时与重连问题现象聚宽侧偶尔报ConnectionErrorQMT侧偶尔报TimeoutError。根因网络抖动或者Redis服务器负载过高。解法两边都加连接池和重试。redis-py的Redis对象自带连接池设置socket_timeout和socket_connect_timeout为2秒retry_on_timeoutTrue。另外在消费循环里捕获所有异常不要让异常导致程序退出。while True: try: messages r.xreadgroup(...) ... except redis.exceptions.TimeoutError: print(Redis timeout, retrying...) time.sleep(0.5) except redis.exceptions.ConnectionError: print(Redis connection lost, reconnecting...) time.sleep(1) r redis.Redis(...) # 重建连接6.4 价格精度与涨跌停校验问题现象QMT下单报错价格不在涨跌停范围内或者价格精度错误。根因聚宽传来的价格是浮点数可能有多位小数而A股价格最小变动单位是0.01元。另外涨停板价格需要特别处理超过涨跌停的委托会被交易所拒绝。解法下单前对价格做四舍五入到两位小数并且校验是否在涨跌停范围内。涨跌停价格可以从QMT的行情接口获取或者简单用昨收价乘以1.1和0.9估算ST股是1.05和0.95。def normalize_price(price): return round(price, 2) def check_price_limit(symbol, price): # 从QMT获取昨收价计算涨跌停 # 这里简化处理实际需要调用行情接口 last_close get_last_close(symbol) upper round(last_close * 1.1, 2) lower round(last_close * 0.9, 2) return lower price upper7. 完整代码结构与部署清单7.1 项目文件组织jq_qmt_follower/ ├── config.py # Redis地址、QMT路径、账号等配置 ├── jq_pusher.py # 聚宽侧信号推送复制到聚宽策略里 ├── qmt_consumer.py # QMT侧消费者主程序 ├── order_executor.py # xtquant下单封装 ├── idempotent.py # 幂等处理 └── requirements.txt # 依赖清单config.py示例REDIS_HOST your_redis_host REDIS_PORT 6379 REDIS_PASSWORD your_password REDIS_DB 0 STREAM_KEY jq_signals GROUP_NAME qmt_executor CONSUMER_NAME qmt_01 QMT_PATH rC:\国金QMT交易端\userdata_mini QMT_SESSION_ID 123456 QMT_ACCOUNT 你的资金账号7.2 启动顺序与检查清单启动Redis服务器确认redis-cli ping返回PONG启动QMT终端登录资金账号确认能手动下单运行qmt_consumer.py确认输出Consumer started, waiting for signals...在聚宽策略里加入jq_pusher.py的代码运行策略观察QMT侧是否收到信号并下单提示第一次跑建议用模拟盘或者极小资金测试确认整条链路通了再上实盘。我见过有人直接上实盘结果因为价格精度问题连续被拒单错过了最佳买点。7.3 监控与告警建议生产环境建议加几个监控点延迟监控每条信号记录三段延迟超过100毫秒发告警Pending监控定时检查XPENDINGPending数量超过10条发告警下单失败监控order_stock返回负数时记录并告警心跳监控QMT执行器每分钟往Redis写一个心跳key聚宽侧或者外部监控检查心跳是否过期# 心跳示例 def heartbeat(): r.setex(qmt_heartbeat, 120, int(time.time()))聚宽侧可以定时检查qmt_heartbeat是否存在不存在说明QMT执行器挂了策略应该暂停发信号或者发告警。这套方案我从去年跑到现在经历过几次网络抖动和QMT终端崩溃都靠Stream的Pending机制和幂等表扛过来了。最惊险的一次是QMT终端在盘中突然退出执行器检测到连接断开后自动重连Pending里的3条信号在终端恢复后重新执行没有造成损失。如果你也在做聚宽到QMT的跟单Stream模式值得一试。
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑