资讯动态

外汇行情接入实战:WebSocket推送与REST查询的分工配合

发布时间:2026/9/16 5:16:09 来源:尧图企业网站定制
1. 为什么实时外汇行情接入不能靠REST轮询硬扛先说一个我自己的经历。去年做一个外汇趋势跟踪的小策略最开始为了省事用 requests 每隔一秒去拉一次报价接口本地拼成 K 线跑了一个星期问题不断接口时不时返回限流错误、行情跳变的那一两秒恰好没抓到、回测数据和实盘数据对不齐。后来把行情接收部分整体切成 WebSocket 推送REST 只用来补历史和查 K 线整个数据链路才真正稳下来。这个经历其实说透了外汇行情接入里最重要的一件事实时行情和K线查询根本是两种不同性质的数据需求不该用同一种方式去拿。实时报价是高频、低延迟、要求不丢数据的推送流REST 是低频、按需、允许一定延迟的请求响应。很多新手上来就把所有操作都做成 REST 请求这是最常见的入门坑。那这两者到底差在哪里我用一个生活化的类比来解释REST 轮询就像你每隔几秒钟去窗口问一次现在价格是多少WebSocket 推送则是柜台主动打电话告诉你价格一变我就通知你。前者能不能用在数据变化慢、实时性不高的场景下能用比如每天收盘后拉一次日线完全没问题。但在外汇这种 24 小时连续交易、价格毫秒级跳动的市场里轮询有两个硬伤。第一是延迟和空隙。秒级轮询意味着你的数据最多只能准确到上次请求的瞬间两次请求之间发生的快速行情波动对你来说完全是盲区。做策略回测的时候这些间隙会被当成价格没变实盘里就会表现为信号触发点不准确。第二是带宽和限流成本。每秒轮询一次24 小时就是 86400 次请求绝大多数公开或半公开数据源都会给你限额。而过高的请求频率还会触发 IP 封锁。反过来WebSocket 一次建立一个长连接价格变化时由服务端主动推给你同样的信息量请求量几乎为零对双方都友好。做外汇行情接入正确姿势是WebSocket 收实时 tickREST 查历史 K 线、补缺口、做初始化。两者配合才能既保证实时性又不把接口请求量打爆。本文后面的内容就是围绕这套组合来展开的。适合读这篇内容的读者我默认是这么一个人熟悉 Python 基础语法知道 requests 和 pandas 的大致用法想给自己的量化策略接入真实外汇行情或者单纯想搞明白 WebSocket 和 REST 在实战里到底怎么分工。如果你完全没写过 Python建议先补一补基础语法再回来看我用到的代码量虽然不大但还是需要一点编程底子的。2. 环境准备与数据源选型先解决从哪拿数据这个根本问题2.1 选一个能真正练手的外汇数据源做外汇行情接入第一步不是写代码而是选数据源。这里我不建议一开始就去接那种需要机构开户、门槛很高的数据终端也建议别用没有真实撮合逻辑的模拟接口。外汇市场没有统一交易所流动性分散在全球各大银行和经纪商手里所以你能拿到的行情本质上都是某个数据服务商自己的报价源。对个人开发者和量化爱好者来说选数据源基本上看三件事有没有免费的练习账户、有没有完整的 REST 和 WebSocket 接口文档、文档里的示例是不是给得够明白。我在项目里用的是 OANDA 的 v20 API。它是少数几个同时提供完整 REST 和 WebSocket 流式接口、并且有免费 practice 账户的零售外汇经纪商。注册过程不复杂填个邮箱拿到 API Token 和账户 ID 就能开始调。当然如果你因为网络条件或注册门槛不方便用选其他提供 WebSocket 流的外汇数据商也完全可以代码逻辑是通用的只是字段名和鉴权方式需要替换。注册好之后你手上有两样关键东西API Token像一把钥匙和 Account ID。鉴权方式建议用请求头里的 Bearer Token而不是把 Token 拼在 URL 里后者容易出现在日志中有泄露风险。2.2 环境依赖与目录结构代码侧我用的是 Python 3.10 以上版本依赖库遵循够用就行原则一共三个websockets异步 WebSocket 客户端库基于 asyncio连接管理做得干净。requests同步 HTTP 库用来查 REST K 线。pandas只在数据处理阶段使用用来整理 K 线 DataFrame。安装就一行命令pip install websockets requests pandas项目结构我建议这样组织后面维护起来不混乱forex_ingest/ ├── config.py # Token、账户ID、行情品种等配置 ├── rest_client.py # REST K线查询封装 ├── ws_client.py # WebSocket实时行情客户端 ├── kline_builder.py # 实时tick转K线REST与WebSocket数据合并 └── main.py # 主程序入口这个结构不复杂但每个文件的职责是清晰的。我做行情接入项目的经验是千万别把 REST 和 WebSocket 代码写在一个文件里因为它们的生命周期完全不一样REST 是一次性请求-响应WebSocket 是长期存在的连接。混在一起后期排查问题会非常痛苦。2.3 先明确要拿哪些行情品种和存储方式外汇品种代码在 OANDA 的格式是基础货币_计价货币比如欧元兑美元写成EUR_USD英镑兑日元写成GBP_JPY。这个下划线的格式很容易被忽略但它和主流交易软件里的EURUSD写法不同拼接 URL 时要注意。我一般把要订阅的品种写进 config.py用列表维护方便增删INSTRUMENTS [EUR_USD, GBP_JPY, AUD_CAD]存储方面练手阶段不建议一开始就上数据库先落 CSV 文件把接收到的原始 tick 以及实时生成的 K 线各存一份。这样做的理由有两个一是 CSV 排查问题直观用文本编辑器打开就能看到数据长什么样二是避免过早引入数据库的维护成本等你确认数据格式稳定了再迁到 SQLite 或 TimescaleDB 都不迟。我见过不少人在第一步就栽跟头数据源没选好、字段英文看不懂、不会看接口文档然后项目就搁浅了。所以这部分虽然不写代码但值得你花时间认真对待。选型选对了后面写代码是水到渠成的事。3. WebSocket 接入的完整实现鉴权、消息解析、心跳与断线重连3.1 建立连接与鉴权WebSocket 接入的核心代码其实不长。下面这一段可以建立到 OANDA 行情流式接口的连接并开始接收数据。import asyncio import json import websockets from config import API_TOKEN, ACCOUNT_ID, INSTRUMENTS STREAM_URL ( wss://stream-fxpractice.oanda.com/v3/accounts/ f{ACCOUNT_ID}/pricing/stream ) async def price_listener(): headers { Authorization: fBearer {API_TOKEN}, Content-Type: application/json, } instruments_str ,.join(INSTRUMENTS) url f{STREAM_URL}?instruments{instruments_str} async with websockets.connect(url, additional_headersheaders) as ws: print(f已连接订阅品种: {instruments_str}) while True: try: message await asyncio.wait_for(ws.recv(), timeout30) data json.loads(message) if data.get(type) PRICE: instrument data[instrument] bid_price float(data[bids][0][price]) ask_price float(data[asks][0][price]) mid_price (bid_price ask_price) / 2 print( f{instrument} bid{bid_price:.5f} fask{ask_price:.5f} mid{mid_price:.5f} ) except asyncio.TimeoutError: # 30秒没有收到任何消息发送一个ping探测连接 await ws.ping() print(发送ping保持连接活跃) if __name__ __main__: asyncio.run(price_listener())几个值得注意的细节第一asyncio.wait_for(ws.recv(), timeout30)是很多 WebSocket 客户端容易漏掉的关键点。外汇行情的 WebSocket 服务端如果一段时间没有价格更新会发送心跳消息来保活但也会有极端情况——比如市场休市或品种流动性枯竭长时间没有消息。如果你不加超时控制程序就会一直阻塞在recv()上连接变成僵尸连接。加了 30 秒超时后一旦发现长时间没消息就主动 ping 一下探测连接状态。第二消息类型的判断不能省。OANDA 的流式接口在连接建立后、实际行情推送前会先发一条类型为HEARTBEAT的消息。如果你不判断type字段就直接按价格消息解析大概率第一个报错就出在这里。我用if data.get(type) PRICE把非价格消息直接跳过心跳消息由服务端定期发来前端不需要专门处理。第三bids 和 asks 是数组而不是单个对象。这是外汇行情的特殊性同一时刻同一个品种可能有多档报价类似股票盘的买卖一档、二档每档有自己的价格和流动性。实际接 tick 时通常只需要第一档即bids[0]和asks[0]但如果你要分析盘口深度就需要遍历整个数组。这也是容易被 API 文档误导的地方——很多文档示例只显示一条 bid实际情况是一组。3.2 断线重连机制指数退避才是上策WebSocket 长连接最怕的就是断线。网络抖动、服务端重启、长时间空闲被中间网关断开都是实际会碰到的情况。如果你的程序没有重连机制一个断线可能导致整夜的数据全部丢失第二天的策略信号直接就残废了。我用的重连方案是指数退避断线后先等 2 秒重连连续失败则等待时间翻倍4 秒、8 秒、16 秒最多封顶 60 秒。这样既能在瞬时抖动后快速恢复又不会在服务端持续不可用时高频轰炸。import asyncio import websockets async def run_with_reconnect(listener_func, max_retries10): retry_delay 2 retries 0 ws_connected False # 先确保WebSocket连接成功再进入监听循环 while True: try: await listener_func() # listener_func正常退出说明连接被服务端关闭 ws_connected True except websockets.exceptions.ConnectionClosed as exc: print(f连接关闭: {exc.code} {exc.reason}) except Exception as exc: print(f其他错误: {type(exc).__name__}: {exc}) if retries max_retries: print(重试次数超过上限退出) break print(f{retry_delay}秒后进行第{retries 1}次重连...) await asyncio.sleep(retry_delay) retries 1 retry_delay min(retry_delay * 2, 60)这里需要注意一个容易忽视的细节重连次数计数器在成功连接后要重置。很多人会把累加逻辑写成全局计数结果程序运行两天后计数触顶一个正常断线就导致程序彻底退出。正确做法是每次成功收到第一条数据后把retries和retry_delay归零重新计算退避周期。另外重连时最好重新创建整个 WebSocket 连接对象而不是尝试复用旧连接。WebSocket 协议里连接一旦关闭旧的连接对象就不可用了任何基于它的收发操作都会抛异常。每次重连实际是新建连接、重新订阅品种的过程代码逻辑上就是一个完整的重新初始化。3.3 为什么用 asyncio 而不是多线程这里说一个很多教程不会讲透的问题为什么 WebSocket 客户端普遍用 asyncio 而不是开一个线程循环收数据最核心的原因是WebSocket 的recv()是一个阻塞操作如果用多线程你得管理线程的生命周期、线程安全的消息队列、线程退出时如何干净地关闭资源复杂度直线上升。而 asyncio 是单线程协作式调度await挂起时事件循环可以去做别的事情天然适合多个连接并行、连接和主程序之间共享数据的场景。实际项目中我通常会在同一个事件循环里跑三个任务一个负责 WebSocket 行情接收一个负责定时把内存中的数据批量写入 CSV还有一个负责监控程序健康状态。三个任务之间用 asyncio.Queue 传递数据完全不会互相阻塞。这个结构等你需要扩展多个品种或者其他数据源时优势非常明显。当然asyncio 也有它的小脾气不能用阻塞式的同步代码比如time.sleep()、requests.get()直接写在async函数里否则会卡住整个事件循环。写 WebSocket 客户端时最容易犯这个错——在接收行情的同时想去请求一下 REST 接口直接调 requests结果行情接收突然延迟。解决办法是切换到httpx.AsyncClient或把耗时操作丢给asyncio.to_thread去执行。4. REST K线查询的实现补历史、建框架、对数据4.1 K线格式与请求参数REST 在整套逻辑里负责的是补和查两个动作。所谓补是指 WebSocket 只推送连接后产生的新数据你刚启动程序时手里是没有历史 K 线的。所谓查是指策略需要定期查看历史走势、计算技术指标这时候不需要实时推送直接按需请求 K 线即可。OANDA 的 REST K 线接口请求方式是 GET三个核心参数需要理解透彻参数说明我常用的值granularityK线周期M1/M5/M15/H1/H4/D 等M5count返回的K线根数最多5000500priceK线价格类型M为中间价、B为买价、A为卖价M这里的price参数特别值得注意。外汇市场的报价天然包含买入价和卖出价两价之间存在点差。如果策略回测时用中间价M实盘交易时实际成交是用买价或卖价点差会造成真实的成本差异。所以选哪一个价格类型不是一个随意的决定做策略研究用中间价比较多做实盘信号验证建议分别拉取买卖价来对比滑点影响。4.2 完整代码把REST响应转成DataFrame下面是我实际在用的 REST 查询函数封装返回结果直接转成 pandas DataFrame时间列设为索引方便后续和 WebSocket 数据合并import requests import pandas as pd from config import API_TOKEN REST_URL https://api-fxpractice.oanda.com/v3/instruments def fetch_candles( instrument: str, granularity: str M5, count: int 500, price: str M, ) - pd.DataFrame: url f{REST_URL}/{instrument}/candles headers {Authorization: fBearer {API_TOKEN}} params { granularity: granularity, count: count, price: price, } resp requests.get(url, headersheaders, paramsparams) resp.raise_for_status() data resp.json() records [] for candle in data[candles]: # OANDA的candle里time是ISO格式 # 需要先判断这根K线是否已经闭合 records.append({ time: candle[time], open: float(candle[mid][o]), high: float(candle[mid][h]), low: float(candle[mid][l]), close: float(candle[mid][c]), volume: int(candle[volume]), complete: candle[complete], }) df pd.DataFrame(records) df[time] pd.to_datetime(df[time]) df.set_index(time, inplaceTrue) return df这个函数有几个容易踩的雷我逐个说。第一resp.raise_for_status()必须有。很多人写 requests 请求不检查状态码结果接口返回 400 或 429限流时程序要么静默失败要么在解析 JSON 时才报错到那时候错误信息已经很不直观了。提前 raise错误信息会包含服务端返回的状态码和原因排查问题能省大量时间。第二K线是否闭合的判断。candle[complete]这个字段表示这根 K 线对应的周期是否已经走完。比如你请求 M5 的当前 K 线如果这一根还没走满 5 分钟它就是一个未完成的临时状态下一分钟数据还会变。做回测和策略计算指标时默认只使用completeTrue的K线否则计算结果会失真。触发信号时用的最新K线可以用未闭合的但要清楚它是不稳定的。第三响应里所有时间都是 UTC 时间。这一点放到第 5 章细讲但这里先埋个伏笔。REST 返回的时间格式是带时区偏移的 ISO 字符串直接拿来做本地时间K线切分点会整体错位。4.3 REST限流策略再简单的请求也要有礼貌REST 接口虽然简单但也不是让你无限制调用的。OANDA practice 账户对 REST 的限流策略会根据当前系统负载动态调整实测下来高频连续请求很容易触发 429。我踩过几次之后总结了两条经验一是在每两次 REST 请求之间加至少 200 毫秒的间隔。如果你的程序需要批量拉取多个品种的 K 线写一个简单的限速器比如用time.sleep(0.2)比一次性并发请求几十个接口要稳得多。二是把数据做本地缓存。同一个品种同一个周期的 K 线数据短时间内不会变化没必要每次策略循环都重新请求。我在本地用一个字典缓存from datetime import datetime, timedelta candle_cache {} def get_candles_cached(instrument, granularityM5, max_age_seconds60): key f{instrument}_{granularity} now datetime.utcnow() if key in candle_cache: cached_time, df candle_cache[key] if (now - cached_time).total_seconds() max_age_seconds: return df df fetch_candles(instrument, granularitygranularity) candle_cache[key] (now, df) return df这个缓存设计很简单但很实用。它保证你在 60 秒内重复查询同一个品种时不会每次都打到服务端限流风险直降。5. 实战踩坑记录时序、时区与数据口径才是真正的拦路虎5.1 UTC时区导致的K线切分错位做外汇行情接入最容易让人半夜挠头的问题不是 WebSocket 连不上而是时间对不上。外汇市场是 24 小时交易K线的切割天然依赖时间基准。OANDA 服务器返回的时间一律是 UTC而大部分本地策略运行时用的是本地时区。举一个我真实遇到的例子程序在本地 18:30 打印了一根 M5 K 线的收盘时间是18:35但我在交易软件里看到对应 K 线的时间是02:35差了整整 16 个小时。原因是交易软件默认把时间显示成了东八区而我直接把 UTC 时间拿来用了。解决办法有两个层面。如果你只是需要数据用于分析建议统一用 UTC 时间存储所有计算都基于 UTC避免时区混淆如果你需要把K线和本地交易时间对齐就用pandas的时区转换df.index df.index.tz_localize(UTC).tz_convert(Asia/Shanghai)这里有个顺序问题很容易弄错如果原始时间已经是带时区信息的ISO 字符串结尾带 Z 或 00:00你需要用tz_convert而不是tz_localize。tz_localize用于给一个无时区的时间指定时区tz_convert用于把已有时区的时间转换到另一个时区。用反了会直接抛异常或产生错误偏移。外汇市场还有一个特性要注意周末休市和节假日休市会造成 K 线时间轴上的缺口。周五纽约收盘后到周一悉尼开盘前市场没有行情REST 查询到的 K线数据在这个时间段是空的WebSocket 也不会推送任何价格。如果用连续时间索引去处理K线这个缺口会导致时间序列出现空值进而让移动平均等技术指标计算出现断层。做数据合并时我用的是只保留有交易的时间点而不是填充空值避免指标被无意义的数据污染。5.2 WebSocket心跳消息与价格消息的混排WebSocket 接入过程中服务端推送的消息并非只有价格。OANDA 的心跳消息结构大致是{ type: HEARTBEAT, time: 2024-11-20T18:00:00.000000000Z }和价格消息的type字段值不同但结构上有重叠。如果你把所有消息一股脑塞给K线构建器心跳消息会被当成异常数据丢弃或产生解析错误。我的处理策略是在消息入口做一个分发器async def handle_message(raw_message): data json.loads(raw_message) msg_type data.get(type) if msg_type PRICE: await process_price(data) elif msg_type HEARTBEAT: # 更新最后心跳时间用于健康检查 last_heartbeat_time data.get(time) else: # 未知类型记录日志便于排查 logger.warning(f未知消息类型: {raw_message})这个分发的价值在于以后接入更多数据源或消息类型时你只需要扩展process_price之外的处理分支主循环的稳定性不会受影响。未知类型不要静默丢弃记录日志能帮你提前发现服务端协议变更的情况。5.3 浮点计算精度够用吗这类问题的处理思路外汇报价通常是 5 位小数如1.08432很多刚接触的人会担心 Python float 的精度不够想引入 decimal.Decimal。我的建议是判断够不够用取决于你的用途。如果只是计算指标、画K线图、做信号判断float 完全够用因为外汇价格小数位数很少float 的误差在这个量级下远小于点差本身对策略结果不会有实质影响。但如果你的程序涉及资金计算、盈亏结算、需要精确到分甚至厘那应该用 Decimal 把它当钱来算和行情数据分开处理。行情侧我用 float资金侧用 Decimal边界清晰互不干涉。这也是防止自己掉进过度设计陷阱的一层意识。5.4 断线重连后K线是否有缺口需要自己判定WebSocket 断线重连期间市场照样在走价格照样在变。断线 30 秒M5 的 K 线可能已经更新了好几个 tick。如果你只依赖 WebSocket 数据重连后就会产生一段数据空洞。这个问题没有银弹但有一个很实用的兜底方案每次 WebSocket 重连成功后立即调用一次 REST 接口拉取最近几根K线和本地已有的K线做一次末根对齐。如果 REST 返回的最后一根完整K线在本地不存在就用便 REST 数据补进去。这样即使断线期间有数据丢失也能在重连后快速修复不至于带着缺口继续跑。对齐操作的核心代码本质上是把两个 DataFrame 做一次索引差集def align_missing_candles(local_df, rest_df): missing_index rest_df.index.difference(local_df.index) missing_df rest_df.loc[missing_index] combined pd.concat([local_df, missing_df]) combined.sort_index(inplaceTrue) # 重复索引去重保留REST数据 return combined[~combined.index.duplicated(keeplast)]这段代码不复杂但它是整个数据链路稳定性的最后一道保险。没有它断线重连只能保证连接恢复不能保证数据完整有了它才能称得上一个合格的实时行情接入工程。6. 进阶实践把 REST 和 WebSocket 拼成一套实时K线维护器6.1 整体架构与数据流到这里实时行情和 K 线查询两个能力都具备了但它们是互相独立的。实战中的最后一步是把它们组织成一个完整的程序启动时先用 REST 拉历史 K 线建立初始框架然后 WebSocket 持续接收 tick把 tick 实时聚合到当前K线上同时定时把数据落盘。数据流可以描述成一句话REST 管过去WebSocket 管现在拼在一起就是连续。我维护的KlineBuilder类大致是下面这个结构class KlineBuilder: def __init__(self, instrument, granularityM5): self.instrument instrument self.granularity granularity self.df self._load_history() self.current_candle_start self._next_candle_start() def _load_history(self): # 启动时从REST拉最近500根已闭合K线 df fetch_candles(self.instrument, self.granularity, count500) return df[df[complete]] def _next_candle_start(self): # 根据当前UTC时间计算下一根K线的开始时间 now pd.Timestamp.utcnow() # 按分钟周期取整 floor now.floor(self.granularity) return floor pd.Timedelta(self.granularity) def update_price(self, bid, ask, tsNone): # 收到一个tick时将价格合并进当前K线 if ts is None: ts pd.Timestamp.utcnow() mid_price (bid ask) / 2 if ts self.current_candle_start: self._finalize_current_candle() self.current_candle_start self._next_candle_start() self._update_current_candle(mid_price)这个类的核心逻辑其实只有一个判断当前 tick 的时间戳是否已经超过当前K线的结束时间。一旦超过就把当前K线定稿finalize开始下一根。外汇市场的K线边界是严格按照 UTC 时间对齐的所以这里不能用接收了多少个 tick来决定K线的切换必须用时间戳判断。这是整个实时K线维护器最关键的一行逻辑也是最容易被搞错的地方。6.2 多品种并行接收时的资源管理如果你要同时订阅多个品种比较合理的做法是每一个品种创建一个独立的KlineBuilder实例然后在 asyncio 里为每个品种创建一个监听任务async def listen_one_instrument(instrument): builder KlineBuilder(instrument) async for price_data in subscribe(instrument): builder.update_price(price_data[bid], price_data[ask]) async def main(): builders {inst: KlineBuilder(inst) for inst in INSTRUMENTS} tasks [asyncio.create_task(listen_one_instrument(inst)) for inst in INSTRUMENTS] await asyncio.gather(*tasks)这样做的好处是单个品种的异常不会影响其他品种的连接。比如GBP_JPY的 WebSocket 连接断线了正在重试EUR_USD的连接还在正常收数据。我曾经把多个品种写在一个连接里用一条循环处理结果一个品种的消息阻塞导致所有品种都停滞最后不得不重构成一品种一连接一任务的模型。后续排查问题的体验好了非常多。6.3 数据落盘与程序退出时要注意的事最终的数据落盘我建议使用追加写入 周期 flush的策略。每收到一个 tick 就写一次文件磁盘 I/O 会成为性能瓶颈但长时间攒在内存里不写一旦程序崩溃缓冲区的数据就全丢了。我的做法是 WebSocket 主循环里每收到一条价格就更新内存中的最新K线但落盘操作放在一个独立的定时任务里每 5 秒执行一次批量写入。程序退出时还有最后一个容易忽视的坑用户按 CtrlC 退出时WebSocket 连接如果需要优雅关闭必须给事件循环一个处理信号的逻辑。否则连接不会正常断开对某些数据源来说异常掉线后短时间内的重连会被拒绝。import signal def request_shutdown(): print(收到退出信号正在保存数据...) # 触发当前任务的安全退出 asyncio.get_running_loop().stop() loop asyncio.get_event_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, request_shutdown)这段信号处理能让程序在退出前完成最后一轮数据落盘避免跑了一整天退出时发现最后10分钟的数据丢了这种憋屈的事情。6.4 一点个人体会把 REST 和 WebSocket 组合在一起做一套能持续运行的实时行情维护器整个过程最花时间的往往不是代码本身而是对细节的打磨时间切分点、消息类型判断、断线数据修复、退出时数据保存。这些细节点单独拿出来都不值一提但它们组合起来就是你从能跑通的Demo走向能稳定运行的工程之间那一段必经的爬坡路。爬过去之后后面接策略、接信号、接交易执行就都顺了。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价