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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 21:30:56