实时行情数据源的选择
股票与期货的实时行情获取,核心在于数据源的稳定性和延迟。券商提供的Level-1行情通常有3秒快照,Level-2行情可达到毫秒级,包含十档盘口和逐笔成交。期货方面,国内四大交易所(上期所、大商所、郑商所、中金所)通过CTP接口提供实时行情,穿透式监管要求下,个人投资者需通过期货公司柜台转发。常见方案包括:直接连接券商柜台(如恒生UFT、华锐ATP),使用第三方数据服务(如Wind、通联、新浪财经),或者自建行情网关解析组播数据。
低延迟场景建议使用FPGA加速或内核旁路技术(DPDK),普通量化交易者可从WebSocket或REST API起步。WebSocket推送tick数据,延迟约50-200ms,满足中低频策略。REST API轮询适合分钟级K线,但频繁请求易被限流。
基于WebSocket的实时数据接收
以下Python示例使用websockets库连接某券商模拟行情,接收股票和期货的tick数据,并写入队列。

import asyncio
import websockets
import json
import queue
async def market_data_client(uri, symbol_list):
async with websockets.connect(uri) as ws:
subscribe_msg = {"action": "subscribe", "symbols": symbol_list}
await ws.send(json.dumps(subscribe_msg))
tick_queue = queue.Queue(maxsize=10000)
while True:
msg = await ws.recv()
data = json.loads(msg)
if data['type'] == 'tick':
tick_queue.put(data)
print(f"{data['symbol']} 最新价 {data['price']} 成交量 {data['volume']}")
if __name__ == "__main__":
symbols = ['600519.SH', 'IF2312.CFE', 'rb2401.SHF']
asyncio.run(market_data_client('wss://api.example.com/ws', symbols))
该代码处理股票(600519.SH)、股指期货(IF2312.CFE)、螺纹钢期货(rb2401.SHF)的实时推送。队列用于解耦接收与策略计算,避免阻塞。
实时数据的清洗与存储
原始tick包含重复、乱序、缺失。清洗步骤:过滤成交量与成交额均为0的快照;按时间戳排序,修正交易所时间与本地时间差;对期货连续合约进行复权拼接。存储选用时序数据库,如InfluxDB或ClickHouse,写入频率可达每秒百万点。
股票行情需注意集合竞价阶段(9:15-9:25)的特殊处理,期货夜盘(21:00-次日2:30)跨日。使用Redis缓存最新tick,供策略快速读取,持久化用Parquet文件按日分区。
量化交易信号生成:VWAP与ATR
日内交易常用VWAP(成交量加权平均价)判断趋势,ATR(平均真实波幅)设置止损。计算逻辑:
VWAP = 累计成交额 / 累计成交量
ATR = 过去N根K线的真实波幅均值,真实波幅 = max(最高-最低, |最高-昨收|, |最低-昨收|)
当实时价格上穿VWAP且ATR放大,产生买入信号;下穿VWAP且ATR放大,产生卖出信号。以下代码基于tick数据实时计算VWAP和ATR。
import numpy as np
from collections import deque
class VWAPATRStrategy:
def __init__(self, atr_period=14):
self.cum_amount = 0.0
self.cum_volume = 0
self.highs = deque(maxlen=atr_period)
self.lows = deque(maxlen=atr_period)
self.closes = deque(maxlen=atr_period)
self.prev_close = None
self.position = 0
def on_tick(self, price, volume, high, low):
self.cum_amount += price * volume
self.cum_volume += volume
vwap = self.cum_amount / self.cum_volume if self.cum_volume > 0 else price
if self.prev_close is not None:
tr = max(high - low, abs(high - self.prev_close), abs(low - self.prev_close))
self.highs.append(high)
self.lows.append(low)
self.closes.append(price)
self.prev_close = price
atr = np.mean([max(h-l, abs(h-pc), abs(l-pc)) for h, l, pc in zip(self.highs, self.lows, self.closes)]) if len(self.highs) >= 2 else 0
if price > vwap and atr > 0 and self.position == 0:
self.position = 1
print(f"买入 价格{price} VWAP{vwap:.2f} ATR{atr:.2f}")
elif price < vwap and atr > 0 and self.position == 1:
self.position = 0
print(f"卖出 价格{price} VWAP{vwap:.2f} ATR{atr:.2f}")
策略运行在tick级别,每笔成交更新VWAP和ATR,触发买卖。实际部署需加入滑点、手续费、仓位管理。
期货特殊处理与风控
期货有杠杆、保证金、强平机制。实时风控需监控:账户权益、可用资金、持仓风险度。使用CTP接口的OnRtnTrade和OnRtnOrder回调实时更新仓位。当风险度超过80%时,自动减仓。
跨期套利需同时订阅近月与远月合约,计算价差。价差突破布林带上轨做空价差,突破下轨做多价差。套利指令使用交易所组合保证金,降低资金占用。
夜盘波动剧烈,ATR阈值应动态调整。例如螺纹钢夜盘ATR通常是日盘的1.5倍,需乘以系数。
回测与实盘一致性
实时行情与历史回测存在差异:历史数据是快照,实时是逐笔。回测需使用tick级历史数据,模拟排队成交。实盘使用限价单,未成交则撤单重发。
使用事件驱动回测框架,将实时tick注入策略,与实盘共用同一套信号逻辑。通过比较回测与实盘的滑点分布,优化下单算法。
性能优化与部署
Python适合策略研发,但实盘延迟高。关键路径用C++或Rust重写,通过pybind11暴露接口。行情接收与策略计算分离,使用共享内存传递tick。容器化部署(Docker)保证环境一致,Kubernetes管理多策略实例。
监控指标:行情延迟、策略计算耗时、订单响应时间。延迟超过阈值自动切换备用数据源。