TradingAgents-CN 历史数据存储优化:基于 MongoDB `stock_daily_quotes` 的三数据源统一方案
TradingAgents-CN 历史数据存储优化基于 MongoDBstock_daily_quotes的三数据源统一方案【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN本文面向需要在 TradingAgents-CN 中管理多数据源Tushare / AKShare / BaoStock历史 K 线数据的开发者系统讲解该仓库围绕stock_daily_quotes集合实施的统一存储方案从分散存储、格式不统一、同步实现不完整的旧问题出发到统一数据模型、索引设计、HistoricalDataService服务层、RESTful API、三数据源接入与部署验证的完整落地路径。阅读后你将掌握该方案的集合结构与索引策略、服务层核心方法及调用链并可直接在本地部署验证。一、背景原有历史数据存储的三类问题在设计文档与历史代码中历史数据存储主要存在三类问题1. 存储分散历史数据散落在stock_data、market_quotes等多个集合中且数据格式不统一——部分以 JSON 字符串形式存放部分为结构化文档导致查询复杂、性能低下也无法对数据进行跨数据源对比。2. 实现不完整Tushare 同步服务历史数据保存功能未实现AKShare 同步服务仅有 TODO 注释无实际保存逻辑BaoStock 同步服务只保存元信息不保存实际 K 线数据。3. 设计与实现脱节设计文档中早已定义stock_daily_quotes集合但实际代码从未使用该集合且缺乏统一的数据管理接口。这三类问题共同导致历史数据无法可靠落地、无法统一查询、无法跨数据源验证质量。优化目标因此非常明确实现三数据源Tushare / AKShare / BaoStock的统一、高效、可靠的历史数据管理即统一集合、统一数据模型、统一服务接口、统一 API。二、统一数据模型stock_daily_quotes集合优化的第一步是创建专门的stock_daily_quotes集合以结构化文档统一承载历史 K 线数据。原优化方案中定义的文档结构如下{ _id: ObjectId(...), symbol: 000001, // 股票代码 full_symbol: 000001.SZ, // 完整代码 market: CN, // 市场类型 trade_date: 2024-01-16, // 交易日期 open: 12.60, // 开盘价 high: 12.85, // 最高价 low: 12.45, // 最低价 close: 12.75, // 收盘价 pre_close: 12.55, // 前收盘价 change: 0.20, // 涨跌额 pct_chg: 1.59, // 涨跌幅 volume: 130000000, // 成交量 amount: 1650000000, // 成交额 data_source: tushare, // 数据源标识 created_at: ISODate(...), // 创建时间 updated_at: ISODate(...), // 更新时间 version: 1 // 版本号 }在仓库实际实现中该模型还经历了进一步演进见 historical_data_service.py 的_standardize_record方法新增period字段支持daily/weekly/monthly多周期数据周期成为记录主键的一部分新增code字段与symbol保持一致用于向后兼容可选指标字段turnover_rate、volume_ratio、pe、pb、ps、adjustflag、tradestatus、isST等字段在数据源提供时才会写入避免了空字段浪费存储change/pct_chg自动计算当close与pre_close均存在时服务层会自动计算涨跌额与涨跌幅并保留 4 位小数。集合初始化脚本 create_historical_data_collection.py 会同时插入一条data_source: example的示例数据平安银行 2024-01-15 的日线便于开发阶段验证查询链路。full_symbol 生成规则从源码 historical_data_service.py 可以确认full_symbol的生成逻辑市场规则示例CNA股6开头 →.SH0/3开头 →.SZ其余默认.SZ000001→000001.SZHK追加.HK00700→00700.HKUS直接使用原代码AAPL→AAPL这一规则使同一套模型可以承载 A 股、港股、美股的 K 线数据与仓库的多市场支持能力保持一致。三、索引设计从“10 个优化索引”到 12 个索引的演进原优化方案设计了 10 个优化索引核心如下// 1. 复合唯一索引防重复 {symbol: 1, trade_date: 1, data_source: 1} // 2. 常用查询索引 {symbol: 1} // 按股票查询 {trade_date: -1} // 按日期查询 {data_source: 1} // 按数据源查询 // 3. 复合查询索引 {symbol: 1, trade_date: -1} // 股票历史数据 {market: 1, trade_date: -1} // 市场数据 {data_source: 1, updated_at: -1} // 同步监控 // 4. 性能优化索引 {volume: -1} // 成交量排序稀疏索引 {updated_at: -1} // 数据维护在实际落地时由于引入了period周期维度唯一索引被升级为四字段复合唯一索引{symbol: 1, trade_date: 1, data_source: 1, period: 1}名为symbol_date_source_period_unique这保证了同一股票、同一交易日、同一数据源、同一周期的记录全局唯一从而支撑后续的upsert幂等写入。create_historical_data_collection.py 中实际创建的 12 个索引包括索引名称用途{symbol, trade_date, data_source, period}唯一symbol_date_source_period_uniqueupsert 防重复{symbol}symbol_index单股查询{trade_date: -1}trade_date_index按日期范围查询{data_source}data_source_index按数据源查询{symbol, trade_date: -1}symbol_date_index股票历史数据常用查询{market}market_index按市场查询{updated_at: -1}updated_at_index数据维护{market, trade_date: -1}market_date_index市场级别查询{data_source, updated_at: -1}source_updated_index数据同步监控{volume: -1}稀疏volume_index筛选活跃股票{period}period_index按周期查询{symbol, period, trade_date: -1}symbol_period_date_index股票周期查询其中volume索引使用sparse: true稀疏索引因为部分数据源记录可能缺少成交量字段。若你的环境中已存在旧版本索引可运行 update_historical_data_indexes.py 完成升级它会先删除旧的symbol_date_source_unique索引为存量数据批量补充period: daily字段再创建新的四字段唯一索引与周期相关索引。此外服务层启动时还会通过_ensure_indexes自动检查并创建 4 个核心索引唯一索引 symbol / trade_date / symbol_date索引创建失败不会阻止服务启动仅记录警告见 historical_data_service.py。四、统一数据管理服务HistoricalDataService仓库通过单例模式提供统一的HistoricalDataService见 historical_data_service.py模块级全局变量_historical_data_service配合get_historical_data_service()工厂函数实现懒加载初始化首次调用时建立数据库连接并确保索引存在后续复用同一实例。核心方法一览方法功能save_historical_data(symbol, data, data_source, market, period)批量保存历史 K 线数据返回保存条数get_historical_data(symbol, start_date, end_date, data_source, period, limit)多维度查询历史数据get_latest_date(symbol, data_source)获取某股票某数据源的最新数据日期get_data_statistics()数据量统计与质量监控保存链路向量化单位转换 批量 upsert 指数退避重试save_historical_data是整套方案的核心其实现包含三个值得借鉴的工程细节1. DataFrame 层面向量化单位转换在逐行标准化之前先在 DataFrame 层面对整列做向量化运算见 historical_data_service.pyTushare 数据成交额amount或turnover× 1000千元 → 元成交量volume或vol× 100手 → 股港股/美股数据若缺少pre_close通过close.shift(1)取前一交易日收盘价作为前收盘价。这避免了逐行转换的开销是批量写入性能的关键之一。2. 批量 upsert 幂等写入逐行调用_standardize_record标准化后构造ReplaceOne操作以{symbol, trade_date, data_source, period}为过滤条件upsertTrue实现“有则更新、无则插入”。每攒满batch_size条源码中为200 条比设计文档初稿的 500 条进一步调小以避免批量写入超时执行一次bulk_write(operations, orderedFalse)。orderedFalse允许 MongoDB 并行执行写入进一步提升吞吐。3. 超时重试指数退避_execute_bulk_write_with_retry对asyncio.TimeoutError及消息含timeout/timed out的异常进行最多5 次重试退避等待时间按3 ** retry_count递增3 秒、9 秒、27 秒、81 秒见 historical_data_service.py。这一机制显著提升了在大批量历史数据同步场景下对 MongoDB 偶发超时的容错能力。使用示例与设计文档一致# 获取服务实例懒加载单例 service await get_historical_data_service() # 保存历史数据 saved_count await service.save_historical_data( symbol000001, datadataframe, # pd.DataFrame含 open/high/low/close/volume/amount 等列 data_sourcetushare, # tushare / akshare / baostock marketCN # CN / HK / US ) # 查询历史数据按 trade_date 倒序返回 results await service.get_historical_data( symbol000001, start_date2024-01-01, end_date2024-01-31, data_sourcetushare ) # 获取统计信息 stats await service.get_data_statistics()get_data_statistics通过聚合管道返回总记录数、去重股票数、按数据源含各源最新日期、按市场的分布统计以及统计生成时间见 historical_data_service.py可直接用于数据量监控看板。五、三数据源同步服务接入统一服务层的价值最终体现在三个同步服务的接入上。当前仓库中三个 worker 均已接入get_historical_data_serviceTushare 同步服务tushare_sync_service.pyasync def _save_historical_data(self, symbol: str, df, period: str daily) - int: 保存历史数据到数据库 if self.historical_service is None: self.historical_service await get_historical_data_service() # 使用统一历史数据服务保存指定周期 saved_count await self.historical_service.save_historical_data( symbolsymbol, datadf, data_sourcetushare, marketCN, periodperiod, ) return saved_count在sync_historical_data主流程中API 拉取与数据保存分别计时api_duration/save_duration保存耗时被计入同步性能统计。AKShare 同步服务akshare_sync_service.pyif hist_data is not None and not hist_data.empty: # 保存到统一历史数据集合 if self.historical_service is None: self.historical_service await get_historical_data_service() saved_count await self.historical_service.save_historical_data( symbolsymbol, datahist_data, data_sourceakshare, marketCN, )AKShare 服务在批量处理流程中调用替代了原先的 TODO 注释占位。BaoStock 同步服务baostock_sync_service.py# 初始化历史数据服务 if self.historical_service is None: self.historical_service await get_historical_data_service() # 保存到统一历史数据集合 saved_count await self.historical_service.save_historical_data( symbolcode, datahist_data, data_sourcebaostock, marketCN, periodperiod, )BaoStock 服务除了保存到统一集合外还会同步更新兼容性元信息保证旧查询路径不受影响。三个服务均采用“延迟初始化”模式self.historical_service None在__init__中声明首次保存时才获取服务实例避免同步任务启动时不必要的数据库开销。六、RESTful API 接口历史数据功能已通过 historical_data.py 暴露为完整的 RESTful API路由前缀为/api/historical-data并在 app/main.py 中随include_router(historical_data.router, tags[historical-data])见 app/main.py集成到主应用。端点方法说明/api/historical-data/query/{symbol}GET查询单只股票历史数据/api/historical-data/queryPOSTPOST 方式查询支持请求体传参/api/historical-data/compare/{symbol}GET跨数据源数据对比需传trade_date/api/historical-data/statisticsGET全局数据统计信息/api/historical-data/latest-date/{symbol}GET获取最新数据日期需传data_source/api/historical-data/healthGET健康检查典型调用示例# 查询历史数据支持 start_date / end_date / data_source / period / limit 筛选 GET /api/historical-data/query/000001?start_date2024-01-01end_date2024-01-31 # 数据对比对比 tushare/akshare/baostock 三源同日数据 GET /api/historical-data/compare/000001?trade_date2024-01-16 # 统计信息 GET /api/historical-data/statistics # 最新日期 GET /api/historical-data/latest-date/000001?data_sourcetushare # 健康检查 GET /api/historical-data/healthcompare端点是本方案独有的质量验证能力它会依次查询tushare、akshare、baostock三个数据源在同一交易日的记录返回各源数据及实际可用的源列表见 historical_data.py。跨数据源对比得到的结果形如 数据对比结果: - tushare: 收盘价12.75, 成交量130000000, 涨跌幅1.59% - akshare: 收盘价12.73, 成交量128000000, 涨跌幅1.43% - baostock: 收盘价12.77, 成交量132000000, 涨跌幅1.75% 收盘价差异: 0.0400所有 GET 查询响应均统一包装为{success, message, data}结构其中data内含symbol、count、query_params回显与records列表limit参数在服务层约束为1 ≤ limit ≤ 1000。七、部署与验证1. 创建集合和索引python scripts/setup/create_historical_data_collection.py脚本连接app.core.config.settings中配置的MONGO_URI/MONGO_DB创建stock_daily_quotes集合、12 个索引并插入示例数据最后打印集合统计与完整索引列表。若需为存量数据升级周期字段与索引运行python scripts/setup/update_historical_data_indexes.py2. 启动服务历史数据 API 已集成到主应用直接启动即可python -m uvicorn app.main:app --reload3. 验证 APIcurl http://localhost:8000/api/historical-data/health curl http://localhost:8000/api/historical-data/statistics健康检查端点返回服务状态、总记录数、去重股票数与检查时间统计端点返回按数据源/按市场的分布明细可作为上线后的自检手段。4. 功能测试验收设计文档给出的验收测试覆盖 6 项核心能力输出示例 历史数据存储优化简单测试 ✅ MongoDB连接: 通过 ✅ 数据插入: 通过 ✅ 数据查询: 通过 ✅ 数据对比: 通过 ✅ 聚合查询: 通过 ✅ 数据清理: 通过 测试完成: 6/6 项测试通过 ✅ 所有测试通过历史数据存储优化成功若需要独立诊断线上同步问题可参考仓库中的 diagnose_historical_data_sync.py涉及stock_daily_quotes的字段与金额/成交量单位正确性可在 tests/test_amount_fix.py 中看到验证思路检查数据库stock_daily_quotes集合中的文档。八、监控指标优化方案为历史数据运维定义了三个维度的监控指标均可通过上述统计 API 与聚合查询落地数据量监控总记录数、各数据源记录数、各股票记录数、最新数据日期get_data_statistics已返回按数据源的最新日期可进一步按{data_source, updated_at}索引排序实现“同步监控”。性能监控查询响应时间、批量写入速度服务日志中记录了每批次写入耗时与总耗时、索引使用率、存储空间使用。质量监控数据完整性检查、跨数据源一致性对比compare端点、异常数据识别、数据更新及时性latest-date端点。九、总结本次历史数据存储优化在 TradingAgents-CN 中形成了一条完整的闭环统一存储创建专门的stock_daily_quotes集合结构化文档统一承载三数据源 K 线数据并扩展period/code/ 可选指标字段完善功能Tushare、AKShare、BaoStock 三个同步服务全部接入统一保存链路修复了原先“未实现 / 仅 TODO / 只存元信息”的问题提升性能唯一索引 复合查询索引 稀疏索引的组合设计配合向量化单位转换、批量 upsert、orderedFalse并行写入与指数退避重试支撑大规模历史数据的高效落地增强监控统计聚合、跨数据源对比、最新日期、健康检查四个 API 提供了完整的数据量与质量监控手段标准接口GET / POST 双模式查询 API 满足按股票、日期、数据源、周期的多维检索需求。需要注意的是方案文档中提及的“1000 条/批次”与“查询时间从秒级降至毫秒级、存储空间减少 50%”属于优化目标表述实际源码将批量大小设定为 200 条/批次以规避写入超时historical_data_service.py真实性能数据应以部署环境实测为准。整体而言该方案为多智能体交易分析框架提供了可靠的历史数据底座——上游同步服务统一写入、中间服务层统一管理、下游 API 统一消费三数据源从“各自为政”走向“同源可对比、可追溯、可监控”。【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →