当工程师在行情数据的十字路口犹豫不决

你有过这种经历吗:凌晨两点,盯着屏幕上那个死活连不上的 WebSocket,手里的咖啡已经凉透,脑子里就一个问题——

我是不是该用 REST 把这段数据拉下来,而不是在这死等长连接?

这个问题的本质,不是 API 怎么用,而是你到底在解决什么问题

拉历史数据和接收实时行情,是两个完全不同的问题域。它们在数据量级、时效要求、连接特性和错误处理逻辑上有着根本性差异。TickDB 同时提供 REST 和 WebSocket 两套接口,正是为了对应这两类场景。如果你用 REST 去高频轮询实时行情,或者用 WebSocket 去拉十年的历史数据,感受到的"难用"不是工具的问题——是用错了工具。

本文从协议原理出发,拆解为什么 REST 和 WebSocket 在行情数据场景下不可互换,然后给出生产级的代码示例,最后附上 TickDB 的场景适配决策表。


一、问题的根源:两种截然不同的数据消费模式

在谈具体接口之前,必须先把底层逻辑理清楚。REST 和 WebSocket 的本质区别,不是"风格不同",而是数据消费模型的根本对立

1.1 REST:请求-响应,"我想看的时候再看"

REST 基于 HTTP,是一个同步的请求-响应模型。你发一个请求,服务器返回一个响应,连接关闭。整个过程是可预测的、无状态的、可缓存的。

这意味着什么?

  • 主动权在你。你需要什么时间窗口、什么数据精度,你说了算。服务端不关心你是谁、你来过没有。
  • 天然适合历史数据。历史数据是"已经发生的事实",不存在时效性。你今天查去年今天的收盘价,和明天查,结果是一样的。
  • HTTP 基础设施成熟。代理、缓存、CDN、日志审计,这些现成工具直接复用。

但反过来,REST 的优点在实时场景下变成了缺陷

# 如果用 REST 轮询实时行情——这不是"难用",这是架构性错误
import requests
import time

API_KEY = os.environ.get("TICKDB_API_KEY")

def poll_realtime_via_rest(symbol: str, interval: float = 0.5):
    """
    ⚠️ 这段代码演示了一个反面模式,请勿在生产环境使用。
    每 0.5 秒发一个请求,每小时产生 7200 次 API 调用。
    """
    count = 0
    while True:
        response = requests.get(
            f"https://api.tickdb.ai/v1/market/kline/latest",
            headers={"X-API-Key": API_KEY},
            params={"symbol": symbol, "interval": "1m"},
            timeout=(3.05, 10)
        )
        data = response.json()
        print(f"[{data.get('ts')}] {symbol}: {data.get('close')}")
        count += 1
        if count % 100 == 0:
            print(f"已发送 {count} 次请求,当前调用频率:{3600/interval:.0f}/小时")
        time.sleep(interval)

这段代码跑起来,五分钟之内你的 API 调用配额就会报警。更糟糕的是,0.5 秒的轮询间隔意味着你最多拿到 2Hz 的数据刷新率,而 1 分钟 K 线在理论上每一秒都可能更新——你用最重的工具,拿到最差的结果。

1.2 WebSocket:长连接推送,"有变化我就告诉你"

WebSocket 在 TCP 连接建立后,升级为全双工通信通道。服务端可以随时向客户端推送数据,不需要客户端反复发起请求。

这意味着什么?

  • 主动权在服务端。数据变化了,服务端主动推给你。你不需要问"现在价格是多少"。
  • 天然适合实时数据。行情是"持续流动的状态",每一笔成交、每一个订单簿更新,都是新的信息。你需要在它们发生的第一时间拿到。
  • 延迟极低。服务端推送 vs. REST 轮询,在时效性上差距不是一个量级。盘口每毫秒都在变,等你下一轮请求发出,前一个价格可能已经滑过三个档位了。

但 WebSocket 的优点在历史数据场景下,同样是缺陷

// ⚠️ 同样是一个反面模式演示
// WebSocket 不适合拉历史数据,原因如下:

// 1. 协议本身不支持范围查询
//    WebSocket 连接建立后,你只能接收服务端"决定发给你"的数据。
//    你无法指定"我要查 AAPL.US 2020-2024 年的 1 小时 K 线"。

// 2. 连接生命周期与数据请求不匹配
//    历史数据请求是一次性需求,数据返回完毕,请求就应该结束。
//    但 WebSocket 连接是为"持续监听"设计的,长期挂着只拉一次数据,浪费资源。

