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

如何将AngelBroking股票API实时数据转为可用格式以计算技术指标

实现思路

你可以通过维护一个全局可更新的DataFrame,结合周期K线聚合逻辑,把实时tick数据转换成可用于TA-Lib计算的格式,具体步骤如下:

  • 首先确认AngelBroking WebSocket返回的tick数据结构,通常message为字典类型,核心字段包含:token(标的标识)、ltp(最新成交价)、open、high、low、close、timestamp(时间戳)
  • 定义全局的实时K线DataFrame,预设列包含timestamp、open、high、low、close、volume等你需要的字段
  • 按需设置K线周期(比如1分钟、5分钟),在on_message回调中收到新tick时,先判断是否属于新的周期:
    • 如果是新周期,就往DataFrame插入一行新的K线记录,OHLC初始值都设为当前tick的成交价
    • 如果属于当前周期,就更新当前行的high(取最大值)、low(取最小值)、close(设为最新成交价)
  • 每次更新完DataFrame后,就可以直接调用TA-Lib的方法计算技术指标

修改后可直接运行的代码示例

from smartapi import SmartWebSocket
import pandas as pd
import datetime
import talib

# 配置参数
FEED_TOKEN="YOUR_FEED_TOKEN"
CLIENT_CODE="YOUR_CLIENT_CODE"
token="nse_cm|2885" # 替换成你要订阅的标的
task="mw"

# 全局变量定义
KLINE_PERIOD = '1T' # 1分钟K线,可改5T、15T、1H等
# 初始化K线DataFrame
kline_df = pd.DataFrame(columns=['timestamp', 'open', 'high', 'low', 'close'])
current_period_start = None

ss = SmartWebSocket(FEED_TOKEN, CLIENT_CODE)

def on_message(ws, message):
    global kline_df, current_period_start
    # 解析tick数据,以实际返回的字段为准,以下为通用示例
    tick = message[0] # 单标的订阅时取第一个元素,多标的需按token过滤
    ltp = float(tick['ltp'])
    tick_time = datetime.datetime.fromtimestamp(tick['exchange_timestamp']/1000)
    
    # 计算当前tick所属的周期起始时间
    period_start = tick_time.floor(KLINE_PERIOD)
    
    if period_start != current_period_start:
        # 新周期,插入新K线
        new_row = pd.Series({
            'timestamp': period_start,
            'open': ltp,
            'high': ltp,
            'low': ltp,
            'close': ltp
        })
        kline_df = pd.concat([kline_df, new_row.to_frame().T], ignore_index=True)
        current_period_start = period_start
    else:
        # 现有周期,更新OHLC
        current_idx = kline_df.index[-1]
        kline_df.loc[current_idx, 'high'] = max(kline_df.loc[current_idx, 'high'], ltp)
        kline_df.loc[current_idx, 'low'] = min(kline_df.loc[current_idx, 'low'], ltp)
        kline_df.loc[current_idx, 'close'] = ltp
    
    # 数据量足够时计算指标,比如MA需要至少20条数据
    if len(kline_df) >= 20:
        kline_df['ma20'] = talib.MA(kline_df['close'], timeperiod=20)
        # 可继续添加其他指标计算逻辑
        print("当前K线数据及指标:\n", kline_df.tail(5))

def on_open(ws):
    print("连接成功,开始订阅行情")
    ss.subscribe(task, token)

def on_error(ws, error):
    print("报错:", error)

def on_close(ws):
    print("连接关闭")
    # 断开连接时可选择保存K线数据到本地
    kline_df.to_csv("realtime_kline.csv", index=False)

# 绑定回调
ss._on_open = on_open
ss._on_message = on_message
ss._on_error = on_error
ss._on_close = on_close

ss.connect()

注意事项

  • 以上代码中tick字段的取值需要根据你实际收到的message结构调整,你可以先打印一次完整的message看清楚层级和字段名再修改
  • 如果订阅多个标的,可以给DataFrame加token列做区分,收到tick时按token过滤后更新对应标的的K线
  • 如果对数据实时性要求高,也可以跳过K线聚合,直接用最新的close值和历史K线拼接后计算指标,不需要等周期走完

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 16:27:03