You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用finnhub.io Websocket API获取实时BTC价格时如何处理迟到数据

流式分钟级VWAP计算与迟到数据处理方案

问题背景

  • 基于finnhub.io Websocket API拉取BINANCE:BTCUSDT实时交易流,按自然分钟周期统计成交量加权平均价(VWAP)
  • 原有实现逻辑缺陷:仅通过判断单条交易时间戳秒数是否为00触发VWAP输出,存在两个核心问题:
    • 整分钟时间点对应多条交易记录,会导致VWAP输出逻辑被重复触发、统计值被反复重置
    • 无法识别网络传输导致的迟到交易数据,跨周期数据会混入统计,导致结果不准
  • 核心要求:对超过统计截止时间才到达的非当前周期数据做丢弃处理,避免重复计算、跨周期数据污染。

原有缺陷代码

import json
import websocket
from datetime import datetime

TOTAL_PRICE = 0
TOTAL_VOLUME = 0


def edit_message(message): 
    for line in json.loads(message)['data']:
        price = line['p']
        volume = line['v']
        calculate_vwap(price, volume)

        print(datetime.utcfromtimestamp(line['t'] / 1000).strftime('%Y-%m-%d %H:%M:%S'),
              "price:", line['p'],
              "volume:", line['v'])


def calculate_vwap(p, v):
    global TOTAL_PRICE
    global TOTAL_VOLUME

    TOTAL_PRICE += p
    TOTAL_VOLUME += v


def print_vwap():
    global TOTAL_PRICE
    global TOTAL_VOLUME

    print("Volume-weighted average price:", (TOTAL_PRICE * TOTAL_VOLUME))
    TOTAL_PRICE = 0
    TOTAL_VOLUME = 0



def on_message(ws, message):
    edit_message(message)

    for data in json.loads(message)['data']:
        time = datetime.utcfromtimestamp(data['t'] / 1000).strftime('%Y-%m-%d %H:%M:%S')
        secs = str(time[-2:])

        if int(secs) == 00: 
            print_vwap()


def on_close(ws):
    print("### closed ###")


def on_open(ws):
    ws.send('{"type":"subscribe","symbol":"BINANCE:BTCUSDT"}')


if __name__ == "__main__":
    websocket.enableTrace(False)
    ws = websocket.WebSocketApp("wss://ws.finnhub.io?token=cahkkgiad3i7auh49eb0",
                                on_message=on_message,
                                on_close=on_close)
    ws.on_open = on_open
    ws.run_forever()

*注:原代码除了触发逻辑问题,VWAP计算公式也存在错误,正确公式为周期总成交额/周期总成交量,修正版中会一并修复。

实现思路

采用流计算标准的固定时间窗口+水位线宽限期方案,从根源解决重复触发和迟到数据问题:

  • 窗口划分:每条交易按自身携带的交易时间戳(而非数据到达本地的时间)归属到对应自然分钟窗口,窗口用整分钟级时间戳作为唯一标识
  • 状态存储:弃用全局单值统计变量,改用字典独立存储每个窗口的累计成交额、累计成交量,避免跨窗口数据干扰
  • 触发逻辑:不在消息接收回调中判断输出时机,单独启动后台巡检线程,每秒检查一次未输出的窗口,当窗口结束时间超出预设宽限期后,统一输出该窗口的VWAP结果,保证每个窗口仅输出一次
  • 迟到数据处理:窗口输出完成后加入已完成集合,后续所有归属到该窗口的迟到数据直接丢弃,不进入统计流程
  • 宽限期可根据实际观测到的数据源延迟调整:加密货币行情场景下设置3-5秒即可覆盖绝大多数网络延迟,在统计准确性和输出实时性之间做平衡。

修正后可运行代码

import json
import websocket
import threading
import time
from datetime import datetime

# 基础配置
SYMBOL = "BINANCE:BTCUSDT"
WS_ADDR = "wss://ws.finnhub.io?token=cahkkgiad3i7auh49eb0"
WINDOW_SECONDS = 60  # 统计周期:1分钟
WATERMARK_DELAY = 5  # 窗口宽限期:窗口结束后等待5秒再输出,容忍迟到数据

# 统计状态存储
window_stats = dict()  # 格式:{窗口起始时间戳: [累计成交额, 累计成交量]}
completed_windows = set()  # 已完成输出的窗口集合,用于过滤迟到数据
state_lock = threading.Lock()  # 线程锁,避免多线程同时修改状态出错


def get_window_key(trade_ts_ms):
    """将毫秒级交易时间戳转换为所属分钟窗口的起始时间戳(秒级,整分钟)"""
    trade_ts_sec = trade_ts_ms / 1000
    return int(trade_ts_sec // WINDOW_SECONDS * WINDOW_SECONDS)


def process_single_trade(price, volume, trade_ts_ms):
    """处理单条交易数据,归入对应窗口"""
    window_key = get_window_key(trade_ts_ms)
    # 归属窗口已完成输出,直接丢弃迟到数据
    if window_key in completed_windows:
        return
    with state_lock:
        if window_key not in window_stats:
            window_stats[window_key] = [0.0, 0.0]
        # 累计成交额 = 价格*成交量 求和
        window_stats[window_key][0] += price * volume
        # 累计成交量直接求和
        window_stats[window_key][1] += volume


def window_inspect_worker():
    """后台巡检线程,定期检查窗口是否到期,到期则计算输出VWAP"""
    while True:
        current_ts = time.time()
        with state_lock:
            # 遍历所有未完成的窗口
            for w_key in list(window_stats.keys()):
                w_end_ts = w_key + WINDOW_SECONDS
                # 达到水位线要求,触发输出
                if current_ts >= w_end_ts + WATERMARK_DELAY:
                    total_turnover, total_vol = window_stats.pop(w_key)
                    completed_windows.add(w_key)
                    # 计算VWAP,兼容成交量为0的异常场景
                    vwap = total_turnover / total_vol if total_vol > 0 else 0
                    w_time_str = datetime.utcfromtimestamp(w_key).strftime('%Y-%m-%d %H:%M')
                    print(f"\n===== {w_time_str} 分钟周期统计结果 =====")
                    print(f"周期累计成交量: {total_vol:.4f} BTC")
                    print(f"成交量加权平均价(VWAP): {vwap:.2f} USDT\n")
        time.sleep(1)


def on_message(ws, message):
    msg = json.loads(message)
    if "data" not in msg:
        return
    for trade in msg["data"]:
        p = trade["p"]
        v = trade["v"]
        ts = trade["t"]
        process_single_trade(p, v, ts)
        # 打印实时交易明细
        trade_time = datetime.utcfromtimestamp(ts/1000).strftime('%Y-%m-%d %H:%M:%S')
        print(f"[{trade_time}] 成交价: {p} USDT | 成交量: {v} BTC")


def on_close(ws, code, reason):
    print("### Websocket连接已关闭 ###")


def on_open(ws):
    ws.send(json.dumps({"type": "subscribe", "symbol": SYMBOL}))
    print(f"已成功订阅 {SYMBOL} 实时交易数据,开始统计分钟级VWAP...")


if __name__ == "__main__":
    websocket.enableTrace(False)
    # 启动窗口巡检后台线程
    inspect_thread = threading.Thread(target=window_inspect_worker, daemon=True)
    inspect_thread.start()
    # 建立Websocket连接
    ws_client = websocket.WebSocketApp(
        WS_ADDR,
        on_message=on_message,
        on_close=on_close
    )
    ws_client.on_open = on_open
    ws_client.run_forever()

内容的提问来源于stack exchange,提问作者Kristiyan Kotomanov

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 09:18:17