如何将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
相关产品推荐
相关产品推荐