// 3. 断线重连逻辑在历史回放场景下是噩梦
//    如果传输中断,你要怎么"断点续传"?WebSocket 没有原生的请求上下文。

所以结论很清晰:不是哪个 API"更好用",而是它们根本就不是为同一个问题设计的。


二、用场景说话:四象限决策矩阵

光讲原理不够,我们直接进入决策场景。以下四种情况,你应该怎么选?

场景 问题本质 推荐接口 原因
查 AAPL.US 近十年 1 小时 K 线做回测 历史数据,一次性拉取 REST 可缓存、无时效要求、量级大
监控英伟达财报发布后的盘口变化 实时状态,需要第一时间感知 WebSocket 推送驱动、毫秒级延迟
在策略回测中获取某一天的开盘价 历史数据,但量小精准 REST 单次请求,REST 更轻量
实时计算买卖压力比,监控异常信号 实时流计算,需要连续数据 WebSocket depth 需要逐帧订单簿数据

核心判断逻辑只有两条:

第一条:这个数据是"已经写好的记录",还是"正在发生的事件"?

历史 K 线是写好的记录。REST 的请求-响应模型完美匹配这个语义——你去"读"一段历史,就像查一本书的某一页,翻完了就走。实时盘口是正在发生的事件。WebSocket 的推送模型才跟得上这个节奏——你坐在交易台前,等的就是那个"买一量突然增加"的瞬间。

第二条:你能接受多少延迟?

REST 轮询的最小延迟 = 你的轮询间隔。这个间隔设置太大,数据过时;设置太小,API 配额和服务器压力同时爆炸。

WebSocket 的延迟 = 网络往返时间,通常在几十毫秒量级。如果你的策略对延迟敏感(比如统计套利、盘口抢单),REST 根本不在考虑范围内。


三、生产级代码:两套接口的正确打开方式

3.1 REST:正确拉取历史数据的完整流程

用 REST 拉历史 K 线,看起来简单,但有几个坑必须避开:

import os
import time
import requests
from typing import List, Dict, Optional

# ============================================================
# TickDB REST API:历史 K 线拉取(生产级)
# ============================================================
# ⚠️ 关键点:
#   1. 用 /v1/market/kline(不是 /kline/latest)拉历史数据
#   2. 用 end_time 反向分页,避免数据漂移
#   3. 处理限频(code 3001)和品种不存在(code 2002)
#   4. 生产环境建议做本地缓存,避免重复拉取
# ============================================================

TICKDB_API_KEY = os.environ.get("TICKDB_API_KEY")
BASE_URL = "https://api.tickdb.ai/v1/market"

def fetch_kline_historical(
    symbol: str,
    interval: str = "1h",
    start_time: int,
    end_time: int,
    limit: int = 1000
) -> List[Dict]:
    """
    拉取指定时间范围的历史 K 线数据。
    
    重要:end_time 参数是毫秒级 Unix 时间戳。
    推荐从后往前拉(用 end_time 定位),然后用返回数据的
    最后一条时间戳作为下一页的 end_time,实现反向分页。
    
    这样可以避免在回测周期边界附近出现数据丢失。
    """
    all_klines = []
    current_end = end_time
    
    while True:
        headers = {"X-API-Key": TICKDB_API_KEY}
        params = {
            "symbol": symbol,
            "interval": interval,
            "start": start_time,
            "end": current_end,
            "limit": limit
        }
        
        response = requests.get(
            f"{BASE_URL}/kline",
            headers=headers,
            params=params,
            timeout=(3.05, 10)
        )
        
        data = response.json()
        code = data.get("code", 0)
        
        if code == 0:
            klines = data.get("data", [])
            if not klines:
                break
            
            # 反向遍历:从最旧的数据开始追加
            all_klines.extend(reversed(klines))
            
            # 更新 end_time 为当前页最旧数据的时间戳,继续往前拉
            # ⚠️ 这里取 klines[0](最早那条),避免重复
            current_end = klines[0]["ts"]
            
            print(f"已拉取 {len(klines)} 条,当前累计:{len(all_klines)} 条,"
                  f"最旧时间:{klines[0]['ts']}")
            
            if len(klines) < limit:
                break
        elif code == 3001:
            # 请求频率超限,读取 Retry-After 头
            retry_after = int(response.headers.get("Retry-After", 5))
            print(f"触发限频,等待 {retry_after} 秒...")
            time.sleep(retry_after)
        elif code == 2002:
            raise KeyError(f"交易品种 {symbol} 不存在,请检查代码")
        else:
            raise RuntimeError(f"API 错误 {code}: {data.get('message')}")
    
    # 排序并去重(分页边界可能有重叠数据)
    all_klines.sort(key=lambda x: x["ts"])
    seen = set()
    unique_klines = []
    for k in all_klines:
        if k["ts"] not in seen:
            seen.add(k["ts"])
            unique_klines.append(k)
    
    return unique_klines


