diff --git a/ai_agents.py b/ai_agents.py index f23080a..1523c7b 100644 --- a/ai_agents.py +++ b/ai_agents.py @@ -1,5 +1,6 @@ from deepseek_client import DeepSeekClient from typing import Dict, Any +import concurrent.futures import time import config @@ -13,7 +14,6 @@ class StockAnalysisAgents: def technical_analyst_agent(self, stock_info: Dict, stock_data: Any, indicators: Dict) -> Dict[str, Any]: """技术面分析智能体""" print("🔍 技术分析师正在分析中...") - time.sleep(1) # 模拟分析时间 analysis = self.deepseek_client.technical_analysis(stock_info, stock_data, indicators) @@ -38,8 +38,6 @@ class StockAnalysisAgents: else: print(" ⚠ 未获取到季报数据,将基于基本财务数据分析") - time.sleep(1) - analysis = self.deepseek_client.fundamental_analysis(stock_info, financial_data, quarterly_data) return { @@ -61,8 +59,6 @@ class StockAnalysisAgents: else: print(" ⚠ 未获取到资金流向数据,将基于技术指标分析") - time.sleep(1) - analysis = self.deepseek_client.fund_flow_analysis(stock_info, indicators, fund_flow_data) return { @@ -84,8 +80,6 @@ class StockAnalysisAgents: else: print(" ⚠ 未获取到风险数据,将基于基本信息分析") - time.sleep(1) - # 构建风险数据文本 risk_data_text = "" if risk_data and risk_data.get('data_success'): @@ -227,8 +221,6 @@ class StockAnalysisAgents: else: print(" ⚠ 未获取到详细情绪数据,将基于基本信息分析") - time.sleep(1) - # 构建带有市场情绪数据的prompt sentiment_data_text = "" if sentiment_data and sentiment_data.get('data_success'): @@ -317,8 +309,6 @@ class StockAnalysisAgents: else: print(" ⚠ 未获取到新闻数据,将基于基本信息分析") - time.sleep(1) - # 构建带有新闻数据的prompt news_text = "" if news_data and news_data.get('data_success'): @@ -436,32 +426,38 @@ class StockAnalysisAgents: print(f"📋 参与分析的分析师: {', '.join(active_analysts)}") print("=" * 50) - # 并行运行各个分析师 + # 并行运行各个分析师(全部同时启动,等待最慢的一个完成) agents_results = {} - # 技术面分析 + agent_jobs = [] if enabled_analysts.get('technical', True): - agents_results["technical"] = self.technical_analyst_agent(stock_info, stock_data, indicators) - - # 基本面分析 + agent_jobs.append(("technical", lambda: self.technical_analyst_agent(stock_info, stock_data, indicators))) if enabled_analysts.get('fundamental', True): - agents_results["fundamental"] = self.fundamental_analyst_agent(stock_info, financial_data, quarterly_data) - - # 资金面分析(传入资金流向数据) + agent_jobs.append(("fundamental", lambda: self.fundamental_analyst_agent(stock_info, financial_data, quarterly_data))) if enabled_analysts.get('fund_flow', True): - agents_results["fund_flow"] = self.fund_flow_analyst_agent(stock_info, indicators, fund_flow_data) - - # 风险管理分析(传入风险数据) + agent_jobs.append(("fund_flow", lambda: self.fund_flow_analyst_agent(stock_info, indicators, fund_flow_data))) if enabled_analysts.get('risk', True): - agents_results["risk_management"] = self.risk_management_agent(stock_info, indicators, risk_data) - - # 市场情绪分析(传入市场情绪数据) + agent_jobs.append(("risk_management", lambda: self.risk_management_agent(stock_info, indicators, risk_data))) if enabled_analysts.get('sentiment', False): - agents_results["market_sentiment"] = self.market_sentiment_agent(stock_info, sentiment_data) - - # 新闻分析(传入新闻数据) + agent_jobs.append(("market_sentiment", lambda: self.market_sentiment_agent(stock_info, sentiment_data))) if enabled_analysts.get('news', False): - agents_results["news"] = self.news_analyst_agent(stock_info, news_data) + agent_jobs.append(("news", lambda: self.news_analyst_agent(stock_info, news_data))) + + with concurrent.futures.ThreadPoolExecutor(max_workers=max(min(len(agent_jobs), 6), 1)) as executor: + future_to_key = {executor.submit(job): key for key, job in agent_jobs} + for future in concurrent.futures.as_completed(future_to_key): + key = future_to_key[future] + try: + agents_results[key] = future.result() + except Exception as e: + print(f"❌ {key} 分析师并行分析失败: {e}") + agents_results[key] = { + "agent_name": key, + "agent_role": "", + "analysis": f"分析失败: {e}", + "focus_areas": [], + "timestamp": time.strftime("%Y-%m-%d %H:%M:%S") + } print("✅ 所有已选择的分析师完成分析") print("=" * 50) @@ -471,7 +467,6 @@ class StockAnalysisAgents: def conduct_team_discussion(self, agents_results: Dict[str, Any], stock_info: Dict) -> str: """进行团队讨论""" print("🤝 分析团队正在进行综合讨论...") - time.sleep(2) # 收集参与分析的分析师名单和报告 participants = [] @@ -538,7 +533,6 @@ class StockAnalysisAgents: def make_final_decision(self, discussion_result: str, stock_info: Dict, indicators: Dict) -> Dict[str, Any]: """制定最终投资决策""" print("📋 正在制定最终投资决策...") - time.sleep(1) decision = self.deepseek_client.final_decision(discussion_result, stock_info, indicators) diff --git a/app.py b/app.py index a9af32f..8728f32 100644 --- a/app.py +++ b/app.py @@ -9,6 +9,10 @@ import base64 import os import config +# 注入所有外部请求(akshare/tushare等)的默认超时 +from http_timeout import install_default_requests_timeout +install_default_requests_timeout() + from stock_data import StockDataFetcher from ai_agents import StockAnalysisAgents from pdf_generator import display_pdf_export_section diff --git a/data_source_manager.py b/data_source_manager.py index 9931871..298c0e0 100644 --- a/data_source_manager.py +++ b/data_source_manager.py @@ -11,6 +11,10 @@ from dotenv import load_dotenv # 加载环境变量 load_dotenv() +# 注入外部请求默认超时(覆盖akshare、tushare等基于requests的调用) +from http_timeout import install_default_requests_timeout +install_default_requests_timeout() + class DataSourceManager: """数据源管理器 - 实现akshare与tushare自动切换""" @@ -25,7 +29,7 @@ class DataSourceManager: try: import tushare as ts ts.set_token(self.tushare_token) - self.tushare_api = ts.pro_api() + self.tushare_api = ts.pro_api(timeout=float(os.getenv('TUSHARE_TIMEOUT', '15'))) self.tushare_available = True print("✅ Tushare数据源初始化成功") except Exception as e: @@ -36,7 +40,7 @@ class DataSourceManager: def get_stock_hist_data(self, symbol, start_date=None, end_date=None, adjust='qfq'): """ - 获取股票历史数据(优先akshare,失败时使用tushare) + 获取股票历史数据(优先tushare,失败时使用akshare) Args: symbol: 股票代码(6位数字) @@ -55,10 +59,63 @@ class DataSourceManager: else: end_date = datetime.now().strftime('%Y%m%d') - # 优先使用akshare + # 优先使用tushare + if self.tushare_available: + try: + import tushare as ts + print(f"[Tushare] 正在获取 {symbol} 的历史数据(主要数据源)...") + + # 转换股票代码格式(添加市场后缀) + ts_code = self._convert_to_ts_code(symbol) + + # 转换复权类型 + adj_dict = {'qfq': 'qfq', 'hfq': 'hfq', '': None} + adj = adj_dict.get(adjust, 'qfq') + + if adj is None: + # 不复权数据直接使用daily接口 + df = self.tushare_api.daily( + ts_code=ts_code, + start_date=start_date, + end_date=end_date + ) + else: + # 复权数据使用pro_bar(daily接口不支持adj参数) + df = ts.pro_bar( + api=self.tushare_api, + ts_code=ts_code, + start_date=start_date, + end_date=end_date, + adj=adj, + retry_count=1 + ) + + if df is not None and not df.empty: + # 标准化列名和数据格式 + df = df.rename(columns={ + 'trade_date': 'date', + 'vol': 'volume', + 'amount': 'amount' + }) + df['date'] = pd.to_datetime(df['date']) + df = df.sort_values('date') + + # 转换成交量单位(tushare单位是手,转换为股) + df['volume'] = df['volume'] * 100 + # 转换成交额单位(tushare单位是千元,转换为元) + df['amount'] = df['amount'] * 1000 + + print(f"[Tushare] ✅ 成功获取 {len(df)} 条数据") + return df + else: + print(f"[Tushare] ❌ 未获取到数据,尝试备用数据源") + except Exception as e: + print(f"[Tushare] ❌ 获取失败: {e}") + + # tushare失败,回退到akshare try: import akshare as ak - print(f"[Akshare] 正在获取 {symbol} 的历史数据...") + print(f"[Akshare] 正在获取 {symbol} 的历史数据(备用数据源)...") df = ak.stock_zh_a_hist( symbol=symbol, @@ -86,60 +143,18 @@ class DataSourceManager: df['date'] = pd.to_datetime(df['date']) print(f"[Akshare] ✅ 成功获取 {len(df)} 条数据") return df + else: + print(f"[Akshare] ❌ 未获取到数据") except Exception as e: print(f"[Akshare] ❌ 获取失败: {e}") - # akshare失败,尝试tushare - if self.tushare_available: - try: - print(f"[Tushare] 正在获取 {symbol} 的历史数据(备用数据源)...") - - # 转换股票代码格式(添加市场后缀) - ts_code = self._convert_to_ts_code(symbol) - - # 转换复权类型 - adj_dict = {'qfq': 'qfq', 'hfq': 'hfq', '': None} - adj = adj_dict.get(adjust, 'qfq') - - # 格式化日期 - start = f"{start_date[:4]}-{start_date[4:6]}-{start_date[6:]}" if start_date else None - end = f"{end_date[:4]}-{end_date[4:6]}-{end_date[6:]}" if end_date else None - - # 获取数据 - df = self.tushare_api.daily( - ts_code=ts_code, - start_date=start_date, - end_date=end_date, - adj=adj - ) - - if df is not None and not df.empty: - # 标准化列名和数据格式 - df = df.rename(columns={ - 'trade_date': 'date', - 'vol': 'volume', - 'amount': 'amount' - }) - df['date'] = pd.to_datetime(df['date']) - df = df.sort_values('date') - - # 转换成交量单位(tushare单位是手,转换为股) - df['volume'] = df['volume'] * 100 - # 转换成交额单位(tushare单位是千元,转换为元) - df['amount'] = df['amount'] * 1000 - - print(f"[Tushare] ✅ 成功获取 {len(df)} 条数据") - return df - except Exception as e: - print(f"[Tushare] ❌ 获取失败: {e}") - # 两个数据源都失败 print("❌ 所有数据源均获取失败") return None def get_stock_basic_info(self, symbol): """ - 获取股票基本信息(优先akshare,失败时使用tushare) + 获取股票基本信息(优先tushare,失败时使用akshare) Args: symbol: 股票代码 @@ -154,10 +169,34 @@ class DataSourceManager: "market": "未知" } - # 优先使用akshare + # 优先使用tushare + if self.tushare_available: + try: + print(f"[Tushare] 正在获取 {symbol} 的基本信息(主要数据源)...") + + ts_code = self._convert_to_ts_code(symbol) + df = self.tushare_api.stock_basic( + ts_code=ts_code, + fields='ts_code,name,area,industry,market,list_date' + ) + + if df is not None and not df.empty: + info['name'] = df.iloc[0]['name'] + info['industry'] = df.iloc[0]['industry'] + info['market'] = df.iloc[0]['market'] + info['list_date'] = df.iloc[0]['list_date'] + + print(f"[Tushare] ✅ 成功获取基本信息") + return info + else: + print(f"[Tushare] ❌ 未获取到基本信息,尝试备用数据源") + except Exception as e: + print(f"[Tushare] ❌ 获取失败: {e}") + + # tushare失败,回退akshare try: import akshare as ak - print(f"[Akshare] 正在获取 {symbol} 的基本信息...") + print(f"[Akshare] 正在获取 {symbol} 的基本信息(备用数据源)...") stock_info = ak.stock_individual_info_em(symbol=symbol) if stock_info is not None and not stock_info.empty: @@ -181,33 +220,11 @@ class DataSourceManager: except Exception as e: print(f"[Akshare] ❌ 获取失败: {e}") - # akshare失败,尝试tushare - if self.tushare_available: - try: - print(f"[Tushare] 正在获取 {symbol} 的基本信息(备用数据源)...") - - ts_code = self._convert_to_ts_code(symbol) - df = self.tushare_api.stock_basic( - ts_code=ts_code, - fields='ts_code,name,area,industry,market,list_date' - ) - - if df is not None and not df.empty: - info['name'] = df.iloc[0]['name'] - info['industry'] = df.iloc[0]['industry'] - info['market'] = df.iloc[0]['market'] - info['list_date'] = df.iloc[0]['list_date'] - - print(f"[Tushare] ✅ 成功获取基本信息") - return info - except Exception as e: - print(f"[Tushare] ❌ 获取失败: {e}") - return info def get_realtime_quotes(self, symbol): """ - 获取实时行情数据(优先akshare,失败时使用tushare) + 获取实时行情数据(优先tushare,失败时使用akshare) Args: symbol: 股票代码 @@ -217,10 +234,51 @@ class DataSourceManager: """ quotes = {} - # 优先使用akshare + # 优先使用tushare + if self.tushare_available: + try: + print(f"[Tushare] 正在获取 {symbol} 的实时行情(主要数据源)...") + + ts_code = self._convert_to_ts_code(symbol) + today = datetime.now().strftime('%Y%m%d') + df = self.tushare_api.daily( + ts_code=ts_code, + start_date=today, + end_date=today + ) + if df is None or df.empty: + # 非交易日时,获取最近10个交易日的最新数据 + start = (datetime.now() - timedelta(days=10)).strftime('%Y%m%d') + df = self.tushare_api.daily( + ts_code=ts_code, + start_date=start, + end_date=today + ) + + if df is not None and not df.empty: + row = df.iloc[0] + quotes = { + 'symbol': symbol, + 'price': row['close'], + 'change_percent': row['pct_chg'], + 'volume': row['vol'] * 100, + 'amount': row['amount'] * 1000, + 'high': row['high'], + 'low': row['low'], + 'open': row['open'], + 'pre_close': row['pre_close'] + } + print(f"[Tushare] ✅ 成功获取实时行情") + return quotes + else: + print(f"[Tushare] ❌ 未获取到实时行情,尝试备用数据源") + except Exception as e: + print(f"[Tushare] ❌ 获取失败: {e}") + + # tushare失败,回退akshare try: import akshare as ak - print(f"[Akshare] 正在获取 {symbol} 的实时行情...") + print(f"[Akshare] 正在获取 {symbol} 的实时行情(备用数据源)...") df = ak.stock_zh_a_spot_em() stock_df = df[df['代码'] == symbol] @@ -245,41 +303,11 @@ class DataSourceManager: except Exception as e: print(f"[Akshare] ❌ 获取失败: {e}") - # akshare失败,尝试tushare - if self.tushare_available: - try: - print(f"[Tushare] 正在获取 {symbol} 的实时行情(备用数据源)...") - - ts_code = self._convert_to_ts_code(symbol) - df = self.tushare_api.daily( - ts_code=ts_code, - start_date=datetime.now().strftime('%Y%m%d'), - end_date=datetime.now().strftime('%Y%m%d') - ) - - if df is not None and not df.empty: - row = df.iloc[0] - quotes = { - 'symbol': symbol, - 'price': row['close'], - 'change_percent': row['pct_chg'], - 'volume': row['vol'] * 100, - 'amount': row['amount'] * 1000, - 'high': row['high'], - 'low': row['low'], - 'open': row['open'], - 'pre_close': row['pre_close'] - } - print(f"[Tushare] ✅ 成功获取实时行情") - return quotes - except Exception as e: - print(f"[Tushare] ❌ 获取失败: {e}") - return quotes def get_financial_data(self, symbol, report_type='income'): """ - 获取财务数据(优先akshare,失败时使用tushare) + 获取财务数据(优先tushare,失败时使用akshare) Args: symbol: 股票代码 @@ -288,30 +316,10 @@ class DataSourceManager: Returns: DataFrame: 财务数据 """ - # 优先使用akshare - try: - import akshare as ak - print(f"[Akshare] 正在获取 {symbol} 的财务数据...") - - if report_type == 'income': - df = ak.stock_financial_report_sina(stock=symbol, symbol="利润表") - elif report_type == 'balance': - df = ak.stock_financial_report_sina(stock=symbol, symbol="资产负债表") - elif report_type == 'cashflow': - df = ak.stock_financial_report_sina(stock=symbol, symbol="现金流量表") - else: - df = None - - if df is not None and not df.empty: - print(f"[Akshare] ✅ 成功获取财务数据") - return df - except Exception as e: - print(f"[Akshare] ❌ 获取失败: {e}") - - # akshare失败,尝试tushare + # 优先使用tushare if self.tushare_available: try: - print(f"[Tushare] 正在获取 {symbol} 的财务数据(备用数据源)...") + print(f"[Tushare] 正在获取 {symbol} 的财务数据(主要数据源)...") ts_code = self._convert_to_ts_code(symbol) @@ -327,9 +335,31 @@ class DataSourceManager: if df is not None and not df.empty: print(f"[Tushare] ✅ 成功获取财务数据") return df + else: + print(f"[Tushare] ❌ 未获取到财务数据,尝试备用数据源") except Exception as e: print(f"[Tushare] ❌ 获取失败: {e}") + # tushare失败,回退akshare + try: + import akshare as ak + print(f"[Akshare] 正在获取 {symbol} 的财务数据(备用数据源)...") + + if report_type == 'income': + df = ak.stock_financial_report_sina(stock=symbol, symbol="利润表") + elif report_type == 'balance': + df = ak.stock_financial_report_sina(stock=symbol, symbol="资产负债表") + elif report_type == 'cashflow': + df = ak.stock_financial_report_sina(stock=symbol, symbol="现金流量表") + else: + df = None + + if df is not None and not df.empty: + print(f"[Akshare] ✅ 成功获取财务数据") + return df + except Exception as e: + print(f"[Akshare] ❌ 获取失败: {e}") + return None def _convert_to_ts_code(self, symbol): @@ -376,4 +406,3 @@ class DataSourceManager: # 全局数据源管理器实例 data_source_manager = DataSourceManager() - diff --git a/deepseek_client.py b/deepseek_client.py index 8fd8460..7c7b13f 100644 --- a/deepseek_client.py +++ b/deepseek_client.py @@ -1,5 +1,6 @@ import openai import json +import os from typing import Dict, List, Any, Optional import config @@ -10,7 +11,9 @@ class DeepSeekClient: self.model = model or config.DEFAULT_MODEL_NAME self.client = openai.OpenAI( api_key=config.DEEPSEEK_API_KEY, - base_url=config.DEEPSEEK_BASE_URL + base_url=config.DEEPSEEK_BASE_URL, + timeout=float(os.getenv('DEEPSEEK_TIMEOUT', '180')), + max_retries=int(os.getenv('DEEPSEEK_MAX_RETRIES', '1')) ) def call_api(self, messages: List[Dict[str, str]], model: Optional[str] = None, diff --git a/fund_flow_akshare.py b/fund_flow_akshare.py index ebe992c..4d68f31 100644 --- a/fund_flow_akshare.py +++ b/fund_flow_akshare.py @@ -109,56 +109,54 @@ class FundFlowAkshareDataFetcher: def _get_individual_fund_flow(self, symbol, market): """获取个股资金流向数据(支持akshare和tushare自动切换)""" try: - # 优先使用akshare的stock_individual_fund_flow接口 - print(f" [Akshare] 正在获取资金流向 (市场: {market})...") - - df = ak.stock_individual_fund_flow(stock=symbol, market=market) - + # 优先使用tushare + if data_source_manager.tushare_available: + try: + print(f" [Tushare] 正在获取资金流向数据(主要数据源)...") + ts_code = data_source_manager._convert_to_ts_code(symbol) + + # 计算日期范围(最近N个交易日) + end_date = datetime.now().strftime('%Y%m%d') + start_date = (datetime.now() - timedelta(days=self.days * 2)).strftime('%Y%m%d') + + df = data_source_manager.tushare_api.moneyflow( + ts_code=ts_code, + start_date=start_date, + end_date=end_date + ) + + if df is not None and not df.empty: + # 标准化列名以匹配akshare格式(金额单位:千元→元) + df = df.rename(columns={ + 'trade_date': '日期', + 'close': '收盘价', + 'pct_chg': '涨跌幅', + 'net_mf_amount': '主力净流入-净额' + }) + df['主力净流入-净额'] = df['主力净流入-净额'] * 1000 + df['超大单净流入-净额'] = (df['buy_elg_amount'] - df['sell_elg_amount']) * 1000 + df['大单净流入-净额'] = (df['buy_lg_amount'] - df['sell_lg_amount']) * 1000 + df['中单净流入-净额'] = (df['buy_md_amount'] - df['sell_md_amount']) * 1000 + df['小单净流入-净额'] = (df['buy_sm_amount'] - df['sell_sm_amount']) * 1000 + # tushare按日期倒序返回,取最近N天后转为正序(与akshare一致) + df = df.head(self.days) + df = df.iloc[::-1].reset_index(drop=True) + print(f" [Tushare] ✅ 成功获取 {len(df)} 条资金流向数据") + else: + print(f" [Tushare] ❌ 未找到资金流向数据,尝试备用数据源") + df = None + except Exception as te: + print(f" [Tushare] ❌ 获取失败: {te}") + df = None + else: + df = None + + # tushare不可用或失败时,回退akshare if df is None or df.empty: - print(f" [Akshare] 未找到资金流向数据,尝试备用数据源...") - - # akshare失败,尝试tushare - if data_source_manager.tushare_available: - try: - print(f" [Tushare] 正在获取资金流向数据(备用数据源)...") - ts_code = data_source_manager._convert_to_ts_code(symbol) - - # 计算日期范围(最近N个交易日) - end_date = datetime.now().strftime('%Y%m%d') - start_date = (datetime.now() - timedelta(days=self.days * 2)).strftime('%Y%m%d') - - # 获取资金流向数据 - df = data_source_manager.tushare_api.moneyflow( - ts_code=ts_code, - start_date=start_date, - end_date=end_date - ) - - if df is not None and not df.empty: - # 标准化列名以匹配akshare格式 - df = df.rename(columns={ - 'trade_date': '日期', - 'buy_sm_amount': '小单买入', - 'sell_sm_amount': '小单卖出', - 'buy_md_amount': '中单买入', - 'sell_md_amount': '中单卖出', - 'buy_lg_amount': '大单买入', - 'sell_lg_amount': '大单卖出', - 'buy_elg_amount': '超大单买入', - 'sell_elg_amount': '超大单卖出', - 'net_mf_amount': '净额' - }) - - # 限制为最近N天 - df = df.head(self.days) - print(f" [Tushare] ✅ 成功获取 {len(df)} 条资金流向数据") - else: - print(f" [Tushare] ❌ 未找到资金流向数据") - return None - except Exception as te: - print(f" [Tushare] ❌ 获取失败: {te}") - return None - else: + print(f" [Akshare] 正在获取资金流向 (市场: {market})(备用数据源)...") + df = ak.stock_individual_fund_flow(stock=symbol, market=market) + if df is None or df.empty: + print(f" [Akshare] 未找到资金流向数据") return None # akshare 返回的数据是按时间正序排列(从旧到新),所以使用 tail() 获取最近N天的数据 @@ -341,4 +339,3 @@ if __name__ == "__main__": print(f"\n获取失败: {data.get('error', '未知错误')}") print("\n") - diff --git a/http_timeout.py b/http_timeout.py new file mode 100644 index 0000000..ef21cd2 --- /dev/null +++ b/http_timeout.py @@ -0,0 +1,77 @@ +""" +统一的外部请求超时控制工具 + +- install_default_requests_timeout: 给所有基于 requests 的外部请求 + (akshare、tushare、pywencai 等)注入默认超时,调用方未指定 timeout 时生效。 +- call_with_timeout: 在独立线程中执行阻塞调用,超过时限立即放弃等待, + 防止个别环节长时间卡死整个程序。 +""" + +import os +import threading + +# 默认连接超时 / 读取超时(单位:秒),可通过环境变量覆盖 +DEFAULT_CONNECT_TIMEOUT = float(os.getenv("HTTP_CONNECT_TIMEOUT", "6")) +DEFAULT_READ_TIMEOUT = float(os.getenv("HTTP_READ_TIMEOUT", "20")) +DEFAULT_CALL_TIMEOUT = float(os.getenv("HTTP_CALL_TIMEOUT", "30")) + +_installed = False +_install_lock = threading.Lock() + + +def install_default_requests_timeout(connect_timeout=None, read_timeout=None): + """给 requests 库注入默认超时(仅当调用方未显式指定 timeout 时生效)。""" + global _installed + with _install_lock: + if _installed: + return + + connect = DEFAULT_CONNECT_TIMEOUT if connect_timeout is None else connect_timeout + read = DEFAULT_READ_TIMEOUT if read_timeout is None else read_timeout + + try: + import requests.sessions + original_request = requests.sessions.Session.request + + def request_with_timeout(self, method, url, **kwargs): + if kwargs.get("timeout") is None: + kwargs["timeout"] = (connect, read) + return original_request(self, method, url, **kwargs) + + requests.sessions.Session.request = request_with_timeout + _installed = True + print(f"✅ 已为外部请求注入默认超时(连接 {connect}s / 读取 {read}s)") + except Exception as e: + print(f"⚠️ 注入 requests 默认超时失败: {e}") + + +def call_with_timeout(func, timeout=None, *args, **kwargs): + """在线程中执行函数,超过 timeout 秒则放弃等待并抛出 TimeoutError。 + + 注意:超时后线程会继续在后台运行直至结束(守护线程),不会影响主程序。 + """ + if timeout is None: + timeout = DEFAULT_CALL_TIMEOUT + if timeout <= 0: + return func(*args, **kwargs) + + result_box = {} + + def _run(): + try: + result_box["ok"] = True + result_box["value"] = func(*args, **kwargs) + except BaseException as e: # noqa: BLE001 + result_box["ok"] = False + result_box["error"] = e + + worker = threading.Thread(target=_run, daemon=True) + worker.start() + worker.join(timeout) + + if worker.is_alive(): + func_name = getattr(func, "__name__", "function") + raise TimeoutError(f"调用 {func_name} 超过 {timeout}s 未返回,已放弃等待") + if result_box.get("ok"): + return result_box["value"] + raise result_box.get("error", RuntimeError("未知错误")) diff --git a/low_price_bull_selector.py b/low_price_bull_selector.py index ec971a5..f41d524 100644 --- a/low_price_bull_selector.py +++ b/low_price_bull_selector.py @@ -10,6 +10,7 @@ import pywencai from datetime import datetime from typing import Tuple, Optional import time +from http_timeout import call_with_timeout class LowPriceBullSelector: @@ -60,7 +61,7 @@ class LowPriceBullSelector: print(f"正在调用问财接口...") # 调用pywencai - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None: return False, None, "问财接口返回None,请检查网络或稍后重试" diff --git a/main_force_selector.py b/main_force_selector.py index 3543009..cf8c541 100644 --- a/main_force_selector.py +++ b/main_force_selector.py @@ -11,6 +11,7 @@ import pywencai from datetime import datetime, timedelta from typing import Dict, List, Tuple import time +from http_timeout import call_with_timeout class MainForceStockSelector: """主力选股类""" @@ -71,7 +72,7 @@ class MainForceStockSelector: print(f"查询语句: {query[:100]}...") try: - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None: print(f" ⚠️ 方案{i}返回None,尝试下一个方案") @@ -387,4 +388,3 @@ class MainForceStockSelector: # 全局实例 main_force_selector = MainForceStockSelector() - diff --git a/market_sentiment_data.py b/market_sentiment_data.py index 03e3b17..516a19d 100644 --- a/market_sentiment_data.py +++ b/market_sentiment_data.py @@ -329,11 +329,56 @@ class MarketSentimentDataFetcher: } def _get_turnover_rate(self, symbol): - """获取换手率数据(支持akshare和tushare自动切换)""" + """获取换手率数据(优先tushare,失败时使用akshare)""" try: - # 优先使用akshare获取最近的换手率数据 - print(f" [Akshare] 正在获取换手率数据...") - # 获取A股实时行情数据(不需要参数) + # 优先使用tushare(daily_basic,取最近10个交易日保证非交易日也有数据) + if data_source_manager.tushare_available: + try: + print(f" [Tushare] 正在获取换手率数据(主要数据源)...") + ts_code = data_source_manager._convert_to_ts_code(symbol) + + end_date = datetime.now().strftime('%Y%m%d') + start_date = (datetime.now() - timedelta(days=10)).strftime('%Y%m%d') + df = data_source_manager.tushare_api.daily_basic( + ts_code=ts_code, + start_date=start_date, + end_date=end_date + ) + + if df is not None and not df.empty: + row = df.iloc[0] + turnover_rate = row.get('turnover_rate', 'N/A') + + # 解读换手率 + interpretation = "" + if turnover_rate != 'N/A': + try: + turnover = float(turnover_rate) + if turnover > 20: + interpretation = "换手率极高(>20%),资金活跃度极高,可能存在炒作" + elif turnover > 10: + interpretation = "换手率较高(>10%),交易活跃" + elif turnover > 5: + interpretation = "换手率正常(5%-10%),交易适中" + elif turnover > 2: + interpretation = "换手率偏低(2%-5%),交易相对清淡" + else: + interpretation = "换手率很低(<2%),交易清淡" + except: + pass + + print(f" [Tushare] ✅ 成功获取换手率: {turnover_rate}%") + return { + "current_turnover_rate": turnover_rate, + "interpretation": interpretation + } + else: + print(f" [Tushare] ❌ 未获取到换手率,尝试备用数据源") + except Exception as te: + print(f" [Tushare] ❌ 获取失败: {te}") + + # tushare失败,回退akshare + print(f" [Akshare] 正在获取换手率数据(备用数据源)...") df = ak.stock_zh_a_spot_em() if df is not None and not df.empty: stock_data = df[df['代码'] == symbol] @@ -366,57 +411,42 @@ class MarketSentimentDataFetcher: } except Exception as e: print(f" [Akshare] ❌ 获取换手率失败: {e}") - - # akshare失败,尝试tushare - if data_source_manager.tushare_available: - try: - print(f" [Tushare] 正在获取换手率数据(备用数据源)...") - ts_code = data_source_manager._convert_to_ts_code(symbol) - - # 获取最近一个交易日的数据 - df = data_source_manager.tushare_api.daily_basic( - ts_code=ts_code, - trade_date=datetime.now().strftime('%Y%m%d') - ) - - if df is not None and not df.empty: - row = df.iloc[0] - turnover_rate = row.get('turnover_rate', 'N/A') - - # 解读换手率 - interpretation = "" - if turnover_rate != 'N/A': - try: - turnover = float(turnover_rate) - if turnover > 20: - interpretation = "换手率极高(>20%),资金活跃度极高,可能存在炒作" - elif turnover > 10: - interpretation = "换手率较高(>10%),交易活跃" - elif turnover > 5: - interpretation = "换手率正常(5%-10%),交易适中" - elif turnover > 2: - interpretation = "换手率偏低(2%-5%),交易相对清淡" - else: - interpretation = "换手率很低(<2%),交易清淡" - except: - pass - - print(f" [Tushare] ✅ 成功获取换手率: {turnover_rate}%") - return { - "current_turnover_rate": turnover_rate, - "interpretation": interpretation - } - except Exception as te: - print(f" [Tushare] ❌ 获取失败: {te}") return None def _get_market_index_sentiment(self): """获取大盘指数情绪(支持akshare和tushare自动切换)""" try: - # 优先使用akshare获取上证指数实时数据 - print(f" [Akshare] 正在获取大盘指数数据...") - # 使用正确的symbol参数 + # 优先使用tushare(index_daily,取最近10个交易日保证非交易日也有数据) + if data_source_manager.tushare_available: + try: + print(f" [Tushare] 正在获取大盘指数数据(主要数据源)...") + + # 获取上证指数数据 + end_date = datetime.now().strftime('%Y%m%d') + start_date = (datetime.now() - timedelta(days=10)).strftime('%Y%m%d') + df = data_source_manager.tushare_api.index_daily( + ts_code='000001.SH', + start_date=start_date, + end_date=end_date + ) + + if df is not None and not df.empty: + row = df.iloc[0] + change_pct = row.get('pct_chg', 0) + + print(f" [Tushare] ✅ 成功获取大盘指数涨跌幅: {change_pct}%") + return { + "index_name": "上证指数", + "change_percent": change_pct + } + else: + print(f" [Tushare] ❌ 未获取到大盘指数,尝试备用数据源") + except Exception as te: + print(f" [Tushare] ❌ 获取失败: {te}") + + # tushare失败,回退akshare + print(f" [Akshare] 正在获取大盘指数数据(备用数据源)...") df = ak.stock_zh_index_spot_em(symbol="上证系列指数") if df is not None and not df.empty: # 查找上证指数(代码为000001) @@ -470,30 +500,6 @@ class MarketSentimentDataFetcher: } except Exception as e: print(f" [Akshare] ❌ 获取大盘指数失败: {e}") - - # akshare失败,尝试tushare - if data_source_manager.tushare_available: - try: - print(f" [Tushare] 正在获取大盘指数数据(备用数据源)...") - - # 获取上证指数数据 - df = data_source_manager.tushare_api.index_daily( - ts_code='000001.SH', - start_date=datetime.now().strftime('%Y%m%d'), - end_date=datetime.now().strftime('%Y%m%d') - ) - - if df is not None and not df.empty: - row = df.iloc[0] - change_pct = row.get('pct_chg', 0) - - print(f" [Tushare] ✅ 成功获取大盘指数涨跌幅: {change_pct}%") - return { - "index_name": "上证指数", - "change_percent": change_pct - } - except Exception as te: - print(f" [Tushare] ❌ 获取失败: {te}") return None @@ -503,19 +509,50 @@ class MarketSentimentDataFetcher: # 获取今日涨停和跌停统计 today = datetime.now().strftime('%Y%m%d') - # 获取涨停股票 - try: - limit_up_df = ak.stock_zt_pool_em(date=today) - limit_up_count = len(limit_up_df) if limit_up_df is not None and not limit_up_df.empty else 0 - except: - limit_up_count = 0 + limit_up_count = 0 + limit_down_count = 0 - # 获取跌停股票 - try: - limit_down_df = ak.stock_zt_pool_dtgc_em(date=today) - limit_down_count = len(limit_down_df) if limit_down_df is not None and not limit_down_df.empty else 0 - except: - limit_down_count = 0 + # 优先使用tushare的涨跌停列表 + if data_source_manager.tushare_available: + try: + print(f" [Tushare] 正在获取涨跌停数据(主要数据源)...") + df_ll = data_source_manager.tushare_api.limit_list_d(trade_date=today) + if df_ll is not None and not df_ll.empty: + if 'limit_type' in df_ll.columns: + limit_up_count = int((df_ll['limit_type'].fillna('') == 'U').sum()) + limit_down_count = int((df_ll['limit_type'].fillna('') == 'D').sum()) + print(f" [Tushare] ✅ 成功获取涨跌停: 涨停{limit_up_count} / 跌停{limit_down_count}") + elif 'pct_chg' in df_ll.columns: + limit_up_count = int((df_ll['pct_chg'] >= 9.5).sum()) + limit_down_count = int((df_ll['pct_chg'] <= -9.5).sum()) + print(f" [Tushare] ✅ 成功获取涨跌停: 涨停{limit_up_count} / 跌停{limit_down_count}") + elif '涨跌幅' in df_ll.columns: + limit_up_count = int((df_ll['涨跌幅'] >= 9.5).sum()) + limit_down_count = int((df_ll['涨跌幅'] <= -9.5).sum()) + print(f" [Tushare] ✅ 成功获取涨跌停: 涨停{limit_up_count} / 跌停{limit_down_count}") + else: + # 无法识别的列结构,按0处理并回退akshare + print(f" [Tushare] ⚠ 涨跌停返回列无法识别: {list(df_ll.columns)[:10]},尝试备用数据源") + else: + print(f" [Tushare] ❌ 未获取到涨跌停数据,尝试备用数据源") + except Exception as e: + print(f" [Tushare] ❌ 获取涨跌停数据失败: {e}") + + # tushare不可用或失败时,回退akshare + if limit_up_count == 0 and limit_down_count == 0: + # 获取涨停股票 + try: + limit_up_df = ak.stock_zt_pool_em(date=today) + limit_up_count = len(limit_up_df) if limit_up_df is not None and not limit_up_df.empty else 0 + except: + limit_up_count = 0 + + # 获取跌停股票 + try: + limit_down_df = ak.stock_zt_pool_dtgc_em(date=today) + limit_down_count = len(limit_down_df) if limit_down_df is not None and not limit_down_df.empty else 0 + except: + limit_down_count = 0 # 计算涨跌停比例 if limit_up_count + limit_down_count > 0: @@ -549,6 +586,44 @@ class MarketSentimentDataFetcher: def _get_margin_trading_data(self, symbol): """获取融资融券数据""" try: + # 优先使用tushare的个股融资融券明细 + if data_source_manager.tushare_available: + try: + print(f" [Tushare] 正在获取融资融券数据(主要数据源)...") + ts_code = data_source_manager._convert_to_ts_code(symbol) + end_date = datetime.now().strftime('%Y%m%d') + start_date = (datetime.now() - timedelta(days=15)).strftime('%Y%m%d') + df = data_source_manager.tushare_api.margin_detail( + ts_code=ts_code, + start_date=start_date, + end_date=end_date + ) + if df is not None and not df.empty: + latest = df.iloc[0] + margin_balance = latest.get('rzye', 0) or 0 + short_balance = latest.get('rqye', 0) or 0 + + # 解读融资融券 + interpretation = [] + if margin_balance > short_balance * 10: + interpretation.append("融资余额远大于融券余额,投资者看多情绪强") + elif margin_balance > short_balance * 3: + interpretation.append("融资余额大于融券余额,投资者偏看多") + else: + interpretation.append("融资融券相对平衡") + + print(f" [Tushare] ✅ 成功获取融资融券数据") + return { + "margin_balance": margin_balance, + "short_balance": short_balance, + "interpretation": interpretation, + "date": str(latest.get('trade_date', datetime.now().strftime('%Y-%m-%d'))) + } + else: + print(f" [Tushare] ❌ 未获取到融资融券数据,尝试备用数据源") + except Exception as e: + print(f" [Tushare] ❌ 获取融资融券数据失败: {e}") + # 获取个股融资融券数据(尝试多个API) try: # 方法1:获取沪深融资融券明细 @@ -762,4 +837,3 @@ if __name__ == "__main__": print(formatted_text) else: print(f"\n获取失败: {sentiment_data.get('error', '未知错误')}") - diff --git a/news_announcement_data.py b/news_announcement_data.py index 6163b0b..bba9390 100644 --- a/news_announcement_data.py +++ b/news_announcement_data.py @@ -9,6 +9,7 @@ import sys import io import warnings from datetime import datetime +from http_timeout import call_with_timeout warnings.filterwarnings('ignore') @@ -100,7 +101,7 @@ class NewsAnnouncementDataFetcher: print(f" 使用问财查询: {query}") # 使用pywencai查询 - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None: print(f" 问财查询返回None") @@ -191,7 +192,7 @@ class NewsAnnouncementDataFetcher: print(f" 使用问财查询: {query}") # 使用pywencai查询 - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None: print(f" 问财查询返回None") @@ -342,4 +343,3 @@ if __name__ == "__main__": print(formatted_text) else: print(f"\n获取失败: {data.get('error', '未知错误')}") - diff --git a/profit_growth_selector.py b/profit_growth_selector.py index 530f313..f900b90 100644 --- a/profit_growth_selector.py +++ b/profit_growth_selector.py @@ -8,6 +8,7 @@ import logging from typing import Tuple, Optional import pandas as pd +from http_timeout import call_with_timeout class ProfitGrowthSelector: @@ -50,7 +51,7 @@ class ProfitGrowthSelector: self.logger.info(f"开始执行净利增长选股,查询条件: {query}") # 调用pywencai - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None or result.empty: self.logger.warning("未获取到符合条件的股票") diff --git a/qstock_news_data.py b/qstock_news_data.py index 1f78a2f..7e32c73 100644 --- a/qstock_news_data.py +++ b/qstock_news_data.py @@ -1,6 +1,6 @@ """ 新闻数据获取模块 -使用akshare获取股票的最新新闻信息(替代qstock) +优先使用tushare,失败时使用akshare获取股票的最新新闻信息 """ import pandas as pd @@ -9,6 +9,7 @@ import io import warnings from datetime import datetime, timedelta import akshare as ak +from data_source_manager import data_source_manager warnings.filterwarnings('ignore') @@ -30,6 +31,9 @@ def _setup_stdout_encoding(): _setup_stdout_encoding() +# 记录tushare news接口是否无权限(避免每次分析都重复请求失败) +_tushare_news_unavailable = False + class QStockNewsDataFetcher: """新闻数据获取类(使用akshare作为数据源)""" @@ -37,7 +41,7 @@ class QStockNewsDataFetcher: def __init__(self): self.max_items = 30 # 最多获取的新闻数量 self.available = True - print("✓ 新闻数据获取器初始化成功(akshare数据源)") + print("✓ 新闻数据获取器初始化成功(tushare优先/akshare备用)") def get_stock_news(self, symbol): """ @@ -67,7 +71,7 @@ class QStockNewsDataFetcher: try: # 获取新闻数据 - print(f"📰 正在使用qstock获取 {symbol} 的最新新闻...") + print(f"📰 正在获取 {symbol} 的最新新闻...") news_data = self._get_news_data(symbol) if news_data: @@ -89,9 +93,20 @@ class QStockNewsDataFetcher: return symbol.isdigit() and len(symbol) == 6 def _get_news_data(self, symbol): - """获取新闻数据(使用akshare)""" + """获取新闻数据(优先tushare,失败时使用akshare)""" try: - print(f" 使用 akshare 获取新闻...") + # 优先使用tushare新闻接口 + tushare_items = self._get_news_from_tushare(symbol) + if tushare_items: + print(f" ✓ 从tushare获取到 {len(tushare_items)} 条相关新闻") + return { + "items": tushare_items, + "count": len(tushare_items), + "query_time": datetime.now().strftime('%Y-%m-%d %H:%M:%S'), + "date_range": "最近新闻" + } + + print(f" 使用 akshare 获取新闻(备用数据源)...") news_items = [] @@ -223,6 +238,69 @@ class QStockNewsDataFetcher: traceback.print_exc() return None + def _get_news_from_tushare(self, symbol): + """从tushare获取个股新闻(按股票名称/代码过滤)""" + global _tushare_news_unavailable + try: + if _tushare_news_unavailable or not data_source_manager.tushare_available: + return None + + # 获取股票名称 + stock_name = None + try: + basic = data_source_manager.get_stock_basic_info(symbol) + if basic and basic.get('name') and basic['name'] != '未知': + stock_name = basic['name'] + except Exception as e: + print(f" 获取股票名称失败: {e}") + + # 查询最近7天的全市场新闻(东方财富源) + end_date = datetime.now().strftime('%Y-%m-%d') + start_date = (datetime.now() - timedelta(days=7)).strftime('%Y-%m-%d') + df = data_source_manager.tushare_api.news( + src='eastmoney', + start_date=start_date, + end_date=end_date + ) + if df is None or df.empty: + return None + + # 按股票代码或名称过滤 + mask = df['title'].str.contains(symbol, na=False) | df['title'].str.contains(stock_name, na=False) if stock_name else df['title'].str.contains(symbol, na=False) + if 'content' in df.columns: + mask = mask | df['content'].str.contains(symbol, na=False) + if stock_name: + mask = mask | df['content'].str.contains(stock_name, na=False) + + df_filtered = df[mask] + if df_filtered.empty: + return None + + news_items = [] + for _, row in df_filtered.head(self.max_items).iterrows(): + item = {'source': 'tushare-东方财富'} + for col in ['title', 'content', 'pub_time']: + if col in df_filtered.columns: + value = row.get(col) + if value is None or (isinstance(value, float) and pd.isna(value)): + continue + try: + item[col] = str(value) + except: + item[col] = "无法解析" + if len(item) > 1: + news_items.append(item) + return news_items or None + + except Exception as e: + error_msg = str(e) + if "权限" in error_msg or "积分" in error_msg: + _tushare_news_unavailable = True + print(" ⚠ tushare news 接口需要较高积分,当前账号无权限,已自动使用 akshare 获取新闻") + else: + print(f" ⚠ 从tushare获取新闻失败: {error_msg}") + return None + def format_news_for_ai(self, data): """ 将新闻数据格式化为适合AI阅读的文本 @@ -236,7 +314,7 @@ class QStockNewsDataFetcher: if data.get("news_data"): news_data = data["news_data"] text_parts.append(f""" -【最新新闻 - akshare数据源】 +【最新新闻 - tushare/akshare自动切换】 查询时间:{news_data.get('query_time', 'N/A')} 时间范围:{news_data.get('date_range', 'N/A')} 新闻数量:{news_data.get('count', 0)}条 @@ -303,4 +381,3 @@ if __name__ == "__main__": print(f"\n获取失败: {data.get('error', '未知错误')}") print("\n") - diff --git a/quarterly_report_data.py b/quarterly_report_data.py index 9874737..b6b52a7 100644 --- a/quarterly_report_data.py +++ b/quarterly_report_data.py @@ -9,6 +9,8 @@ import io import warnings from datetime import datetime import akshare as ak +from http_timeout import call_with_timeout +from data_source_manager import data_source_manager warnings.filterwarnings('ignore') @@ -67,26 +69,34 @@ class QuarterlyReportDataFetcher: try: print(f"📊 正在获取 {symbol} 的季报数据...") - # 获取利润表 - income_data = self._get_income_statement(symbol) + # 获取利润表(优先tushare,失败时回退akshare) + income_data = self._get_income_statement_from_tushare(symbol) + if income_data is None: + income_data = self._get_income_statement(symbol) if income_data: data["income_statement"] = income_data print(f" ✓ 成功获取 {len(income_data.get('data', []))} 期利润表数据") - # 获取资产负债表 - balance_data = self._get_balance_sheet(symbol) + # 获取资产负债表(优先tushare,失败时回退akshare) + balance_data = self._get_balance_sheet_from_tushare(symbol) + if balance_data is None: + balance_data = self._get_balance_sheet(symbol) if balance_data: data["balance_sheet"] = balance_data print(f" ✓ 成功获取 {len(balance_data.get('data', []))} 期资产负债表数据") - # 获取现金流量表 - cash_flow_data = self._get_cash_flow(symbol) + # 获取现金流量表(优先tushare,失败时回退akshare) + cash_flow_data = self._get_cash_flow_from_tushare(symbol) + if cash_flow_data is None: + cash_flow_data = self._get_cash_flow(symbol) if cash_flow_data: data["cash_flow"] = cash_flow_data print(f" ✓ 成功获取 {len(cash_flow_data.get('data', []))} 期现金流量表数据") - # 获取财务指标 - indicators_data = self._get_financial_indicators(symbol) + # 获取财务指标(优先tushare,失败时回退akshare) + indicators_data = self._get_financial_indicators_from_tushare(symbol) + if indicators_data is None: + indicators_data = self._get_financial_indicators(symbol) if indicators_data: data["financial_indicators"] = indicators_data print(f" ✓ 成功获取 {len(indicators_data.get('data', []))} 期财务指标数据") @@ -108,11 +118,160 @@ class QuarterlyReportDataFetcher: """判断是否为中国股票""" return symbol.isdigit() and len(symbol) == 6 + def _convert_ts_records(self, df, field_map, periods): + """将tushare返回的财务表转换为统一的记录结构""" + try: + data_list = [] + for _, row in df.head(periods).iterrows(): + item = {} + for ch_name, ts_name in field_map.items(): + if ts_name in df.columns: + value = row.get(ts_name) + if value is None or (isinstance(value, float) and pd.isna(value)): + continue + try: + item[ch_name] = str(value) + except: + item[ch_name] = "N/A" + if item: + data_list.append(item) + return { + "data": data_list, + "periods": len(data_list), + "columns": list(field_map.keys()), + "query_time": datetime.now().strftime('%Y-%m-%d %H:%M:%S') + } + except Exception as e: + print(f" 转换tushare财务数据异常: {e}") + return None + + def _get_income_statement_from_tushare(self, symbol): + """从tushare获取利润表(优先数据源)""" + try: + if not data_source_manager.tushare_available: + return None + ts_code = data_source_manager._convert_to_ts_code(symbol) + df = data_source_manager.tushare_api.income(ts_code=ts_code) + if df is None or df.empty: + return None + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + field_map = { + '报告期': 'end_date', + '营业总收入': 'total_revenue', + '营业收入': 'revenue', + '营业总成本': 'total_operate_cost', + '营业利润': 'operate_profit', + '利润总额': 'total_profit', + '净利润': 'n_income', + '归属于母公司所有者的净利润': 'n_income_attr_p', + '基本每股收益': 'basic_eps', + '稀释每股收益': 'diluted_eps', + '销售费用': 'sell_exp', + '管理费用': 'admin_exp', + '财务费用': 'fin_exp', + '研发费用': 'rd_exp', + } + result = self._convert_ts_records(df, field_map, self.periods) + if result and result.get('periods'): + print(f" ✓ tushare成功获取 {result['periods']} 期利润表数据") + return result + except Exception as e: + print(f" tushare获取利润表异常: {e}") + return None + + def _get_balance_sheet_from_tushare(self, symbol): + """从tushare获取资产负债表(优先数据源)""" + try: + if not data_source_manager.tushare_available: + return None + ts_code = data_source_manager._convert_to_ts_code(symbol) + df = data_source_manager.tushare_api.balancesheet(ts_code=ts_code) + if df is None or df.empty: + return None + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + field_map = { + '报告期': 'end_date', + '资产总计': 'total_assets', + '流动资产合计': 'total_cur_assets', + '非流动资产合计': 'total_ncur_assets', + '负债合计': 'total_liab', + '流动负债合计': 'total_cur_liab', + '非流动负债合计': 'total_ncur_liab', + '所有者权益合计': 'total_hldr_eqy_inc_min_int', + '归属于母公司股东权益合计': 'total_hldr_eqy_exc_min_int', + } + result = self._convert_ts_records(df, field_map, self.periods) + if result and result.get('periods'): + print(f" ✓ tushare成功获取 {result['periods']} 期资产负债表数据") + return result + except Exception as e: + print(f" tushare获取资产负债表异常: {e}") + return None + + def _get_cash_flow_from_tushare(self, symbol): + """从tushare获取现金流量表(优先数据源)""" + try: + if not data_source_manager.tushare_available: + return None + ts_code = data_source_manager._convert_to_ts_code(symbol) + df = data_source_manager.tushare_api.cashflow(ts_code=ts_code) + if df is None or df.empty: + return None + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + field_map = { + '报告期': 'end_date', + '经营活动产生的现金流量净额': 'n_cashflow_act', + '投资活动产生的现金流量净额': 'n_cashflow_inv_act', + '筹资活动产生的现金流量净额': 'n_cash_flows_fnc_act', + } + result = self._convert_ts_records(df, field_map, self.periods) + if result and result.get('periods'): + print(f" ✓ tushare成功获取 {result['periods']} 期现金流量表数据") + return result + except Exception as e: + print(f" tushare获取现金流量表异常: {e}") + return None + + def _get_financial_indicators_from_tushare(self, symbol): + """从tushare获取财务指标(优先数据源)""" + try: + if not data_source_manager.tushare_available: + return None + ts_code = data_source_manager._convert_to_ts_code(symbol) + df = data_source_manager.tushare_api.fina_indicator(ts_code=ts_code) + if df is None or df.empty: + return None + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + field_map = { + '报告期': 'end_date', + '净资产收益率': 'roe', + '总资产净利率': 'roa', + '销售净利率': 'netprofit_margin', + '销售毛利率': 'grossprofit_margin', + '资产负债率': 'debt_to_assets', + '流动比率': 'current_ratio', + '速动比率': 'quick_ratio', + '应收账款周转率': 'ar_turnover', + '存货周转率': 'inventory_turnover', + '总资产周转率': 'assets_turnover', + '每股收益': 'eps', + '每股净资产': 'bps', + '每股经营现金流': 'cfps', + } + result = self._convert_ts_records(df, field_map, self.periods) + if result and result.get('periods'): + print(f" ✓ tushare成功获取 {result['periods']} 期财务指标数据") + return result + except Exception as e: + print(f" tushare获取财务指标异常: {e}") + return None + def _get_income_statement(self, symbol): """获取利润表数据""" try: # stock_financial_report_sina - 新浪财经季度利润表 - df = ak.stock_financial_report_sina(stock=symbol, symbol="利润表") + df = call_with_timeout(ak.stock_financial_report_sina, timeout=25, + stock=symbol, symbol="利润表") if df is None or df.empty: print(f" 未找到利润表数据") @@ -151,7 +310,8 @@ class QuarterlyReportDataFetcher: """获取资产负债表数据""" try: # stock_financial_report_sina - 新浪财经季度资产负债表 - df = ak.stock_financial_report_sina(stock=symbol, symbol="资产负债表") + df = call_with_timeout(ak.stock_financial_report_sina, timeout=25, + stock=symbol, symbol="资产负债表") if df is None or df.empty: print(f" 未找到资产负债表数据") @@ -190,7 +350,8 @@ class QuarterlyReportDataFetcher: """获取现金流量表数据""" try: # stock_financial_report_sina - 新浪财经季度现金流量表 - df = ak.stock_financial_report_sina(stock=symbol, symbol="现金流量表") + df = call_with_timeout(ak.stock_financial_report_sina, timeout=25, + stock=symbol, symbol="现金流量表") if df is None or df.empty: print(f" 未找到现金流量表数据") @@ -229,7 +390,7 @@ class QuarterlyReportDataFetcher: """获取财务指标数据""" try: # 使用stock_financial_abstract替代已失效的stock_financial_analysis_indicator - df = ak.stock_financial_abstract(symbol=symbol) + df = call_with_timeout(ak.stock_financial_abstract, timeout=25, symbol=symbol) if df is None or df.empty: print(f" 未找到财务指标数据") @@ -292,7 +453,7 @@ class QuarterlyReportDataFetcher: text_parts = [] text_parts.append(f""" -【季度财务报告数据 - akshare数据源】 +【季度财务报告数据 - tushare/akshare自动切换】 股票代码:{data.get('symbol', 'N/A')} 数据期数:最近{self.periods}期季报 @@ -423,4 +584,3 @@ if __name__ == "__main__": print(f"\n获取失败: {data.get('error', '未知错误')}") print("\n") - diff --git a/risk_data_fetcher.py b/risk_data_fetcher.py index bd6af7e..387f936 100644 --- a/risk_data_fetcher.py +++ b/risk_data_fetcher.py @@ -12,6 +12,7 @@ from typing import Dict, Any import time import warnings import os +from http_timeout import call_with_timeout # 屏蔽pywencai的Node.js警告信息(不影响功能) warnings.filterwarnings('ignore', category=DeprecationWarning) @@ -108,7 +109,7 @@ class RiskDataFetcher: query = f"{symbol}限售解禁" # 使用pywencai查询 - response = pywencai.get(query=query, loop=True) + response = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if response is None: return result @@ -171,7 +172,7 @@ class RiskDataFetcher: query = f"{symbol}大股东减持公告" # 使用pywencai查询 - response = pywencai.get(query=query, loop=True) + response = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if response is None: return result @@ -234,7 +235,7 @@ class RiskDataFetcher: query = f"{symbol}近期重要事件" # 使用pywencai查询 - response = pywencai.get(query=query, loop=True) + response = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if response is None: return result @@ -466,4 +467,3 @@ if __name__ == "__main__": if risk_data['data_success']: print("\n格式化的风险数据:") print(fetcher.format_risk_data_for_ai(risk_data)) - diff --git a/sector_strategy_data.py b/sector_strategy_data.py index 0df19c1..61c4c07 100644 --- a/sector_strategy_data.py +++ b/sector_strategy_data.py @@ -1,6 +1,6 @@ """ 智策板块数据采集模块 -使用AKShare获取板块相关数据 +优先使用Tushare,失败时使用AKShare获取板块相关数据 """ import akshare as ak @@ -12,6 +12,7 @@ import logging import os from dotenv import load_dotenv from sector_strategy_db import SectorStrategyDatabase +from data_source_manager import data_source_manager # 加载环境变量 load_dotenv() @@ -254,40 +255,52 @@ class SectorStrategyDataFetcher: except: pass - # 大盘指数 + # 大盘指数(优先tushare,失败回退akshare) + def _get_index_data(index_code, ts_code, name): + if data_source_manager.tushare_available: + try: + end = datetime.now().strftime('%Y%m%d') + start = (datetime.now() - timedelta(days=10)).strftime('%Y%m%d') + df = data_source_manager.tushare_api.index_daily( + ts_code=ts_code, start_date=start, end_date=end + ) + if df is not None and not df.empty: + row = df.iloc[0] + return { + "code": index_code, + "name": name, + "close": row.get('close', 0), + "change_pct": row.get('pct_chg', 0), + "change": row.get('change', 0) + } + print(f" tushare未获取到{name},尝试备用数据源") + except Exception as e: + print(f" tushare获取{name}失败: {e}") + try: + df = self._safe_request(ak.stock_zh_index_spot_em, symbol=name) + if df is not None and not df.empty: + row = df.iloc[0] + return { + "code": index_code, + "name": name, + "close": row.get('最新价', 0), + "change_pct": row.get('涨跌幅', 0), + "change": row.get('涨跌额', 0) + } + except Exception as e: + print(f" akshare获取{name}失败: {e}") + return None + try: - # 上证指数 - df_sh = ak.stock_zh_index_spot_em(symbol="上证指数") - if df_sh is not None and not df_sh.empty: - overview["sh_index"] = { - "code": "000001", - "name": "上证指数", - "close": df_sh.iloc[0].get('最新价', 0), - "change_pct": df_sh.iloc[0].get('涨跌幅', 0), - "change": df_sh.iloc[0].get('涨跌额', 0) - } - - # 深证成指 - df_sz = self._safe_request(ak.stock_zh_index_spot_em, symbol="深证成指") - if df_sz is not None and not df_sz.empty: - overview["sz_index"] = { - "code": "399001", - "name": "深证成指", - "close": df_sz.iloc[0].get('最新价', 0), - "change_pct": df_sz.iloc[0].get('涨跌幅', 0), - "change": df_sz.iloc[0].get('涨跌额', 0) - } - - # 创业板指 - df_cyb = self._safe_request(ak.stock_zh_index_spot_em, symbol="创业板指") - if df_cyb is not None and not df_cyb.empty: - overview["cyb_index"] = { - "code": "399006", - "name": "创业板指", - "close": df_cyb.iloc[0].get('最新价', 0), - "change_pct": df_cyb.iloc[0].get('涨跌幅', 0), - "change": df_cyb.iloc[0].get('涨跌额', 0) - } + sh_index = _get_index_data("000001", "000001.SH", "上证指数") + if sh_index: + overview["sh_index"] = sh_index + sz_index = _get_index_data("399001", "399001.SZ", "深证成指") + if sz_index: + overview["sz_index"] = sz_index + cyb_index = _get_index_data("399006", "399006.SZ", "创业板指") + if cyb_index: + overview["cyb_index"] = cyb_index except: pass @@ -310,7 +323,7 @@ class SectorStrategyDataFetcher: try: import tushare as ts ts.set_token(tushare_token) - self.ts_pro = ts.pro_api() + self.ts_pro = ts.pro_api(timeout=float(os.getenv('TUSHARE_TIMEOUT', '15'))) print(" [Tushare] ✅ 初始化成功") except Exception as e: print(f" [Tushare] 初始化失败: {e}") @@ -757,4 +770,3 @@ if __name__ == "__main__": print(f"\n... (总长度: {len(formatted_text)} 字符)") else: print(f"\n数据采集失败: {data.get('error', '未知错误')}") - diff --git a/small_cap_selector.py b/small_cap_selector.py index bb03c2d..c13b67e 100644 --- a/small_cap_selector.py +++ b/small_cap_selector.py @@ -8,6 +8,7 @@ import logging from typing import Tuple, Optional import pandas as pd +from http_timeout import call_with_timeout class SmallCapSelector: @@ -54,7 +55,7 @@ class SmallCapSelector: self.logger.info(f"开始执行小市值策略选股,查询条件: {query}") # 调用pywencai - result = pywencai.get(query=query, loop=True) + result = call_with_timeout(pywencai.get, timeout=30, query=query, loop=True) if result is None or result.empty: self.logger.warning("未获取到符合条件的股票") diff --git a/smart_monitor.db b/smart_monitor.db index 81035a3..cd483be 100644 Binary files a/smart_monitor.db and b/smart_monitor.db differ diff --git a/smart_monitor_data.py b/smart_monitor_data.py index 3837169..451aad9 100644 --- a/smart_monitor_data.py +++ b/smart_monitor_data.py @@ -54,7 +54,7 @@ class SmartMonitorDataFetcher: try: import tushare as ts ts.set_token(tushare_token) - self.ts_pro = ts.pro_api() + self.ts_pro = ts.pro_api(timeout=float(os.getenv('TUSHARE_TIMEOUT', '15'))) self.logger.info("Tushare备用数据源初始化成功") except Exception as e: self.logger.warning(f"Tushare初始化失败: {e}") @@ -64,7 +64,7 @@ class SmartMonitorDataFetcher: def get_realtime_quote(self, stock_code: str, retry: int = 1) -> Optional[Dict]: """ 获取实时行情(带重试和降级机制) - 优先使用TDX,失败时降级到AKShare,最后降级到Tushare + 优先使用TDX,失败时降级到Tushare,最后降级到AKShare Args: stock_code: 股票代码(如:600519) @@ -82,11 +82,20 @@ class SmartMonitorDataFetcher: if quote: return quote else: - self.logger.warning(f"TDX获取失败 {stock_code},尝试降级到AKShare") + self.logger.warning(f"TDX获取失败 {stock_code},尝试降级到Tushare") except Exception as e: - self.logger.warning(f"TDX获取异常 {stock_code}: {e},尝试降级到AKShare") + self.logger.warning(f"TDX获取异常 {stock_code}: {e},尝试降级到Tushare") - # 方法2: 组合使用AKShare分钟行情 + 基本信息 + # 方法2: 降级到Tushare(优先于AKShare) + if self.ts_pro: + quote = self._get_realtime_quote_from_tushare(stock_code) + if quote: + return quote + self.logger.warning(f"Tushare获取失败 {stock_code},尝试降级到AKShare") + else: + self.logger.warning(f"未配置Tushare,尝试使用AKShare") + + # 方法3: 组合使用AKShare分钟行情 + 基本信息 for attempt in range(retry): try: # 1.1 获取股票基本信息(名称) @@ -167,18 +176,14 @@ class SmartMonitorDataFetcher: else: self.logger.warning(f"AKShare获取失败 {stock_code}(已重试{retry}次),尝试降级") - # 降级到Tushare - if self.ts_pro: - self.logger.info(f"降级到Tushare获取 {stock_code}...") - return self._get_realtime_quote_from_tushare(stock_code) - else: - self.logger.error(f"AKShare失败且未配置Tushare,无法获取 {stock_code} 行情") - return None + # 所有数据源都失败 + self.logger.error(f"所有数据源都无法获取 {stock_code} 行情") + return None def get_technical_indicators(self, stock_code: str, period: str = 'daily', retry: int = 1) -> Optional[Dict]: """ 计算技术指标(带降级机制) - 优先使用TDX,失败时降级到AKShare,最后降级到Tushare + 优先使用TDX,失败时降级到Tushare,最后降级到AKShare Args: stock_code: 股票代码 @@ -197,11 +202,21 @@ class SmartMonitorDataFetcher: if indicators: return indicators else: - self.logger.warning(f"TDX计算技术指标失败 {stock_code},尝试降级到AKShare") + self.logger.warning(f"TDX计算技术指标失败 {stock_code},尝试降级到Tushare") except Exception as e: - self.logger.warning(f"TDX计算技术指标异常 {stock_code}: {e},尝试降级到AKShare") + self.logger.warning(f"TDX计算技术指标异常 {stock_code}: {e},尝试降级到Tushare") - # 方法2: 尝试使用AKShare + # 方法2: 降级到Tushare(优先于AKShare) + if self.ts_pro: + self.logger.info(f"降级到Tushare获取 {stock_code} 历史数据...") + indicators = self._get_technical_indicators_from_tushare(stock_code, period) + if indicators: + return indicators + self.logger.warning(f"Tushare获取技术指标失败 {stock_code},尝试降级到AKShare") + else: + self.logger.warning(f"未配置Tushare,尝试使用AKShare") + + # 方法3: 尝试使用AKShare for attempt in range(retry): try: # 获取历史数据(最近200个交易日,用于计算指标) @@ -237,13 +252,9 @@ class SmartMonitorDataFetcher: self.logger.warning(f"AKShare获取历史数据失败 {stock_code}(已重试{retry}次),尝试降级到Tushare") break - # 方法3: 降级到Tushare - if self.ts_pro: - self.logger.info(f"降级到Tushare获取 {stock_code} 历史数据...") - return self._get_technical_indicators_from_tushare(stock_code, period) - else: - self.logger.error(f"AKShare失败且未配置Tushare,无法获取 {stock_code} 技术指标") - return None + # 所有数据源都失败 + self.logger.error(f"所有数据源都无法获取 {stock_code} 技术指标") + return None def _calculate_all_indicators(self, df: pd.DataFrame, stock_code: str) -> Optional[Dict]: """ @@ -439,6 +450,17 @@ class SmartMonitorDataFetcher: """ import time + # 优先使用Tushare个股资金流向 + if self.ts_pro: + try: + result = self._get_main_force_from_tushare(stock_code) + if result: + self.logger.info(f"✅ Tushare成功获取 {stock_code} 资金流向(主要数据源)") + return result + self.logger.warning(f"Tushare未获取到资金流向 {stock_code},尝试AKShare") + except Exception as e: + self.logger.warning(f"Tushare获取资金流向失败 {stock_code}: {e}") + for attempt in range(retry): try: # 获取个股资金流(新版AKShare API参数调整) @@ -495,12 +517,9 @@ class SmartMonitorDataFetcher: self.logger.warning(f"AKShare获取资金流向失败 {stock_code}(已重试{retry}次),尝试降级到Tushare") break - # 降级到Tushare - if self.ts_pro: - return self._get_main_force_from_tushare(stock_code) - else: - self.logger.error(f"AKShare失败且未配置Tushare,无法获取 {stock_code} 资金流向") - return None + # 所有数据源都失败 + self.logger.error(f"所有数据源都无法获取 {stock_code} 资金流向") + return None def get_comprehensive_data(self, stock_code: str) -> Dict: """ @@ -822,4 +841,3 @@ if __name__ == '__main__': print("\n主力资金:") print(f" 主力净额: {data['main_force']['main_net']:.2f}万") print(f" 主力动向: {data['main_force']['trend']}") - diff --git a/smart_monitor_kline.py b/smart_monitor_kline.py index 96df3d5..90dd91d 100644 --- a/smart_monitor_kline.py +++ b/smart_monitor_kline.py @@ -333,15 +333,26 @@ class SmartMonitorKline: self.logger.info(f"✅ TDX获取K线数据成功 {stock_code},共{len(df)}条") return df else: - self.logger.warning(f"TDX未返回K线数据 {stock_code},尝试降级到AKShare") + self.logger.warning(f"TDX未返回K线数据 {stock_code},尝试降级到Tushare") except Exception as e: - self.logger.warning(f"TDX获取K线数据失败 {stock_code}: {type(e).__name__}, 尝试降级到AKShare") + self.logger.warning(f"TDX获取K线数据失败 {stock_code}: {type(e).__name__}, 尝试降级到Tushare") # 计算日期范围 end_date = datetime.now().strftime('%Y%m%d') start_date = (datetime.now() - timedelta(days=days + 30)).strftime('%Y%m%d') # 多取30天以确保足够数据 - # 方法2: 尝试使用AKShare获取(只尝试1次,避免IP封禁) + # 方法2: 降级到Tushare(优先于AKShare) + if data_fetcher and data_fetcher.ts_pro: + self.logger.info(f"降级使用Tushare获取K线数据 {stock_code}") + df = self._get_kline_from_tushare(stock_code, days, data_fetcher.ts_pro) + if df is not None and not df.empty: + self.logger.info(f"✅ Tushare获取K线数据成功 {stock_code},共{len(df)}条") + return df + self.logger.warning(f"Tushare未返回K线数据 {stock_code},尝试降级到AKShare") + else: + self.logger.warning(f"未配置Tushare,尝试使用AKShare") + + # 方法3: 尝试使用AKShare获取(只尝试1次,避免IP封禁) try: import akshare as ak df = ak.stock_zh_a_hist( @@ -358,17 +369,9 @@ class SmartMonitorKline: self.logger.info(f"✅ AKShare获取K线数据成功 {stock_code},共{len(df)}条") return df else: - self.logger.warning(f"AKShare未返回K线数据 {stock_code},尝试降级到Tushare") + self.logger.warning(f"AKShare未返回K线数据 {stock_code}") except Exception as e: - self.logger.warning(f"AKShare获取K线数据失败 {stock_code}: {type(e).__name__}, 尝试降级到Tushare") - - # 方法3: 降级到Tushare - if data_fetcher and data_fetcher.ts_pro: - self.logger.info(f"降级使用Tushare获取K线数据 {stock_code}") - df = self._get_kline_from_tushare(stock_code, days, data_fetcher.ts_pro) - if df is not None and not df.empty: - self.logger.info(f"✅ Tushare获取K线数据成功 {stock_code},共{len(df)}条") - return df + self.logger.warning(f"AKShare获取K线数据失败 {stock_code}: {type(e).__name__}") self.logger.error(f"所有数据源都无法获取K线数据 {stock_code}") return None @@ -404,12 +407,15 @@ class SmartMonitorKline: end_date = datetime.now().strftime('%Y%m%d') start_date = (datetime.now() - timedelta(days=days + 60)).strftime('%Y%m%d') - # 获取日K线数据(前复权) - df = ts_pro.daily( + # 获取日K线数据(前复权,daily接口不支持adj参数,需用pro_bar) + import tushare as ts + df = ts.pro_bar( + api=ts_pro, ts_code=ts_code, start_date=start_date, end_date=end_date, - adj='qfq' + adj='qfq', + retry_count=1 ) if df is None or df.empty: @@ -493,4 +499,3 @@ if __name__ == '__main__': print("K线图已保存到 test_kline.html") else: print("获取K线数据失败") - diff --git a/stock_analysis.db b/stock_analysis.db index 1de4e1c..e640d4c 100644 Binary files a/stock_analysis.db and b/stock_analysis.db differ diff --git a/stock_data.py b/stock_data.py index c5b3268..fad29dd 100644 --- a/stock_data.py +++ b/stock_data.py @@ -8,10 +8,46 @@ import requests import json import pywencai from data_source_manager import data_source_manager +from http_timeout import call_with_timeout class StockDataFetcher: """股票数据获取类""" + # tushare财务字段 -> 中文映射(供AI分析展示) + INCOME_TS_MAP = { + '报告期': 'end_date', + '营业总收入': 'total_revenue', + '营业收入': 'revenue', + '营业总成本': 'total_operate_cost', + '营业利润': 'operate_profit', + '利润总额': 'total_profit', + '净利润': 'n_income', + '归属于母公司所有者的净利润': 'n_income_attr_p', + '基本每股收益': 'basic_eps', + '稀释每股收益': 'diluted_eps', + '销售费用': 'sell_exp', + '管理费用': 'admin_exp', + '财务费用': 'fin_exp', + '研发费用': 'rd_exp', + } + BALANCE_TS_MAP = { + '报告期': 'end_date', + '资产总计': 'total_assets', + '流动资产合计': 'total_cur_assets', + '非流动资产合计': 'total_ncur_assets', + '负债合计': 'total_liab', + '流动负债合计': 'total_cur_liab', + '非流动负债合计': 'total_ncur_liab', + '所有者权益合计': 'total_hldr_eqy_inc_min_int', + '归属于母公司股东权益合计': 'total_hldr_eqy_exc_min_int', + } + CASHFLOW_TS_MAP = { + '报告期': 'end_date', + '经营活动产生的现金流量净额': 'n_cashflow_act', + '投资活动产生的现金流量净额': 'n_cashflow_inv_act', + '筹资活动产生的现金流量净额': 'n_cash_flows_fnc_act', + } + def __init__(self): self.data = None self.info = None @@ -89,57 +125,70 @@ class StockDataFetcher: if basic_info: info.update(basic_info) - # 方法1: 尝试获取个股详细信息(akshare) - try: - stock_info = ak.stock_individual_info_em(symbol=symbol) - if stock_info is not None and not stock_info.empty: - for _, row in stock_info.iterrows(): - key = row['item'] - value = row['value'] - - if key == '股票简称': - info['name'] = value - elif key == '总市值': - try: - if value and value != '-': - info['market_cap'] = float(value) - except: - pass - elif key == '市盈率-动态': - try: - if value and value != '-': - pe_value = float(value) - if 0 < pe_value <= 1000: - info['pe_ratio'] = pe_value - except: - pass - elif key == '市净率': - try: - if value and value != '-': - pb_value = float(value) - if 0 < pb_value <= 100: - info['pb_ratio'] = pb_value - except: - pass - except Exception as e: - print(f"[Akshare] 获取个股详细信息失败: {e}") - # 如果akshare失败,尝试从tushare获取 - if self.data_source_manager.tushare_available and info['name'] == '未知': - print(f"[Tushare] 尝试获取基本信息(tushare)...") + # 方法1: 获取详细估值信息(优先tushare,失败时回退akshare) + if (info.get('name') == '未知' or info.get('pe_ratio') == 'N/A' or + info.get('pb_ratio') == 'N/A' or info.get('market_cap') == 'N/A'): + # 优先使用tushare daily_basic(一次获取PE/PB/市值) + if self.data_source_manager.tushare_available: try: + print(f"[Tushare] 正在获取 {symbol} 的估值信息(主要数据源)...") ts_code = self.data_source_manager._convert_to_ts_code(symbol) df = self.data_source_manager.tushare_api.daily_basic( ts_code=ts_code, - trade_date=datetime.now().strftime('%Y%m%d') + start_date=(datetime.now() - timedelta(days=10)).strftime('%Y%m%d'), + end_date=datetime.now().strftime('%Y%m%d') ) if df is not None and not df.empty: row = df.iloc[0] - info['pe_ratio'] = row.get('pe', 'N/A') - info['pb_ratio'] = row.get('pb', 'N/A') - info['market_cap'] = row.get('total_mv', 'N/A') - print(f"[Tushare] ✅ 成功获取部分信息") - except Exception as te: - print(f"[Tushare] ❌ 获取失败: {te}") + if info.get('pe_ratio') == 'N/A' and 'pe' in df.columns: + info['pe_ratio'] = row.get('pe', 'N/A') + if info.get('pb_ratio') == 'N/A' and 'pb' in df.columns: + info['pb_ratio'] = row.get('pb', 'N/A') + if info.get('market_cap') == 'N/A' and 'total_mv' in df.columns: + info['market_cap'] = row.get('total_mv', 'N/A') + print(f"[Tushare] ✅ 成功获取估值信息") + else: + print(f"[Tushare] ❌ 未获取到估值信息,尝试备用数据源") + except Exception as e: + print(f"[Tushare] ❌ 获取估值信息失败: {e}") + + # tushare未获取到时,回退akshare + if (info.get('name') == '未知' or info.get('pe_ratio') == 'N/A' or + info.get('pb_ratio') == 'N/A' or info.get('market_cap') == 'N/A'): + try: + print(f"[Akshare] 正在获取 {symbol} 的详细信息(备用数据源)...") + stock_info = ak.stock_individual_info_em(symbol=symbol) + if stock_info is not None and not stock_info.empty: + for _, row in stock_info.iterrows(): + key = row['item'] + value = row['value'] + + if key == '股票简称': + info['name'] = value + elif key == '总市值': + try: + if value and value != '-': + info['market_cap'] = float(value) + except: + pass + elif key == '市盈率-动态': + try: + if value and value != '-': + pe_value = float(value) + if 0 < pe_value <= 1000: + info['pe_ratio'] = pe_value + except: + pass + elif key == '市净率': + try: + if value and value != '-': + pb_value = float(value) + if 0 < pb_value <= 100: + info['pb_ratio'] = pb_value + except: + pass + except Exception as e: + print(f"[Akshare] 获取个股详细信息失败: {e}") # 方法2: 尝试获取历史价格和涨跌幅(如果网络允许) # try: @@ -204,7 +253,10 @@ class StockDataFetcher: # 方法3: 使用百度估值数据获取市盈率和市净率 if info['pe_ratio'] == 'N/A': try: - pe_data = ak.stock_zh_valuation_baidu(symbol=symbol, indicator="市盈率(TTM)") + pe_data = call_with_timeout( + ak.stock_zh_valuation_baidu, timeout=15, + symbol=symbol, indicator="市盈率(TTM)" + ) if pe_data is not None and not pe_data.empty: latest_pe = pe_data.iloc[-1]['value'] if latest_pe and latest_pe != '-': @@ -216,7 +268,10 @@ class StockDataFetcher: if info['pb_ratio'] == 'N/A': try: - pb_data = ak.stock_zh_valuation_baidu(symbol=symbol, indicator="市净率") + pb_data = call_with_timeout( + ak.stock_zh_valuation_baidu, timeout=15, + symbol=symbol, indicator="市净率" + ) if pb_data is not None and not pb_data.empty: latest_pb = pb_data.iloc[-1]['value'] if latest_pb and latest_pb != '-': @@ -262,6 +317,32 @@ class StockDataFetcher: "exchange": "香港交易所" } + # 优先使用tushare(hk_basic + hk_daily) + if self.data_source_manager.tushare_available: + try: + print(f"[Tushare] 正在获取港股信息(主要数据源)...") + ts_code = f"{hk_code}.HK" + bdf = self.data_source_manager.tushare_api.hk_basic(ts_code=ts_code) + if bdf is not None and not bdf.empty and 'name' in bdf.columns: + info['name'] = bdf.iloc[0].get('name', '未知') + + end = datetime.now().strftime('%Y%m%d') + start = (datetime.now() - timedelta(days=10)).strftime('%Y%m%d') + hdf = self.data_source_manager.tushare_api.hk_daily( + ts_code=ts_code, start_date=start, end_date=end + ) + if hdf is not None and not hdf.empty: + latest = hdf.iloc[0] + info['current_price'] = latest.get('close', 'N/A') + info['change_percent'] = latest.get('pct_chg', 'N/A') + + if info['current_price'] != 'N/A': + print(f"[Tushare] ✅ 成功获取港股信息") + return info + print(f"[Tushare] ❌ 未获取到港股行情,尝试备用数据源") + except Exception as e: + print(f"[Tushare] 获取港股信息失败: {e}") + # 方法1: 获取港股实时行情 try: # 使用akshare获取港股实时数据 @@ -331,12 +412,7 @@ class StockDataFetcher: def _get_us_stock_info(self, symbol): """获取美股基本信息""" - import time - try: - # 添加延迟避免频率限制 - time.sleep(1) - ticker = yf.Ticker(symbol) # 先尝试获取历史数据(通常更稳定) @@ -487,7 +563,36 @@ class StockDataFetcher: else: start_date = (datetime.now() - timedelta(days=365)).strftime('%Y%m%d') - # 获取港股历史数据 + # 优先使用tushare(hk_daily) + if self.data_source_manager.tushare_available: + try: + print(f"[Tushare] 正在获取港股历史数据(主要数据源)...") + ts_code = f"{hk_code}.HK" + df = self.data_source_manager.tushare_api.hk_daily( + ts_code=ts_code, + start_date=start_date, + end_date=end_date + ) + if df is not None and not df.empty: + df = df.rename(columns={ + 'trade_date': 'Date', + 'open': 'Open', + 'high': 'High', + 'low': 'Low', + 'close': 'Close', + 'vol': 'Volume' + }) + df['Date'] = pd.to_datetime(df['Date']) + df = df.sort_values('Date') + df.set_index('Date', inplace=True) + print(f"[Tushare] ✅ 成功获取港股历史数据") + return df + else: + print(f"[Tushare] ❌ 未获取到港股历史数据,尝试备用数据源") + except Exception as e: + print(f"[Tushare] 获取港股历史数据失败: {e}") + + # 回退akshare df = ak.stock_hk_hist(symbol=hk_code, period="daily", start_date=start_date, end_date=end_date, adjust="qfq") @@ -600,6 +705,112 @@ class StockDataFetcher: except Exception as e: return {"error": f"获取财务数据失败: {str(e)}"} + def _convert_ts_financial_records(self, df, field_map, limit=8): + """将tushare财务表转换为统一的记录列表""" + try: + data_list = [] + for _, row in df.head(limit).iterrows(): + item = {} + for ch_name, ts_name in field_map.items(): + if ts_name in df.columns: + value = row.get(ts_name) + if value is None or (isinstance(value, float) and pd.isna(value)): + continue + try: + item[ch_name] = str(value) + except: + item[ch_name] = "N/A" + if item: + data_list.append(item) + return data_list + except Exception as e: + print(f"转换tushare财务数据异常: {e}") + return None + + def _convert_ts_ratios(self, df): + """将tushare财务指标转换为AI使用的比率字典""" + try: + df = df.sort_values('end_date', ascending=False) + if df.empty: + return {} + row = df.iloc[0] + mapping = { + '净资产收益率(ROE)': 'roe', + '总资产报酬率(ROA)': 'roa', + '销售毛利率': 'grossprofit_margin', + '销售净利率': 'netprofit_margin', + '资产负债率': 'debt_to_assets', + '流动比率': 'current_ratio', + '速动比率': 'quick_ratio', + '存货周转率': 'inventory_turnover', + '应收账款周转率': 'ar_turnover', + '总资产周转率': 'assets_turnover', + '营业收入同比增长': 'or_yoy', + '净利润同比增长': 'netprofit_yoy', + 'EPS': 'eps', + } + ratios = {'报告期': str(row.get('end_date', 'N/A'))} + for ch_name, ts_name in mapping.items(): + if ts_name in df.columns: + value = row.get(ts_name) + if value is None or (isinstance(value, float) and pd.isna(value)): + ratios[ch_name] = "N/A" + else: + try: + ratios[ch_name] = str(value) + except: + ratios[ch_name] = "N/A" + return ratios + except Exception as e: + print(f"转换tushare财务指标异常: {e}") + return {} + + def _get_financial_data_from_tushare(self, symbol): + """优先从tushare获取财务三表与财务指标""" + try: + if not self.data_source_manager.tushare_available: + return None + ts_code = self.data_source_manager._convert_to_ts_code(symbol) + result = {} + + # 利润表 + df = self.data_source_manager.tushare_api.income(ts_code=ts_code) + if df is not None and not df.empty: + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + records = self._convert_ts_financial_records(df, self.INCOME_TS_MAP) + if records: + result["income_statement"] = records + + # 资产负债表 + df = self.data_source_manager.tushare_api.balancesheet(ts_code=ts_code) + if df is not None and not df.empty: + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + records = self._convert_ts_financial_records(df, self.BALANCE_TS_MAP) + if records: + result["balance_sheet"] = records + + # 现金流量表 + df = self.data_source_manager.tushare_api.cashflow(ts_code=ts_code) + if df is not None and not df.empty: + df = df.drop_duplicates(subset=['end_date']).sort_values('end_date', ascending=False) + records = self._convert_ts_financial_records(df, self.CASHFLOW_TS_MAP) + if records: + result["cash_flow"] = records + + # 财务指标 + df = self.data_source_manager.tushare_api.fina_indicator(ts_code=ts_code) + if df is not None and not df.empty: + ratios = self._convert_ts_ratios(df) + if ratios: + result["financial_ratios"] = ratios + + if result: + print(f"[Tushare] ✅ 成功获取财务数据(主要数据源)") + return result or None + except Exception as e: + print(f"[Tushare] ❌ 获取财务数据失败: {e}") + return None + def _get_chinese_financial_data(self, symbol): """获取中国股票财务数据""" financial_data = { @@ -612,68 +823,77 @@ class StockDataFetcher: } try: - # 1. 获取资产负债表 - try: - balance_sheet = ak.stock_financial_abstract_ths(symbol=symbol, indicator="资产负债表") - if balance_sheet is not None and not balance_sheet.empty: - financial_data["balance_sheet"] = balance_sheet.head(8).to_dict('records') - except Exception as e: - print(f"获取资产负债表失败: {e}") + # 0. 优先使用tushare获取财务三表与财务指标 + ts_financial = self._get_financial_data_from_tushare(symbol) + if ts_financial: + financial_data.update(ts_financial) - # 2. 获取利润表 - try: - income_statement = ak.stock_financial_abstract_ths(symbol=symbol, indicator="利润表") - if income_statement is not None and not income_statement.empty: - financial_data["income_statement"] = income_statement.head(8).to_dict('records') - except Exception as e: - print(f"获取利润表失败: {e}") + # 1. 获取资产负债表(tushare未获取到时回退akshare) + if financial_data["balance_sheet"] is None: + try: + balance_sheet = ak.stock_financial_abstract_ths(symbol=symbol, indicator="资产负债表") + if balance_sheet is not None and not balance_sheet.empty: + financial_data["balance_sheet"] = balance_sheet.head(8).to_dict('records') + except Exception as e: + print(f"获取资产负债表失败: {e}") - # 3. 获取现金流量表 - try: - cash_flow = ak.stock_financial_abstract_ths(symbol=symbol, indicator="现金流量表") - if cash_flow is not None and not cash_flow.empty: - financial_data["cash_flow"] = cash_flow.head(8).to_dict('records') - except Exception as e: - print(f"获取现金流量表失败: {e}") + # 2. 获取利润表(tushare未获取到时回退akshare) + if financial_data["income_statement"] is None: + try: + income_statement = ak.stock_financial_abstract_ths(symbol=symbol, indicator="利润表") + if income_statement is not None and not income_statement.empty: + financial_data["income_statement"] = income_statement.head(8).to_dict('records') + except Exception as e: + print(f"获取利润表失败: {e}") - # 4. 获取主要财务指标 - try: - financial_abstract = ak.stock_financial_abstract(symbol=symbol) - if financial_abstract is not None and not financial_abstract.empty: - # 提取关键财务指标 - key_indicators = [ - '净资产收益率(ROE)', '总资产报酬率(ROA)', '销售毛利率', '销售净利率', - '资产负债率', '流动比率', '速动比率', '存货周转率', '应收账款周转率', - '总资产周转率', '营业收入同比增长', '净利润同比增长' - ] - - # 筛选出包含关键指标的行 - indicator_rows = financial_abstract[financial_abstract['指标'].isin(key_indicators)] - - if not indicator_rows.empty: - # 获取最新的报告期数据(第一列日期) - date_columns = [col for col in financial_abstract.columns if col not in ['选项', '指标']] - if date_columns: - latest_date = date_columns[0] # 最新日期列 - - # 构建财务比率字典 - financial_ratios = {"报告期": latest_date} - - # 提取每个指标的最新值 - for _, row in indicator_rows.iterrows(): - indicator_name = row['指标'] - value = row.get(latest_date, 'N/A') - if value is not None and not (isinstance(value, float) and pd.isna(value)): - try: - financial_ratios[indicator_name] = str(value) - except: + # 3. 获取现金流量表(tushare未获取到时回退akshare) + if financial_data["cash_flow"] is None: + try: + cash_flow = ak.stock_financial_abstract_ths(symbol=symbol, indicator="现金流量表") + if cash_flow is not None and not cash_flow.empty: + financial_data["cash_flow"] = cash_flow.head(8).to_dict('records') + except Exception as e: + print(f"获取现金流量表失败: {e}") + + # 4. 获取主要财务指标(tushare未获取到时回退akshare) + if not financial_data["financial_ratios"]: + try: + financial_abstract = ak.stock_financial_abstract(symbol=symbol) + if financial_abstract is not None and not financial_abstract.empty: + # 提取关键财务指标 + key_indicators = [ + '净资产收益率(ROE)', '总资产报酬率(ROA)', '销售毛利率', '销售净利率', + '资产负债率', '流动比率', '速动比率', '存货周转率', '应收账款周转率', + '总资产周转率', '营业收入同比增长', '净利润同比增长' + ] + + # 筛选出包含关键指标的行 + indicator_rows = financial_abstract[financial_abstract['指标'].isin(key_indicators)] + + if not indicator_rows.empty: + # 获取最新的报告期数据(第一列日期) + date_columns = [col for col in financial_abstract.columns if col not in ['选项', '指标']] + if date_columns: + latest_date = date_columns[0] # 最新日期列 + + # 构建财务比率字典 + financial_ratios = {"报告期": latest_date} + + # 提取每个指标的最新值 + for _, row in indicator_rows.iterrows(): + indicator_name = row['指标'] + value = row.get(latest_date, 'N/A') + if value is not None and not (isinstance(value, float) and pd.isna(value)): + try: + financial_ratios[indicator_name] = str(value) + except: + financial_ratios[indicator_name] = "N/A" + else: financial_ratios[indicator_name] = "N/A" - else: - financial_ratios[indicator_name] = "N/A" - - financial_data["financial_ratios"] = financial_ratios - except Exception as e: - print(f"获取财务指标失败: {e}") + + financial_data["financial_ratios"] = financial_ratios + except Exception as e: + print(f"获取财务指标失败: {e}") # 注意:季报数据现在由 quarterly_report_data.py 模块使用 akshare 获取(8期完整季报) # 不再使用问财获取季报,避免重复 diff --git a/test_tdx_api.py b/test_tdx_api.py index c777723..a80bd8f 100644 --- a/test_tdx_api.py +++ b/test_tdx_api.py @@ -14,7 +14,7 @@ from dotenv import load_dotenv load_dotenv() # 获取TDX API URL -TDX_API_URL = os.getenv('TDX_BASE_URL', 'http://127.0.0.1:5000') +TDX_API_URL = os.getenv('TDX_BASE_URL', 'https://tdx.javagood.top') print("=" * 60) print("TDX API配置测试") diff --git a/value_stock_strategy.py b/value_stock_strategy.py index 1b8c4cb..16f25e9 100644 --- a/value_stock_strategy.py +++ b/value_stock_strategy.py @@ -6,10 +6,10 @@ """ import pandas as pd -import akshare as ak from datetime import datetime, timedelta from typing import Dict, List, Optional import logging +from data_source_manager import data_source_manager class ValueStockStrategy: @@ -146,10 +146,9 @@ class ValueStockStrategy: RSI值 或 None """ try: - # 获取近60天日线数据 - df = ak.stock_zh_a_hist( + # 获取近60天日线数据(tushare优先,akshare备用) + df = data_source_manager.get_stock_hist_data( symbol=stock_code, - period="daily", start_date=(datetime.now() - timedelta(days=90)).strftime("%Y%m%d"), end_date=datetime.now().strftime("%Y%m%d"), adjust="qfq" @@ -159,7 +158,7 @@ class ValueStockStrategy: return None # 计算RSI - close = df['收盘'].astype(float) + close = df['close'].astype(float) delta = close.diff() gain = delta.where(delta > 0, 0) loss = (-delta).where(delta < 0, 0) diff --git a/verify_data_sources.py b/verify_data_sources.py new file mode 100644 index 0000000..ecd84df --- /dev/null +++ b/verify_data_sources.py @@ -0,0 +1,109 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +""" +数据源改造实盘验证脚本 +请在【有网络】的环境运行(例如在 PyCharm 的 aiagents-stock 解释器中执行), +它会逐项调用各模块的数据获取接口,并打印 PASS/FAIL 与实际使用的数据源日志。 +""" + +import os +import sys + +PROJECT = os.environ.get("AIAGENTS_STOCK_HOME", "/Users/songzhuoyuan/Desktop/code/python/aiagents-stock") +if not os.path.isdir(PROJECT): + PROJECT = input("请输入项目绝对路径: ").strip() +sys.path.insert(0, PROJECT) +os.chdir(PROJECT) + +print("=" * 60) +print("数据源改造验证 - 股票代码统一使用 600637") +print("=" * 60) + +from data_source_manager import data_source_manager +print("Tushare Token 已配置:", bool(data_source_manager.tushare_token)) +print("Tushare 可用:", data_source_manager.tushare_available) +print() + +passed = 0 +failed = 0 + +def test(name, fn): + global passed, failed + try: + result = fn() + if result: + passed += 1 + print(f"PASS | {name}") + else: + failed += 1 + print(f"FAIL | {name}(返回空/失败)") + except Exception as e: + failed += 1 + print(f"FAIL | {name}: {type(e).__name__}: {str(e)[:120]}") + +# 1. 数据源管理器(核心) +test("历史数据(1年,前复权)", lambda: data_source_manager.get_stock_hist_data("600637", start_date="20250801", end_date="20260810", adjust="qfq") is not None) +test("基本信息", lambda: data_source_manager.get_stock_basic_info("600637").get("name") != "未知") +test("实时行情", lambda: bool(data_source_manager.get_realtime_quotes("600637"))) +test("财务数据(利润表)", lambda: data_source_manager.get_financial_data("600637", "income") is not None) + +# 2. 资金流向 +def fund_flow_test(): + from fund_flow_akshare import FundFlowAkshareDataFetcher + return FundFlowAkshareDataFetcher().get_fund_flow_data("600637").get("data_success") +test("资金流向", fund_flow_test) + +# 3. 季报 +def quarterly_test(): + from quarterly_report_data import QuarterlyReportDataFetcher + return QuarterlyReportDataFetcher().get_quarterly_reports("600637").get("data_success") +test("季报(三表+指标)", quarterly_test) + +# 4. 新闻 +def news_test(): + from qstock_news_data import QStockNewsDataFetcher + return QStockNewsDataFetcher().get_stock_news("600637").get("data_success") +test("个股新闻", news_test) + +# 5. 市场情绪 +def sentiment_tests(): + from market_sentiment_data import MarketSentimentDataFetcher + f = MarketSentimentDataFetcher() + return (f._get_turnover_rate("600637") is not None and + f._get_market_index_sentiment() is not None) +test("情绪-换手率/大盘指数", sentiment_tests) + +def limit_test(): + from market_sentiment_data import MarketSentimentDataFetcher + return MarketSentimentDataFetcher()._get_limit_up_down_stats() is not None +test("情绪-涨跌停统计", limit_test) + +def margin_test(): + from market_sentiment_data import MarketSentimentDataFetcher + return MarketSentimentDataFetcher()._get_margin_trading_data("600637") is not None +test("情绪-融资融券", margin_test) + +# 6. 综合股票数据(主分析链路) +def stock_data_test(): + from stock_data import StockDataFetcher + f = StockDataFetcher() + info = f.get_stock_info("600637") + data = f.get_stock_data("600637", "1y") + fin = f.get_financial_data("600637") + return (isinstance(data, dict) and "error" not in data) and len(fin) > 0 +test("主分析链路(信息/行情/财务)", stock_data_test) + +# 7. 港股(可选,若积分不足会自动回退akshare) +def hk_test(): + from stock_data import StockDataFetcher + f = StockDataFetcher() + data = f.get_stock_data("00700", "1mo") + return isinstance(data, dict) and "error" not in data +test("港股日线(00700)", hk_test) + +print() +print("=" * 60) +print(f"验证完成:通过 {passed} 项,失败 {failed} 项") +if failed: + print("注意:失败项通常表示对应 tushare 接口积分不足或网络问题,程序会自动回退 akshare,不影响使用。") +print("=" * 60)