Python流式API输出的字典数据如何追加保存到CSV文件
TD流式股票行情数据追加写入CSV实现方案
核心逻辑是给流式客户端绑定消息回调函数,每次服务端推送数据时自动触发解析、追加写入操作,不需要手动轮询终端输出。
前置准备
首先明确返回数据中数字键对应的行情字段含义,方便后续导出可读的CSV列:
1:对应K线开盘价2:对应K线最高价3:对应K线最低价4:对应K线收盘价5:对应单根K线成交量6:对应单根K线成交笔数7:对应K线生成的毫秒级时间戳8:对应当日累计成交量key:对应标的代码seq:对应消息序列号
完整实现代码
先导入依赖模块,保留你原有初始化逻辑,新增回调处理、CSV写入逻辑即可:
import csv import os # 保留你原有TD客户端相关的导入、初始化代码 # -------------------------- # 新增配置项 # -------------------------- # 数字键到CSV列名的映射 FIELD_MAPPING = { "1": "open", "2": "high", "3": "low", "4": "close", "5": "volume", "6": "trade_count", "7": "bar_timestamp", "8": "cumulative_volume", "key": "symbol", "seq": "message_seq" } # CSV文件存储路径 CSV_SAVE_PATH = "qqq_kline_stream.csv" # CSV完整表头,额外加一列服务端响应时间戳 CSV_COLUMNS = list(FIELD_MAPPING.values()) + ["server_response_timestamp"] # 首次运行时初始化CSV,写入表头 if not os.path.exists(CSV_SAVE_PATH): with open(CSV_SAVE_PATH, "w", newline="", encoding="utf-8") as csv_file: writer = csv.DictWriter(csv_file, fieldnames=CSV_COLUMNS) writer.writeheader() # -------------------------- # 新增消息处理回调 # -------------------------- def on_stream_message_received(raw_message): """每次收到流式推送自动执行的处理逻辑""" # 遍历推送里的所有数据块 for data_block in raw_message.get("data", []): # 只处理股票K线服务的推送,跳过心跳、订阅成功通知等其他消息 if data_block.get("service") != "CHART_EQUITY": continue # 如果不需要存订阅成功的SUBS返回,放开下面这行注释即可 # if data_block.get("command") == "SUBS": # continue server_ts = data_block.get("timestamp") # 遍历单块数据里的所有行情条目 for content_item in data_block.get("content", []): # 转换数字键为可读列名 row_data = {} for num_key, col_name in FIELD_MAPPING.items(): row_data[col_name] = content_item.get(num_key) row_data["server_response_timestamp"] = server_ts # 追加模式写入CSV with open(CSV_SAVE_PATH, "a", newline="", encoding="utf-8") as csv_file: writer = csv.DictWriter(csv_file, fieldnames=CSV_COLUMNS) writer.writerow(row_data) # -------------------------- # 修改原有流式启动逻辑,注册回调 # -------------------------- streaming_api_service = td_client.streaming_api_client() # 绑定消息处理回调 streaming_api_service.add_event_handler("message", on_stream_message_received) streaming_services.chart( service=ChartServices.ChartEquity, symbols=['QQQ'], fields=ChartEquity.All ) streaming_api_service.open_stream()
注意事项
- 打开CSV文件时必须带
newline=""参数,否则Windows环境下导出的文件会出现行间空行 - 你的场景下数据更新频率为1分钟1条,上述每次收到消息再打开文件写入的逻辑完全够用,还能避免程序异常退出时文件损坏、内存数据丢失;如果是高频推送场景,可以把文件对象设为全局常驻减少IO开销
- 后续如果新增订阅其他股票代码,不需要修改写入逻辑,
symbol列会自动记录对应标的代码 - 时间戳字段是毫秒级Unix时间,后续做数据分析时可以直接转成标准时间格式
内容的提问来源于stack exchange,提问作者user16950345
相关产品推荐
相关产品推荐