# 使用示例:拉取苹果公司 2023 年全年 1 小时 K 线
if __name__ == "__main__":
    import datetime
    
    start = int(datetime.datetime(2023, 1, 1).timestamp() * 1000)
    end = int(datetime.datetime(2023, 12, 31, 23, 59).timestamp() * 1000)
    
    try:
        klines = fetch_kline_historical(
            symbol="AAPL.US",
            interval="1h",
            start_time=start,
            end_time=end
        )
        print(f"\n✅ 共获取 {len(klines)} 条 K 线数据")
        print(f"时间范围:{klines[0]['ts']} ~ {klines[-1]['ts']}")
    except Exception as e:
        print(f"❌ 获取失败:{e}")

3.2 WebSocket:正确接收实时行情的完整流程

实时行情的 WebSocket 实现,远比 REST 复杂。它需要处理连接保活、断线重连、消息解析、限频等待等多个状态。

import os
import json
import time
import random
import threading
import websocket
from typing import Callable, Optional, Dict, List

# ============================================================
# TickDB WebSocket:实时行情订阅(生产级)
# ============================================================
# ⚠️ 关键点:
#   1. WebSocket 鉴权通过 URL 参数传递 api_key
#   2. 必须处理 ping/pong 心跳保活,否则连接会被服务端超时断开
#   3. 断线后用指数退避 + 抖动重连,避免惊群效应
#   4. code 3001 错误码需要从消息体解析,并按 retry_after 等待
#   5. 高频场景建议改用 aiohttp/asyncio 架构
# ============================================================

TICKDB_API_KEY = os.environ.get("TICKDB_API_KEY")
WS_BASE_URL = "wss://api.tickdb.ai/v1/market/ws"

class TickDBWebSocket:
    """TickDB WebSocket 客户端封装(生产级)"""
    
    def __init__(
        self,
        api_key: str,
        on_message: Optional[Callable[[Dict], None]] = None,
        on_error: Optional[Callable[[str], None]] = None
    ):
        self.api_key = api_key
        self.on_message = on_message or self._default_handler
        self.on_error = on_error or (lambda e: print(f"[ERROR] {e}"))
        self.ws: Optional[websocket.WebSocketApp] = None
        self._running = False
        self._reconnect_delay = 1
        self._max_reconnect_delay = 60
        self._thread: Optional[threading.Thread] = None
    
    def _default_handler(self, data: Dict):
        """默认消息处理器,子类可重写"""
        print(f"[TICK] {data.get('s')}: close={data.get('c')}, "
              f"vol={data.get('v')}, depth_bid={data.get('bid_vol')}")
    
    def connect(self):
        """启动 WebSocket 连接(在独立线程中运行)"""
        self._running = True
        self._thread = threading.Thread(target=self._run, daemon=True)
        self._thread.start()
        print(f"[WS] 启动连接线程,URL: {WS_BASE_URL}")
    
    def _run(self):
        """WebSocket 主循环"""
        while self._running:
            try:
                # ⚠️ 鉴权通过 URL 参数传递(不是 Header)
                ws_url = f"{WS_BASE_URL}?api_key={self.api_key}"
                
                self.ws = websocket.WebSocketApp(
                    ws_url,
                    on_message=self._on_ws_message,
                    on_error=self._on_ws_error,
                    on_close=self._on_ws_close,
                    on_open=self._on_ws_open
                )
                
                # run_forever 的 ping_interval 启用自动心跳
                self.ws.run_forever(
                    ping_interval=30,      # 每 30 秒发送一次 ping
                    ping_timeout=10,       # 等待 pong 的超时时间
                    reconnect=False        # 手动控制重连逻辑
                )
                
            except Exception as e:
                self.on_error(f"WebSocket 连接异常:{e}")
            
            if self._running:
                self._schedule_reconnect()
    
    def _on_ws_open(self, ws):
        """连接建立后的回调"""
        print("[WS] 连接已建立,开始订阅...")
        self._reconnect_delay = 1  # 重置退避时间
        
        # 订阅行情数据(示例:苹果股票 + 比特币)
        subscribe_msg = {
            "cmd": "sub",
            "params": ["AAPL.US:ticker", "BTC.US:kline-1m", "AAPL.US:depth-5"]
        }
        ws.send(json.dumps(subscribe_msg))
        print(f"[WS] 已发送订阅请求:{subscribe_msg}")
    
    def _on_ws_message(self, ws, raw_message: str):
        """消息处理入口"""
        try:
            msg = json.loads(raw_message)
            
            # ---- 1. 心跳响应 ----
            if msg.get("cmd") == "pong":
                print(f"[WS] 心跳响应正常")
                return
            
            # ---- 2. 限频错误码处理 ----
            if msg.get("code") == 3001:
                retry_after = int(msg.get("retry_after", 5))
                print(f"[WS] 触发限频(3001),等待 {retry_after} 秒...")
                time.sleep(retry_after)
                return
            
            # ---- 3. 订阅确认 ----
            if msg.get("cmd") == "sub_ack":
                print(f"[WS] 订阅确认:{msg}")
                return
            
            # ---- 4. 正常数据消息 ----
            if "s" in msg:  # 标准行情消息格式,含交易品种字段
                self.on_message(msg)
            
        except json.JSONDecodeError:
            self.on_error(f"消息解析失败:{raw_message}")
        except Exception as e:
            self.on_error(f"消息处理异常:{e}")
    
    def _on_ws_error(self, ws, error):
        self.on_error(f"WebSocket 错误:{error}")
    
    def _on_ws_close(self, ws, close_status_code, close_msg):
        print(f"[WS] 连接关闭(code={close_status_code},msg={close_msg})")
    
    def _schedule_reconnect(self):
        """指数退避 + 抖动重连"""
        # 指数退避:每次重连等待时间翻倍
        delay = min(
            self._reconnect_delay * (2 ** random.randint(0, 1)),
            self._max_reconnect_delay
        )
        # 抖动:避免大量客户端同时重连造成服务端压力
        jitter = random.uniform(0, delay * 0.1)
        total_delay = delay + jitter
        
        print(f"[WS] {total_delay:.1f} 秒后准备重连(当前退避时间:{delay:.1f}s)")
        time.sleep(total_delay)
        
        self._reconnect_delay = min(self._reconnect_delay * 2, self._max_reconnect_delay)
    
    def close(self):
        """主动关闭连接"""
        self._running = False
        if self.ws:
            self.ws.close()
        print("[WS] 连接已主动关闭")


