ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

OKX行情接口实战:REST与WebSocket双通道接入及断线重连

OKX行情接口实战:REST与WebSocket双通道接入及断线重连 1. 从零拆解 OKXapi_1一个行情接口项目到底在做什么很多人第一次看到“OKXapi_1”这个标题第一反应是“这不就是个接口封装吗”。但真上手做过行情对接的人都知道一个能跑、能扛、能长期稳定运行的行情接口项目远不是调几个 REST 接口那么简单。它要解决的核心问题其实有三个第一怎么把交易所的实时行情拉下来第二怎么在本地高效地分发和消费这些数据第三怎么在断线、限频、数据缺失这些异常情况下保证程序不崩、数据不断。我做的这个 OKXapi_1定位就是一个轻量级的行情接入层。它同时支持 REST 和 WebSocket 两条通道REST 负责拉取历史 K 线、交易对列表、账户余额这类“请求-响应”式的数据WebSocket 负责订阅实时行情推送比如 ticker、深度、成交记录。整个项目用 Python 写核心依赖就是 requests、websocket-client 和 asyncio 这几样没有引入过重的框架目的是让一个刚学 Python 的人也能看懂、能改、能跑起来。适合谁来参考如果你正在学 Python想找一个真实项目练手或者你在做量化策略的原型验证需要一个稳定的行情数据源又或者你只是想搞清楚 WebSocket 实时推送到底是怎么一回事那这个项目的内容应该能帮到你。下面我会把整个项目的设计思路、核心代码、踩过的坑和排查技巧全部摊开来讲尽量做到你照着做就能复现。2. 整体架构设计与技术选型思路2.1 为什么 REST 和 WebSocket 要同时用这是整个项目最核心的设计决策。很多人会问既然 WebSocket 能实时推送为什么还要 REST反过来既然 REST 能拿到数据为什么还要折腾 WebSocket答案在于两者的适用场景完全不同。REST 是“你问一次它答一次”适合拉取非实时、低频、结构化的数据。比如你要获取 BTC-USDT 最近 100 根 1 小时 K 线这种请求一天可能只跑几次用 REST 最合适。而 WebSocket 是“你订阅一次它持续推”适合高频、实时、流式的数据。比如你想实时监控 ticker 价格变化用 REST 轮询的话每秒发一次请求既浪费带宽又容易触发限频而 WebSocket 只需要建立一次连接之后数据自动推过来。我在项目里的分工是这样的通道负责内容调用频率典型场景RESTK线历史、交易对列表、账户信息低频按需调用策略初始化、回测数据准备WebSocketticker、深度、成交推送高频持续订阅实时监控、信号触发这个分工不是拍脑袋定的。OKX 的 REST 接口对频率有明确限制比如某些公共接口是每秒 20 次超过就会被限流甚至临时封禁。而 WebSocket 订阅之后数据推送是服务端主动发起的你只需要处理好消息就行。所以把高频需求交给 WebSocket把低频需求交给 REST是最合理的资源分配方式。2.2 Python 技术栈的选择理由项目用 Python 而不是 Go 或 Java原因很直接上手快、生态全、调试方便。对于行情接入这种 IO 密集型任务Python 的 asyncio 配合 websocket-client 已经足够应付。而且 Python 的 requests 库在处理 REST 请求时非常成熟异常处理、超时设置、重试逻辑都有现成方案。具体依赖清单如下pip install requests websocket-client如果你用的是 Python 3.10 以上版本asyncio 的原生支持已经很好不需要额外装 uvloop 之类的加速库。对于初学者来说少一个依赖就少一个坑。注意不要用websockets库和websocket-client库混着写两者的 API 风格完全不同。websocket-client是同步风格的配合线程用websockets是纯异步的配合 asyncio 用。项目里我选的是websocket-client因为它的回调式写法对新手更友好不需要理解 async/await 的完整体系就能跑起来。2.3 项目目录结构一个清晰的项目结构能省掉后面很多找文件的时间。我的组织方式是这样的okxapi_1/ ├── config.py # API Key、密钥、 passphrase 配置 ├── rest_client.py # REST 接口封装 ├── ws_client.py # WebSocket 连接与订阅 ├── data_handler.py # 数据解析与存储 ├── main.py # 入口启动逻辑 └── utils.py # 日志、时间戳转换等工具函数这个结构的好处是职责分离。REST 出问题了只看rest_client.pyWebSocket 断线了只看ws_client.py不会在一大坨代码里翻来翻去。对于后续要加数据库、加策略模块的情况直接新增文件就行不用动核心逻辑。3. 核心细节解析与实操要点3.1 API Key 的申请与安全配置OKX 的 API Key 需要三个东西API Key、Secret Key、Passphrase。前两个是标准的密钥对第三个是你在创建 API Key 时自己设置的密码。这三个缺一不可少一个都会报签名错误。申请流程这里不展开重点说配置。绝对不要把密钥硬编码在代码里然后上传到公开仓库这是新手最容易犯的致命错误。我的做法是放在config.py里然后用环境变量覆盖import os API_KEY os.getenv(OKX_API_KEY, your_api_key_here) SECRET_KEY os.getenv(OKX_SECRET_KEY, your_secret_key_here) PASSPHRASE os.getenv(OKX_PASSPHRASE, your_passphrase_here) BASE_URL https://www.okx.com这样本地开发时可以直接改默认值部署到服务器时用环境变量注入两边不冲突。另外创建 API Key 时权限要最小化如果只是拉行情就只勾选“读取”权限不要开交易和提现。这是安全底线。提示OKX 的模拟盘和实盘用的是不同的域名和不同的 API Key。模拟盘的 REST 基础地址是https://www.okx.com加上特定的模拟盘标识WebSocket 也有对应的模拟盘地址。如果你用实盘的 Key 去连模拟盘会直接返回 401 错误。3.2 REST 请求的签名机制OKX 的私有接口需要签名签名算法是 HMAC-SHA256。具体步骤是先把时间戳、请求方法、请求路径、请求体拼成一个字符串然后用 Secret Key 做 HMAC-SHA256 加密最后 Base64 编码。import hmac import base64 import hashlib import datetime def get_signature(timestamp, method, request_path, body, secret_key): message timestamp method.upper() request_path body mac hmac.new( secret_key.encode(utf-8), message.encode(utf-8), hashlib.sha256 ) return base64.b64encode(mac.digest()).decode(utf-8)时间戳的格式必须是 ISO 8601比如2024-01-15T10:30:00.000Z。这里有个坑时间戳和服务器时间偏差不能超过 30 秒否则会报Invalid timestamp错误。如果你的服务器时间不准先去同步 NTP。def get_timestamp(): now datetime.datetime.utcnow() return now.strftime(%Y-%m-%dT%H:%M:%S.%f)[:-3] Z注意[:-3]是取毫秒部分OKX 要求精确到毫秒。这个细节文档里写得不明显但错了就是签名失败。3.3 WebSocket 连接与订阅流程WebSocket 的流程比 REST 多一步“订阅”。连接建立之后你需要发送一个订阅消息告诉服务器你要哪些频道的数据。订阅消息的格式是 JSONsubscribe_msg { op: subscribe, args: [ {channel: tickers, instId: BTC-USDT}, {channel: books5, instId: ETH-USDT} ] }op是操作类型subscribe表示订阅unsubscribe表示取消订阅。args是一个数组可以一次订阅多个频道和交易对。channel是频道名instId是交易对标识。连接建立后服务器会先返回一个确认消息告诉你订阅是否成功。之后就是持续的数据推送。每条推送消息都包含arg频道信息和data实际数据。import websocket import json def on_message(ws, message): data json.loads(message) if event in data: print(f事件: {data[event]}) elif arg in data: channel data[arg][channel] print(f收到 {channel} 数据: {data[data][0]})on_message是回调函数每收到一条消息就会被调用一次。这里要注意消息可能是订阅确认、错误事件或者正常数据需要用event和arg字段来区分。3.4 心跳保活机制WebSocket 连接如果长时间没有数据传输中间的网络设备可能会主动断开。OKX 的 WebSocket 要求客户端定期发送心跳否则服务器会在 30 秒后断开连接。心跳的发送方式很简单发送字符串ping就行def on_open(ws): def run(): while True: time.sleep(25) ws.send(ping) thread.start_new_thread(run, ())这里用 25 秒而不是 30 秒是留了 5 秒的余量。网络抖动的时候25 秒的间隔能保证心跳不会迟到。服务器收到 ping 后会返回 pong如果你在on_message里看到pong说明心跳正常。注意不要用 WebSocket 协议层面的 ping 帧OKX 要求的是应用层的文本ping。用错了服务器不认照样断线。4. 实操过程与核心环节实现4.1 环境准备与依赖安装先把 Python 环境搞定。如果你还没装 Python去官网下载 3.10 以上的版本。安装的时候记得勾选“Add Python to PATH”否则后面在命令行里敲python会提示找不到命令。装完之后验证一下python --version pip --version两个命令都能正常输出版本号说明环境没问题。然后安装依赖pip install requests websocket-client如果你用的是 PyCharm 或 VS Code也可以在 IDE 的终端里执行同样的命令。VS Code 的话记得在左下角选对 Python 解释器不然装到了别的环境里代码跑起来会报ModuleNotFoundError。4.2 REST 拉取 K 线数据的完整实现先写一个最简单的 REST 请求拉取 BTC-USDT 的 1 小时 K 线import requests def get_candles(instId, bar1H, limit100): url f{BASE_URL}/api/v5/market/candles params { instId: instId, bar: bar, limit: limit } resp requests.get(url, paramsparams, timeout10) if resp.status_code 200: data resp.json() if data[code] 0: return data[data] else: print(f接口返回错误: {data[msg]}) return None else: print(fHTTP 错误: {resp.status_code}) return None返回的数据是一个二维数组每个元素代表一根 K 线格式是[时间戳, 开盘价, 最高价, 最低价, 收盘价, 成交量, 成交额, 收盘时间]。注意时间戳是毫秒级的字符串用的时候要转成整数再转 datetime。import datetime def parse_candle(candle): ts int(candle[0]) / 1000 dt datetime.datetime.fromtimestamp(ts) return { time: dt, open: float(candle[1]), high: float(candle[2]), low: float(candle[3]), close: float(candle[4]), volume: float(candle[5]) }这里有个细节OKX 返回的 K 线数据是按时间倒序排列的最新的在最前面。如果你要画图或者做回测记得先反转一下顺序。4.3 WebSocket 实时行情订阅的完整实现WebSocket 的完整实现包含连接、订阅、消息处理、断线重连四个部分。先看连接和订阅import websocket import json import threading import time WS_URL wss://ws.okx.com:8443/ws/v5/public def on_open(ws): print(连接已建立) sub_msg { op: subscribe, args: [{channel: tickers, instId: BTC-USDT}] } ws.send(json.dumps(sub_msg)) # 启动心跳线程 def heartbeat(): while True: time.sleep(25) try: ws.send(ping) except: break threading.Thread(targetheartbeat, daemonTrue).start() def on_message(ws, message): if message pong: return data json.loads(message) if event in data: print(f事件: {data[event]}, 频道: {data.get(arg, {}).get(channel, N/A)}) elif arg in data: ticker data[data][0] print(fBTC-USDT 最新价: {ticker[last]}) def on_error(ws, error): print(f错误: {error}) def on_close(ws, close_status_code, close_msg): print(连接已关闭) ws websocket.WebSocketApp( WS_URL, on_openon_open, on_messageon_message, on_erroron_error, on_closeon_close ) ws.run_forever()这段代码跑起来之后你会看到 BTC-USDT 的最新价不断刷新。tickers频道推送的数据里包含last最新成交价、bidPx买一价、askPx卖一价、vol24h24 小时成交量等字段。4.4 断线重连的实现WebSocket 连接不可能永远不断。网络波动、服务器维护、心跳超时都会导致断线。一个健壮的客户端必须能自动重连。def run_forever_with_reconnect(): while True: try: ws websocket.WebSocketApp( WS_URL, on_openon_open, on_messageon_message, on_erroron_error, on_closeon_close ) ws.run_forever() except Exception as e: print(f连接异常: {e}) print(5 秒后重连...) time.sleep(5)这个重连逻辑很简单run_forever()返回说明连接断了等 5 秒再建一个新的。实际用的时候可以加一个重连次数上限避免无限重试。另外重连之后要重新发送订阅消息因为新的连接是全新的会话之前的订阅不会自动恢复。提示重连间隔不要设得太短比如 1 秒。如果服务器正在维护频繁重连可能会被临时限制。5 秒是一个比较稳妥的值。4.5 数据存储与消费拿到数据之后通常有两种消费方式一种是直接打印或推送到前端另一种是存到本地数据库或文件里。对于原型验证我建议先用 CSV 存简单直观import csv def save_ticker_to_csv(ticker, filenametickers.csv): with open(filename, a, newline) as f: writer csv.writer(f) writer.writerow([ ticker[ts], ticker[instId], ticker[last], ticker[bidPx], ticker[askPx] ])如果要长期运行建议换成 SQLite 或者时序数据库。CSV 文件大了之后读写都会变慢而且并发写入容易出问题。5. 常见问题与排查技巧实录5.1 签名错误排查签名错误是 REST 私有接口最常见的报错返回码通常是50113或50111。排查顺序如下排查项检查方法常见问题时间戳对比本地时间和服务器时间偏差超过 30 秒请求方法确认是 GET 还是 POST方法大小写不一致请求路径确认包含/api/v5/前缀路径拼写错误请求体POST 请求的 body 必须是 JSON 字符串body 为空时也要传空字符串Secret Key确认没有多余空格复制时带了换行符我踩过最坑的一次是时间戳格式。OKX 要求的是2024-01-15T10:30:00.000Z我一开始用了2024-01-15 10:30:00怎么调都是签名错误。后来对着文档一个字符一个字符比对才发现问题。5.2 WebSocket 连接被断开WebSocket 断线的原因主要有三个心跳没发、订阅消息格式错误、网络问题。排查方法看日志里有没有pong返回没有说明心跳没发出去检查订阅消息的 JSON 格式args必须是数组channel和instId不能拼错用websocket test client工具先手动连一下确认服务器地址和端口没问题如果连接建立后马上就被断开大概率是订阅消息格式不对。OKX 的 WebSocket 对格式要求很严格多一个空格都可能被拒绝。5.3 数据推送延迟或缺失有时候你会发现 WebSocket 推送的数据比实际行情慢了几秒或者某些时间段的数据直接缺失。这种情况通常是网络问题不是代码问题。可以这样排查检查本地网络到 OKX 服务器的延迟用ping或traceroute确认没有开代理代理会增加额外的网络跳转如果延迟持续超过 5 秒考虑换一个网络环境另外books5频道推送的是 5 档深度数据量比tickers大很多。如果你同时订阅了多个交易对的深度数据消息处理可能会跟不上导致消息堆积。解决办法是在on_message里只做最轻量的解析把耗时的操作放到单独的线程或队列里处理。5.4 限频触发与应对REST 接口触发限频后返回码是50011提示Rate limit reached。OKX 的限频规则是按接口和用户维度计算的公共接口的限频通常比较宽松但私有接口的限频更严格。应对策略把不紧急的请求放到队列里控制发送频率用 WebSocket 替代高频的 REST 轮询在代码里加一个简单的令牌桶限流器import time class RateLimiter: def __init__(self, max_calls, period): self.max_calls max_calls self.period period self.calls [] def acquire(self): now time.time() self.calls [t for t in self.calls if now - t self.period] if len(self.calls) self.max_calls: sleep_time self.period - (now - self.calls[0]) time.sleep(sleep_time) self.calls.append(time.time())这个限流器很简单但能有效避免因为突发请求导致的限频。5.5 常见问题速查表问题现象可能原因解决方法返回 401API Key 无效或权限不足检查 Key 是否正确权限是否勾选返回 50113签名错误检查时间戳、请求方法、路径、body返回 50011触发限频降低请求频率加限流器WebSocket 连不上地址或端口错误确认用的是wss://ws.okx.com:8443WebSocket 频繁断线心跳未发送每 25 秒发送一次ping数据解析报错字段缺失或格式变化加 try-except打印原始消息排查程序跑一段时间就崩内存泄漏或线程未回收检查心跳线程和重连逻辑6. 几个我实际踩过的坑和独家建议第一个坑是时间戳的时区问题。我一开始用datetime.now()生成时间戳结果签名一直失败。后来才发现 OKX 要求的是 UTC 时间不是本地时间。改成datetime.utcnow()之后问题解决。这个坑很隐蔽因为报错信息只说签名错误不会告诉你时间戳有问题。第二个坑是WebSocket 的订阅数量限制。OKX 对单个连接的订阅频道数量有限制具体数字文档里有写。我一开始想在一个连接里订阅几十个交易对的 ticker结果部分订阅失败。后来改成分批订阅或者用多个连接分担问题就解决了。第三个坑是消息处理的阻塞。on_message回调是在 WebSocket 的接收线程里执行的如果你在里面做了耗时的操作比如写数据库、发 HTTP 请求会阻塞后续消息的接收。我的做法是把消息丢到一个queue.Queue里然后用单独的消费者线程去处理。这样接收和处理解耦不会因为处理慢而丢消息。import queue import threading msg_queue queue.Queue() def on_message(ws, message): msg_queue.put(message) def consumer(): while True: msg msg_queue.get() # 在这里做耗时的处理 process_message(msg) threading.Thread(targetconsumer, daemonTrue).start()这个模式在实际项目里非常有用尤其是当你需要同时处理多个频道的数据时。最后一个建议日志一定要打全。我见过太多人调试的时候只打印“出错了”但不知道错在哪。我的做法是在关键节点都打日志连接建立、订阅发送、消息接收、异常抛出。日志级别用 INFO 打正常流程用 ERROR 打异常。这样出问题的时候翻日志就能定位到具体是哪一步出的错。import logging logging.basicConfig( levellogging.INFO, format%(asctime)s - %(levelname)s - %(message)s ) logger logging.getLogger(__name__) logger.info(WebSocket 连接已建立) logger.error(f签名失败: {resp.text})这套日志配置跑起来之后控制台输出清晰明了排查问题的时候能省很多时间。
返回列表