ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

TradingAgents-CN 多数据源隔离架构与实时行情增强实践指南

2026/9/10 7:38:44 拓冰建站 浏览量
TradingAgents-CN 多数据源隔离架构与实时行情增强实践指南 TradingAgents-CN 多数据源隔离架构与实时行情增强实践指南【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN日期: 2025-10-28标签:多数据源实时数据PE/PB计算K线图数据隔离 概述本文是 TradingAgents-CN基于多智能体 LLM 的中文金融交易框架增强版于 2025 年 10 月 28 日完成的一次重大数据层架构升级的技术总结。本次升级通过25 个提交、60 个文件修改、约 3,500 行新增代码完成了多数据源隔离存储设计、实时 PE/PB 计算优化、K线图实时数据支持、实时行情同步状态追踪等多项核心能力显著提升了系统在数据完整性、实时性和可靠性三个维度的表现。阅读本文后你将掌握如何在 MongoDB 中通过联合唯一索引实现多数据源隔离存储如何设计与实现 PE/PB 计算的多层回退策略如何让 K 线图在盘中实时融合当日行情以及如何为实时行情同步建立用户可见的状态追踪机制。文中所有代码均可与仓库源码一一对应验证。 一、多数据源隔离存储架构1.1 问题背景TradingAgents-CN 的数据层同时支持 Tushare、AKShare、BaoStock 三个 A 股数据源但在升级前存在严重的数据覆盖问题数据覆盖问题stock_basic_info集合使用code作为唯一索引后运行的同步任务会覆盖先运行的数据源写入的记录无法保留不同数据源的独立数据。数据源优先级不统一不同模块各自使用不同的数据源查询导致同一只股票在不同页面展示的行业、市值等字段不一致用户体验混乱。索引冲突多数据源并发同步时频繁出现E11000 duplicate key error同步任务直接失败、数据不完整。典型的报错信息如下E11000 duplicate key error collection: tradingagents.stock_basic_info index: code_1 dup key: { code: 000001 }1.2 解决方案联合唯一索引 优先级查询核心思路在同一个集合中通过(code, source)联合唯一索引实现数据源隔离。同一只股票允许存在多条记录每条记录通过source字段标记其来源。// 联合唯一索引 db.stock_basic_info.createIndex( { code: 1, source: 1 }, { unique: true } ); // 辅助索引 db.stock_basic_info.createIndex({ code: 1 }); // 查询所有数据源 db.stock_basic_info.createIndex({ source: 1 }); // 按数据源查询升级后的文档结构示例{ code: 000001, source: tushare, name: 平安银行, industry: 银行, list_date: 19910403, ... }该索引设计在仓库中可从 scripts/docker_deployment_init.py 得到验证——部署初始化脚本正是以(code, 1), (source, 1)的uniqueTrue联合索引初始化stock_basic_info集合并同时建立code与source的辅助索引。1.3 索引迁移脚本对于已经存在旧索引的存量环境需要迁移脚本完成删旧建新# scripts/migrations/migrate_stock_basic_info_add_source_index.py async def migrate_stock_basic_info_indexes(): 迁移 stock_basic_info 集合的索引 # 1. 删除旧的 code 唯一索引 try: await db.stock_basic_info.drop_index(code_1) logger.info(✅ 已删除旧索引: code_1) except Exception as e: logger.warning(f⚠️ 删除旧索引失败可能不存在: {e}) # 2. 创建新的联合唯一索引 await db.stock_basic_info.create_index( [(code, 1), (source, 1)], uniqueTrue, namecode_source_unique ) logger.info(✅ 已创建联合唯一索引: (code, source)) # 3. 创建辅助索引 await db.stock_basic_info.create_index([(code, 1)]) await db.stock_basic_info.create_index([(source, 1)]) logger.info(✅ 已创建辅助索引)1.4 统一数据源优先级查询查询侧的核心改动在 app/services/stock_data_service.py当调用方未显式指定数据源时按tushare multi_source akshare baostock的固定优先级依次尝试命中即返回# app/services/stock_data_service.py async def get_stock_basic_info( self, symbol: str, source: Optional[str] None ) - Optional[StockBasicInfoExtended]: 获取股票基础信息 Args: symbol: 6位股票代码 source: 数据源 (tushare/akshare/baostock/multi_source) 默认优先级tushare multi_source akshare baostock symbol6 symbol.lstrip(shsz).zfill(6) if source: # 指定数据源 query {code: symbol6, source: source} doc await db[stock_basic_info].find_one(query, {_id: 0}) else: # 未指定数据源按优先级查询 source_priority [tushare, multi_source, akshare, baostock] doc None for src in source_priority: query {code: symbol6, source: src} doc await db[stock_basic_info].find_one(query, {_id: 0}) if doc: logger.debug(f✅ 使用数据源: {src}) break if not doc: logger.warning(f⚠️ 未找到股票信息: {symbol}) return None return StockBasicInfoExtended(**doc)1.5 修复多数据源同步服务同步侧写入侧的实现位于 app/services/multi_source_basics_sync_service.py其关键点是使用(code, source)作为UpdateOne的联合查询条件并配合upsertTrue从而实现按数据源独立写入、互不覆盖# app/services/multi_source_basics_sync_service.py async def sync_from_source(self, source: str): 从指定数据源同步股票基础信息 # 获取数据 stocks_data await self._fetch_data_from_source(source) # 批量更新使用 upsert operations [] for stock in stocks_data: operations.append( UpdateOne( {code: stock[code], source: source}, # 联合查询条件 {$set: stock}, upsertTrue ) ) # 执行批量操作 if operations: result await db.stock_basic_info.bulk_write(operations) logger.info(f✅ {source}: 更新 {result.modified_count} 条插入 {result.upserted_count} 条)从源码层面还可以看到两个与数据一致性修复直接相关的实现细节分批写入run_full_sync中以batch_size 500为界分批执行bulk_write避免大批量写入导致 MongoDB 超时见 multi_source_basics_sync_service.py。带重试的批量写入_execute_bulk_write_with_retry对bulk_write提供最多 3 次重试采用2 ** retry_count的指数退避2 秒、4 秒、8 秒显著缓解了网络抖动和超时导致的同步失败见 multi_source_basics_sync_service.py。数据源标识文档中的source字段直接取自实际成功拉取数据的适配器source_used不再使用模糊的默认值保证每条记录都有明确来源归属。升级效果✅ 同一股票可以有多条记录不同数据源✅ 保证(code, source)组合唯一✅ 支持灵活查询指定数据源或按优先级✅ 彻底解决索引冲突问题 二、实时 PE/PB 计算优化2.1 问题背景实时 PE/PB 计算依赖多个数据源但升级前存在三类问题数据缺失实时股价可能为空、财务数据可能未同步、总股本数据可能缺失。计算错误单位转换错误元/万元/亿元混用、除零错误、负值处理不当。无回退机制计算失败直接返回None用户看不到任何数据体验不佳。2.2 多层回退策略设计升级后的实时估值模块位于 tradingagents/dataflows/realtime_metrics.py其对外主入口get_pe_pb_with_fallback实现了动态计算优先、静态数据兜底的智能降级策略# tradingagents/dataflows/realtime_metrics.py async def get_realtime_pe_pb( self, symbol: str, source: str tushare ) - Dict[str, Optional[float]]: 获取实时PE/PB多层回退策略 回退策略 1. 优先使用实时股价计算 2. 降级使用数据库缓存值 3. 最后使用历史数据 result { pe: None, pb: None, total_mv: None, data_source: None } # 策略1使用实时股价计算 try: realtime_quote await self._get_realtime_quote(symbol) if realtime_quote and realtime_quote.get(close): pe, pb, total_mv await self._calculate_from_realtime( symbol, realtime_quote[close], source ) if pe or pb: result.update({ pe: pe, pb: pb, total_mv: total_mv, data_source: realtime_calculated }) return result except Exception as e: logger.warning(f⚠️ 实时计算失败: {e}) # 策略2使用数据库缓存值 try: cached_data await self._get_cached_pe_pb(symbol, source) if cached_data and (cached_data.get(pe) or cached_data.get(pb)): result.update({ pe: cached_data.get(pe), pb: cached_data.get(pb), total_mv: cached_data.get(total_mv), data_source: database_cached }) return result except Exception as e: logger.warning(f⚠️ 缓存查询失败: {e}) # 策略3使用历史数据 try: historical_data await self._get_historical_pe_pb(symbol, source) if historical_data and (historical_data.get(pe) or historical_data.get(pb)): result.update({ pe: historical_data.get(pe), pb: historical_data.get(pb), total_mv: historical_data.get(total_mv), data_source: historical_data }) return result except Exception as e: logger.warning(f⚠️ 历史数据查询失败: {e}) logger.warning(f⚠️ {symbol}: 所有策略均失败返回空值) return result2.3 实时市值与 PE/PB 计算修复实时计算的单位换算逻辑是本次修复的重点。仓库中 realtime_metrics.py 的calculate_realtime_pe_pb给出了完整的推导链# tradingagents/dataflows/realtime_metrics.py async def _calculate_from_realtime( self, symbol: str, current_price: float, source: str ) - Tuple[Optional[float], Optional[float], Optional[float]]: 使用实时股价计算PE/PB/市值 # 获取财务数据 financial_data await self._get_financial_data(symbol, source) if not financial_data: return None, None, None # 获取总股本单位万股 total_share financial_data.get(total_share) if not total_share or total_share 0: logger.warning(f⚠️ {symbol}: 总股本数据缺失或无效) return None, None, None # 计算实时市值单位亿元 # total_share 单位万股 # current_price 单位元 # 市值 总股本(万股) * 股价(元) / 10000 亿元 total_mv (total_share * current_price) / 10000 # 计算PE市盈率 net_profit financial_data.get(net_profit) # 单位元 if net_profit and net_profit 0: # 市值(亿元) / 净利润(亿元) PE net_profit_billion net_profit / 100000000 pe total_mv / net_profit_billion else: pe None # 计算PB市净率 net_assets financial_data.get(net_assets) # 单位元 if net_assets and net_assets 0: # 市值(亿元) / 净资产(亿元) PB net_assets_billion net_assets / 100000000 pb total_mv / net_assets_billion else: pb None return pe, pb, total_mv源码中的实际实现比博文示例更进一步——calculate_realtime_pe_pb采用了一条更精妙的链路见 realtime_metrics.py从stock_basic_info获取 Tushare 官方pe_ttm基于昨日收盘价用昨日市值反推TTM 净利润 昨日市值 / pe_ttm用实时股价 × 总股本计算实时市值动态 PE_TTM 实时市值 / TTM 净利润从而让估值指标实时反映盘中股价波动。同时在计算前会判断stock_basic_info是否已在今天 15:00 收盘后更新过——若已更新则直接采用其收盘后最新数据source stock_basic_info_latest避免重复计算。2.4 合理性校验为了防止脏数据进入前端模块还提供了validate_pe_pb校验函数见 realtime_metrics.pydef validate_pe_pb(pe: Optional[float], pb: Optional[float]) - bool: 验证PE/PB是否在合理范围内 # PE合理范围-100 到 1000允许负值因为亏损企业PE为负 if pe is not None and (pe -100 or pe 1000): logger.warning(fPE异常: {pe}) return False # PB合理范围0.1 到 100 if pb is not None and (pb 0.1 or pb 100): logger.warning(fPB异常: {pb}) return False return Trueget_pe_pb_with_fallback在拿到动态计算结果后会先经过该校验超出合理范围如亏损导致的异常负 PE则自动降级到 Tushare 静态 PEsource daily_basic。效果✅ 三层回退策略保证数据可用性✅ 实时市值计算准确✅ PE/PB 单位转换正确✅ 详细的数据来源标识realtime_calculated/database_cached/historical_data/daily_basic 三、K线图实时数据支持3.1 当天实时 K 线数据K 线接口自动从market_quotes集合获取当天实时数据实现盘中实时更新实现在 app/routers/stocks.py 的get_kline端点中。该端点同时兼容 A 股/港股/美股港股、美股走ForeignStockServiceA 股部分融合了实时行情逻辑。核心判定逻辑日线 A 股获取北京时间当前日期today_str检查历史 K 线中是否已存在当天的数据has_today_data判断是否处于交易时间09:30 now 15:00且weekday 5当处于交易时间内或收盘后但历史数据缺少当天数据时从market_quotes拉取实时行情并构造当日 K 线# app/routers/stocks.py router.get(/{code}/kline, response_modeldict) async def get_kline( code: str, period: str day, limit: int 120, adj: str none ): 获取K线数据支持当天实时数据 # 获取历史K线数据 items await historical_service.get_kline_data( symbolcode, periodperiod, limitlimit, adjadj ) # 检查是否需要添加当天实时数据仅针对日线 if period day and items: # 获取当前时间北京时间 tz ZoneInfo(settings.TIMEZONE) now datetime.now(tz) today_str now.strftime(%Y%m%d) current_time now.time() # 检查历史数据中是否已有当天的数据 has_today_data any( item.get(time) today_str for item in items ) # 判断是否在交易时间内 is_trading_time ( dtime(9, 30) current_time dtime(15, 0) and now.weekday() 5 # 周一到周五 ) # 如果在交易时间内或者收盘后但历史数据没有当天数据则从 market_quotes 获取 should_fetch_realtime is_trading_time or not has_today_data if should_fetch_realtime: # 从 market_quotes 获取实时行情 code_padded code.zfill(6) realtime_quote await market_quotes_coll.find_one( {code: code_padded}, {_id: 0} ) if realtime_quote: # 构造当天的K线数据 today_kline { time: today_str, open: float(realtime_quote.get(open, 0)), high: float(realtime_quote.get(high, 0)), low: float(realtime_quote.get(low, 0)), close: float(realtime_quote.get(close, 0)), volume: float(realtime_quote.get(volume, 0)), amount: float(realtime_quote.get(amount, 0)), } # 添加到结果中 if has_today_data: # 替换已有的当天数据 items [item for item in items if item.get(time) ! today_str] items.append(today_kline) items.sort(keylambda x: x[time]) logger.info(f✅ {code}: 添加当天实时K线数据) return { code: code, period: period, limit: limit, adj: adj, source: mongodbmarket_quotes, items: items }从 app/routers/stocks.py 的最新实现可以看到一个细节改进构造当日 K 线时统一使用YYYY-MM-DD日期格式today_str_formatted与历史数据的时间字段格式保持一致并在已有当天数据时替换末条、否则追加避免出现重复的当日 K 线。效果✅ 交易时间内显示实时 K 线✅ 收盘后自动补充当天数据✅ 无需等待历史数据同步✅ 用户体验显著提升 四、实时行情同步状态追踪4.1 同步状态追踪与收盘后缓冲期实时行情同步服务的实现在 app/services/quotes_ingestion_service.py。服务类维护last_sync_time、sync_statusidle/syncing/success/error、sync_error三个状态字段并提供_is_sync_time判断同步窗口# app/services/quotes_ingestion_service.py class QuotesIngestionService: 实时行情同步服务 def __init__(self): self.last_sync_time: Optional[datetime] None self.sync_status: str idle # idle, syncing, success, error self.sync_error: Optional[str] None async def sync_realtime_quotes(self): 同步实时行情 try: self.sync_status syncing self.sync_error None # 判断是否在交易时间内含收盘后30分钟缓冲期 if not self._is_sync_time(): logger.info(⏸️ 非交易时间跳过同步) self.sync_status idle return # 同步数据 await self._fetch_and_save_quotes() # 更新状态 self.last_sync_time datetime.now(ZoneInfo(Asia/Shanghai)) self.sync_status success logger.info(f✅ 实时行情同步成功: {self.last_sync_time}) except Exception as e: self.sync_status error self.sync_error str(e) logger.error(f❌ 实时行情同步失败: {e}) def _is_sync_time(self) - bool: 判断是否在同步时间内交易时间 收盘后30分钟缓冲期 now datetime.now(ZoneInfo(Asia/Shanghai)) current_time now.time() # 周末不同步 if now.weekday() 5: return False # 交易时间9:30-15:00 # 缓冲期15:00-15:30收盘后30分钟 return dtime(9, 30) current_time dtime(15, 30) def get_sync_status(self) - Dict: 获取同步状态 return { status: self.sync_status, last_sync_time: self.last_sync_time.isoformat() if self.last_sync_time else None, error: self.sync_error, is_trading_time: self._is_sync_time() }源码实现quotes_ingestion_service.py比博文示例包含更多工程细节值得展开说明调度频率由settings.QUOTES_INGEST_INTERVAL_SECONDS控制默认 360 秒6 分钟一次全市场近实时行情入库。接口轮换_rotation_sources [tushare, akshare_eastmoney, akshare_sina]按序轮换调用避免单一接口被限流。智能限流Tushare 免费用户每小时最多 2 次调用_tushare_hourly_limit 2付费用户自动切换到高频模式使用deque记录调用时间实现滑动窗口限流。状态持久化_record_sync_status将last_sync_time、data_source、records_count等写入quotes_ingestion_status集合便于前端与运维查看历史状态。该状态通过 app/routers/stock_data.py 的GET /sync-status/quotes端点对外暴露多数据源基础信息同步状态则通过 app/routers/multi_source_sync.py 的GET /status端点暴露数据持久化在sync_status集合job 键为stock_basics_multi_source。4.2 前端状态显示前端在个股详情页展示同步状态徽标实现在 frontend/src/views/Stocks/Detail.vue!-- frontend/src/views/Stocks/Detail.vue -- template el-card template #header div styledisplay: flex; justify-content: space-between; align-items: center; span实时行情/span !-- 同步状态指示器 -- el-tag :typesyncStatusType sizesmall effectplain el-icon stylemargin-right: 4px; component :issyncStatusIcon / /el-icon {{ syncStatusText }} /el-tag /div /template !-- 行情数据 -- div classquote-data ... /div /el-card /template script setup langts import { ref, computed, onMounted } from vue import { stockApi } from /api/stocks const syncStatus refany(null) // 获取同步状态 const fetchSyncStatus async () { try { const response await stockApi.getQuotesSyncStatus() syncStatus.value response.data } catch (error) { console.error(获取同步状态失败:, error) } } // 状态显示 const syncStatusType computed(() { if (!syncStatus.value) return info switch (syncStatus.value.status) { case success: return success case syncing: return warning case error: return danger default: return info } }) const syncStatusText computed(() { if (!syncStatus.value) return 未知 const lastSyncTime syncStatus.value.last_sync_time ? new Date(syncStatus.value.last_sync_time).toLocaleTimeString(zh-CN) : 从未同步 switch (syncStatus.value.status) { case success: return 已同步 (${lastSyncTime}) case syncing: return 同步中... case error: return 同步失败 default: return 空闲 } }) onMounted(() { fetchSyncStatus() // 每30秒刷新一次状态 setInterval(fetchSyncStatus, 30000) }) /script效果✅ 用户可以看到数据同步状态✅ 显示最后同步时间✅ 收盘后 30 分钟缓冲期✅ 前端每 30 秒自动刷新状态 五、其他优化5.1 为 stock_basic_info 添加 symbol 字段为stock_basic_info集合补充symbol字段带市场前缀的完整代码解决不同模块对股票代码表示法理解不一致的问题。从 app/services/multi_source_basics_sync_service.py 可以看到同步文档同时携带code6 位、symbol标准化 6 位、full_symbol带交易所后缀如000001.SZ三个字段其中full_symbol由_generate_full_symbol按代码前缀推断交易所60/68/90开头归上交所00/30/20开头归深交所8/4开头归北交所。# 示例 { code: 000001, # 6位代码 symbol: sz000001, # 带市场前缀 source: tushare, ... }5.2 基本面快照接口增强添加市销率PS动态计算# 使用实时市值和TTM营业收入计算 ps total_mv / revenue_ttm if revenue_ttm else None添加更多财务指标营业收入TTM、净利润TTM、净资产、ROE净资产收益率。这些指标在仓库的估值输出链路中均有体现例如 tradingagents/dataflows/optimized_china_data.py 的估值报告中同时输出市盈率PE、市盈率TTMPE_TTM、市净率PB、市销率PS、股息收益率。5.3 统一导出报告文件名格式// 统一格式TradingAgents_报告类型_股票代码_日期时间.pdf const filename TradingAgents_${reportType}_${stockCode}_${timestamp}.pdf // 示例 // TradingAgents_分析报告_000001_20251028_143052.pdf // TradingAgents_批量分析_20251028_143052.pdf5.4 股票名称获取增强为股票名称获取增加多层降级逻辑# 多层降级策略 # 1. 从 stock_basic_info 获取 # 2. 从 market_quotes 获取 # 3. 使用股票代码作为后备该策略在 app/routers/stocks.py 的行情查询中亦有体现实时数据如turnover_rate优先从market_quotes获取缺失时再降级到stock_basic_info的日度数据。 统计数据提交统计2025-10-28总提交数: 25 个修改文件数: 60 个新增代码: ~3,500 行删除代码: ~500 行净增代码: ~3,000 行功能分类多数据源架构: 6 项改进实时数据: 8 项增强PE/PB计算: 3 项优化K线图: 1 项新功能其他优化: 7 项改进 技术亮点总结1. 多数据源隔离存储设计核心思路联合唯一索引 数据源优先级查询。索引与查询策略在 scripts/docker_deployment_init.py初始化与 app/services/stock_data_service.py查询两处形成闭环// 索引设计 db.stock_basic_info.createIndex({ code: 1, source: 1 }, { unique: true }); // 查询策略 source_priority [tushare, multi_source, akshare, baostock]2. 实时 PE/PB 三层回退策略实时股价计算最准确数据库缓存值次优历史数据保底再叠加validate_pe_pb合理性校验与收盘后已更新则直接复用的短路逻辑保证任何时刻都能返回可用的估值数据。3. K 线图实时数据融合交易时间内从market_quotes获取实时数据收盘后补充当天数据如果历史数据未同步非交易日只显示历史数据4. 同步状态追踪实时状态更新idle/syncing/success/error收盘后 30 分钟缓冲期前端每 30 秒自动刷新 总结今日成果✅25 次提交✅60 个文件修改✅3,500 行新增代码核心价值多数据源架构完善通过(code, source)联合唯一索引彻底解决索引冲突支持数据源隔离存储统一数据源优先级查询配合分批写入与指数退避重试机制从根源上消除了多数据源同步时的数据覆盖与写入失败问题。实时数据能力提升K 线图支持盘中实时数据融合PE/PB 实时计算采用动态计算 → 缓存 → 历史三层回退同步状态全链路可视化。数据准确性改善修复市值计算与单位转换逻辑万股 × 元 / 10000 亿元引入 PE/PB 合理性校验杜绝脏数据进入前端。用户体验优化实时数据展示、同步状态追踪、报告文件名格式统一让用户对数据从哪来、有多新一目了然。这套数据层架构为 TradingAgents-CN 的多智能体分析提供了坚实可靠的数据底座隔离存储保证了多数据源并存不冲突多层回退保证了估值指标在任何数据缺失场景下仍可用实时融合与状态追踪则让盘中分析与决策始终基于最新行情。【免费下载链接】TradingAgents-CN基于多智能体LLM的中文金融交易框架 - TradingAgents中文增强版项目地址: https://gitcode.com/GitHub_Trending/tr/TradingAgents-CN创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考