# 使用示例
if __name__ == "__main__":
    import datetime
    
    def handle_tick(data: Dict):
        """自定义消息处理函数"""
        symbol = data.get("s", "UNKNOWN")
        close = data.get("c", 0)
        ts = data.get("ts", 0)
        bid_vol = data.get("bid_vol", 0)
        ask_vol = data.get("ask_vol", 0)
        
        # 计算买卖压力比(压力比 > 1 表示买盘更积极)
        if ask_vol > 0:
            pressure_ratio = bid_vol / ask_vol
            signal = "🟢 买压主导" if pressure_ratio > 1.2 else \
                     "🔴 卖压主导" if pressure_ratio < 0.8 else "⚪ 均衡"
        else:
            pressure_ratio = 0
            signal = "❓ 无卖盘数据"
        
        dt = datetime.datetime.fromtimestamp(ts / 1000).strftime("%H:%M:%S")
        print(f"[{dt}] {symbol}: ${close} | 压力比: {pressure_ratio:.2f} {signal}")
    
    client = TickDBWebSocket(
        api_key=TICKDB_API_KEY,
        on_message=handle_tick
    )
    
    try:
        client.connect()
        time.sleep(60)  # 监控 60 秒后退出
    except KeyboardInterrupt:
        print("\n收到中断信号,正在关闭...")
    finally:
        client.close()

四、TickDB 的场景适配决策表

有了上面的代码基础,我们把场景决策做成一张可以直接用的对照表。

4.1 接口能力总览

能力维度 REST(/v1/market/kline) WebSocket(/v1/market/ws)
典型用途 历史数据拉取、回测数据准备 实时行情监控、信号触发
数据方向 客户端主动拉取(Pull) 服务端主动推送(Push)
连接特性 无状态,每次请求独立 持久连接,长连接复用
时效性 适合无时效要求的历史数据 适合毫秒级实时数据
数据量级 适合大批量数据(千条/次) 适合小批量高频更新
鉴权方式 Header: X-API-Key URL 参数: ?api_key=
分页支持 支持(limit + start/end 时间戳) 不支持(流式推送)
缓存支持 HTTP 层可缓存 不支持(长连接无法缓存)
重试语义 幂等,可安全重试 非幂等,消息可能有重复
适用 TickDB 接口 /kline/kline/latest/symbols/available sub 订阅 ticker、kline-Nm、depth-N

4.2 实战场景选型

你的目标 推荐方案 说明
做十年美股策略回测 REST /kline 分页拉取,稳定可靠,可本地缓存
监控财报发布瞬间的盘口 WebSocket depth 实时推送订单簿变化,捕捉流动性真空
获取某个期货品种的当前价格 REST /kline/latest 单次请求,轻量快速
做高频统计套利,需要逐笔成交 WebSocket trades(港股/数字货币) 每笔成交实时推送,延迟 < 100ms
盘中实时计算持仓盯市 WebSocket ticker 持续推送最新价、成交量、买卖盘量
查询某数字货币的历史波动率 REST /kline 先拉历史数据本地计算,不走 WebSocket

五、一个真实场景:REST + WebSocket 协同作战

说了这么多分开用,但你实际项目里大概率是两个接口配合着来

举一个完整的例子:做一个基于成交量加权的均值回归策略。

流程拆解:
1. 回测阶段:用 REST 拉 3 年历史 K 线,跑策略逻辑,验证参数
2. 实盘阶段:用 WebSocket 订阅 ticker,实时计算当前持仓的浮盈浮亏
3. 信号触发时:用 REST 查询最新的 `depth` 快照,评估流动性深度再下单

这不是假设的流程,这是 TickDB 用户最常见的使用模式。REST 负责"大量历史数据的初始化",WebSocket 负责"实盘中的持续感知"。两者各有分工,没有谁替代谁。


六、几个高频疑问

Q:REST 轮询能不能替代 WebSocket 做实时监控?

A:可以跑,但代价是:用最高的 API 消耗,换最差的数据时效。轮询间隔设 1 秒,意味着你最多每秒更新一次,而 WebSocket 可以实时推送每一次变化。在行情剧烈波动的时刻,这 1 秒的差距可能就是滑点损失。更别说 API 配额消耗——每 0.5 秒轮询一次,一个交易时段(6.5 小时)就是 46,800 次请求,任何 API 服务商都会对你的账号做特殊关照。

Q:WebSocket 能不能拉历史数据?

A:不能,也不应该。WebSocket 连接是为"持续监听"设计的会话通道。你无法在连接建立后告诉服务端"我要 2020 年到 2024 年的数据",然后期望它一次性推完 10 万条记录。协议层面不支持这种语义,强行实现也是对架构的扭曲。

Q:WebSocket 断线了怎么办?

A:代码中已经实现了指数退避 + 抖动重连。但还有一个更根本的问题:你在断线期间丢失的数据怎么办?

对于 WebSocket 订阅的实时数据,TickDB 不提供数据补推。这意味着如果你监控的是"财报发布瞬间的流动性突变",断线 5 秒可能就错过了最关键的窗口。对于这类场景,正确做法是:先用 REST /kline/latest 查询当前最新状态建立基线,再用 WebSocket 监控后续变化。这样断线重连后,你可以立刻用 REST 拿到当前最新值作为起点,而不是从零开始。


结语

回到文章开头那个凌晨两点的场景。工程师盯着断开的 WebSocket,手足无措——这是用错了工具吗?不一定。更可能的情况是:他在用 REST 思维理解 WebSocket 问题,或者反之。

REST 和 WebSocket 的取舍,本质上是对"数据消费模式"的选择题。你是在读一本已经写完的书(历史数据),还是在现场观看一场直播(实时行情)?选对接口,性能问题自然消失。


下一步行动

如果你需要拉历史数据做策略回测

  1. 访问 tickdb.ai 注册(免费,无需信用卡)
  2. 使用 REST /v1/market/kline 接口,配合本文的分页代码
  3. 免费层级已覆盖个人量化研究需求

如果你需要毫秒级实时行情监控

  1. 使用 WebSocket sub 命令订阅 ticker 或 depth 频道
  2. 参考本文的 TickDBWebSocket 类实现生产级重连逻辑
  3. 高频场景(> 10 个标的)建议使用 aiohttp 异步架构

如果你在设计自己的行情 API 系统
本文的决策框架同样适用:区分"历史记录查询"和"实时状态感知",分别选择 Pull 和 Push 模型。不要用轮询替代长连接,那是用更贵的工具做更差的工作。

如果你习惯用 AI 辅助开发,在 AI 助手中搜索安装 tickdb-market-data SKILL,可以直接用自然语言查询 TickDB 的接口能力。


本文不构成任何投资建议。市场有风险,投资需谨慎。