增加主力选股-批量分析功能

This commit is contained in:
oficcejo
2025-10-21 17:27:29 +08:00
parent 04de9c28bb
commit acc00ec702
5 changed files with 1011 additions and 304 deletions
+467
View File
@@ -13,6 +13,11 @@ import pandas as pd
def display_main_force_selector():
"""显示主力选股界面"""
# 检查是否触发批量分析(不立即删除标志)
if st.session_state.get('main_force_batch_trigger'):
run_main_force_batch_analysis()
return
st.markdown("## 🎯 主力选股 - 智能筛选优质标的")
st.markdown("---")
@@ -270,6 +275,50 @@ def display_analysis_results(result: dict, analyzer):
file_name=f"main_force_stocks_{datetime.now().strftime('%Y%m%d')}.csv",
mime="text/csv"
)
# 批量分析功能区
st.markdown("---")
col_batch1, col_batch2, col_batch3 = st.columns([2, 1, 1])
with col_batch1:
st.markdown("#### 🚀 批量深度分析")
st.caption("对主力资金净流入TOP股票进行完整的AI团队分析,获取投资评级和关键价位")
with col_batch2:
batch_count = st.selectbox(
"分析数量",
options=[10, 20, 30, 50],
index=1, # 默认20只
help="选择分析主力资金净流入前N只股票"
)
with col_batch3:
st.write("") # 占位
if st.button("🚀 开始批量分析", type="primary", use_container_width=True):
# 准备数据:按主力资金净流入排序
df_sorted = analyzer.raw_stocks.copy()
# 确保主力资金列是数值类型并排序
if main_fund_col:
df_sorted[main_fund_col] = pd.to_numeric(df_sorted[main_fund_col], errors='coerce')
df_sorted = df_sorted.sort_values(by=main_fund_col, ascending=False)
# 提取股票代码并去掉市场后缀(.SH, .SZ等)
raw_codes = df_sorted.head(batch_count)['股票代码'].tolist()
stock_codes = []
for code in raw_codes:
# 去掉后缀(如果有的话)
if isinstance(code, str):
# 去掉 .SH, .SZ, .BJ 等后缀
clean_code = code.split('.')[0] if '.' in code else code
stock_codes.append(clean_code)
else:
stock_codes.append(str(code))
# 存储到session_state,触发批量分析
st.session_state.main_force_batch_codes = stock_codes
st.session_state.main_force_batch_trigger = True
st.rerun()
# 显示PDF报告下载区域
if analyzer and result:
@@ -417,3 +466,421 @@ def format_number(value, unit='', suffix=''):
except (ValueError, TypeError):
return str(value)
def run_main_force_batch_analysis():
"""执行主力选股TOP股票批量分析(遵循统一调用规范)"""
import time
import re
st.markdown("## 🚀 主力选股TOP股票批量分析")
st.markdown("---")
# 检查是否已有分析结果
if st.session_state.get('main_force_batch_results'):
display_main_force_batch_results(st.session_state.main_force_batch_results)
# 返回按钮
col_back, col_clear = st.columns(2)
with col_back:
if st.button("🔙 返回主力选股", use_container_width=True):
# 清除所有批量分析相关状态
if 'main_force_batch_trigger' in st.session_state:
del st.session_state.main_force_batch_trigger
if 'main_force_batch_codes' in st.session_state:
del st.session_state.main_force_batch_codes
if 'main_force_batch_results' in st.session_state:
del st.session_state.main_force_batch_results
st.rerun()
with col_clear:
if st.button("🔄 重新分析", use_container_width=True):
# 清除结果,保留触发标志和代码
if 'main_force_batch_results' in st.session_state:
del st.session_state.main_force_batch_results
st.rerun()
return
# 获取股票代码列表
stock_codes = st.session_state.get('main_force_batch_codes', [])
if not stock_codes:
st.error("未找到股票代码列表")
# 清除触发标志
if 'main_force_batch_trigger' in st.session_state:
del st.session_state.main_force_batch_trigger
return
st.info(f"即将分析 {len(stock_codes)} 只股票:{', '.join(stock_codes[:10])}{'...' if len(stock_codes) > 10 else ''}")
# 返回按钮
if st.button("🔙 取消返回", type="secondary"):
# 清除所有批量分析相关状态
if 'main_force_batch_trigger' in st.session_state:
del st.session_state.main_force_batch_trigger
if 'main_force_batch_codes' in st.session_state:
del st.session_state.main_force_batch_codes
st.rerun()
st.markdown("---")
# 分析选项
col1, col2 = st.columns(2)
with col1:
analysis_mode = st.selectbox(
"分析模式",
options=["sequential", "parallel"],
format_func=lambda x: "顺序分析(稳定)" if x == "sequential" else "并行分析(快速)",
help="顺序分析较慢但稳定,并行分析更快但消耗更多资源"
)
with col2:
if analysis_mode == "parallel":
max_workers = st.number_input(
"并行线程数",
min_value=2,
max_value=5,
value=3,
help="同时分析的股票数量"
)
else:
max_workers = 1
st.markdown("---")
# 开始分析按钮
col_confirm, col_cancel = st.columns(2)
start_analysis = False
with col_confirm:
if st.button("🚀 确认开始分析", type="primary", use_container_width=True):
start_analysis = True
with col_cancel:
if st.button("❌ 取消", type="secondary", use_container_width=True):
# 清除所有批量分析相关状态
if 'main_force_batch_trigger' in st.session_state:
del st.session_state.main_force_batch_trigger
if 'main_force_batch_codes' in st.session_state:
del st.session_state.main_force_batch_codes
st.rerun()
if start_analysis:
# 导入统一分析函数(遵循统一规范)
from app import analyze_single_stock_for_batch
import concurrent.futures
import time
st.markdown("---")
st.info("⏳ 正在执行批量分析,请稍候...")
# 显示即将分析的股票代码(调试用)
with st.expander("🔍 调试信息", expanded=True):
st.write(f"**股票代码数量**: {len(stock_codes)}")
st.write(f"**股票代码列表**: {stock_codes}")
st.write(f"**代码格式检查**: {'✅ 无后缀,格式正确' if all('.' not in str(c) for c in stock_codes) else '❌ 包含后缀,可能有问题'}")
st.write(f"**分析模式**: {analysis_mode}")
st.write(f"**线程数**: {max_workers if analysis_mode == 'parallel' else 1}")
# 配置分析师参数
enabled_analysts_config = {
'technical': True,
'fundamental': True,
'fund_flow': True,
'risk': True,
'sentiment': False, # 禁用以提升速度
'news': False # 禁用以提升速度
}
selected_model = 'deepseek-chat'
period = '1y'
# 创建进度显示
progress_bar = st.progress(0)
status_text = st.empty()
# 存储结果
results = []
# 记录开始时间
start_time = time.time()
if analysis_mode == "sequential":
# 顺序分析
for i, code in enumerate(stock_codes):
status_text.text(f"正在分析 {code} ({i+1}/{len(stock_codes)})")
progress_bar.progress((i + 1) / len(stock_codes))
try:
# 调用统一分析函数
result = analyze_single_stock_for_batch(
symbol=code,
period=period,
enabled_analysts_config=enabled_analysts_config,
selected_model=selected_model
)
results.append(result)
except Exception as e:
results.append({
"symbol": code,
"success": False,
"error": str(e)
})
else:
# 并行分析
status_text.text(f"并行分析 {len(stock_codes)} 只股票({max_workers}线程)...")
def analyze_one(code):
try:
result = analyze_single_stock_for_batch(
symbol=code,
period=period,
enabled_analysts_config=enabled_analysts_config,
selected_model=selected_model
)
return result
except Exception as e:
return {"symbol": code, "success": False, "error": str(e)}
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = {executor.submit(analyze_one, code): code for code in stock_codes}
completed = 0
for future in concurrent.futures.as_completed(futures):
code = futures[future] # 获取对应的股票代码
completed += 1
progress_bar.progress(completed / len(stock_codes))
status_text.text(f"已完成 {completed}/{len(stock_codes)} ({code})")
try:
result = future.result()
results.append(result)
except Exception as e:
results.append({"symbol": code, "success": False, "error": str(e)})
# 清除进度
progress_bar.empty()
status_text.empty()
# 计算统计
elapsed_time = time.time() - start_time
success_count = sum(1 for r in results if r.get("success", False))
failed_count = len(results) - success_count
# 显示完成信息
if success_count > 0:
st.success(f"✅ 批量分析完成!成功 {success_count} 只,失败 {failed_count} 只,耗时 {elapsed_time/60:.1f} 分钟")
else:
st.error(f"❌ 批量分析完成,但所有 {failed_count} 只股票都分析失败!")
# 显示失败原因(调试用)
with st.expander("❌ 查看失败原因", expanded=True):
for r in results:
if not r.get("success", False):
st.error(f"**{r.get('symbol', 'N/A')}**: {r.get('error', '未知错误')}")
# 保存结果到session_state
st.session_state.main_force_batch_results = {
"results": results,
"total": len(results),
"success": success_count,
"failed": failed_count,
"elapsed_time": elapsed_time,
"analysis_mode": analysis_mode
}
time.sleep(1)
# 重新渲染以显示结果
st.rerun()
def display_main_force_batch_results(batch_results):
"""显示主力选股批量分析结果"""
import re
results = batch_results['results']
total = batch_results['total']
success = batch_results['success']
failed = batch_results['failed']
elapsed_time = batch_results['elapsed_time']
st.markdown("## 📊 批量分析结果")
st.markdown("---")
# 统计信息
col1, col2, col3, col4 = st.columns(4)
with col1:
st.metric("总计分析", f"{total}")
with col2:
st.metric("成功分析", f"{success}", delta=f"{success/total*100:.1f}%")
with col3:
st.metric("失败分析", f"{failed}")
with col4:
st.metric("总耗时", f"{elapsed_time/60:.1f} 分钟")
st.markdown("---")
# 成功分析的股票
successful_results = [r for r in results if r['success']]
if successful_results:
st.markdown(f"### ✅ 成功分析的股票 ({len(successful_results)}只)")
# 创建DataFrame展示
display_data = []
for result in successful_results:
stock_info = result.get('stock_info', {})
final_decision = result.get('final_decision', {})
# 提取评级emoji
rating = final_decision.get('rating', '未知')
rating_emoji = {
'强烈买入': '🔥',
'买入': '',
'持有': '⏸️',
'卖出': '⚠️',
'强烈卖出': '🚫'
}.get(rating, '')
display_data.append({
'股票代码': stock_info.get('symbol', ''),
'股票名称': stock_info.get('name', ''),
'评级': f"{rating_emoji} {rating}",
'信心度': final_decision.get('confidence_level', 'N/A'),
'进场区间': final_decision.get('entry_range', 'N/A'),
'止盈位': final_decision.get('take_profit', 'N/A'),
'止损位': final_decision.get('stop_loss', 'N/A'),
'目标价': final_decision.get('target_price', 'N/A')
})
df_display = pd.DataFrame(display_data)
st.dataframe(df_display, use_container_width=True, height=400)
# 详细分析结果(可展开)
st.markdown("---")
st.markdown("### 📋 详细分析报告")
for result in successful_results:
stock_info = result.get('stock_info', {})
final_decision = result.get('final_decision', {})
symbol = stock_info.get('symbol', '')
name = stock_info.get('name', '')
rating = final_decision.get('rating', '未知')
rating_emoji = {
'强烈买入': '🔥',
'买入': '',
'持有': '⏸️',
'卖出': '⚠️',
'强烈卖出': '🚫'
}.get(rating, '')
with st.expander(f"{rating_emoji} {symbol} - {name} | {rating}"):
# 关键信息
col1, col2, col3 = st.columns(3)
with col1:
st.metric("信心度", final_decision.get('confidence_level', 'N/A'))
with col2:
st.metric("进场区间", final_decision.get('entry_range', 'N/A'))
with col3:
st.metric("目标价", final_decision.get('target_price', 'N/A'))
# 止盈止损
col1, col2 = st.columns(2)
with col1:
st.metric("止盈位", final_decision.get('take_profit', 'N/A'))
with col2:
st.metric("止损位", final_decision.get('stop_loss', 'N/A'))
# 投资建议
st.markdown("#### 💡 投资建议")
advice = final_decision.get('advice', '暂无建议')
st.info(advice)
# 加入监测按钮
if st.button(f" 加入监测列表", key=f"monitor_{symbol}"):
# 解析进场区间
entry_range = final_decision.get('entry_range', '')
entry_min, entry_max = None, None
if entry_range and isinstance(entry_range, str) and "-" in entry_range:
try:
parts = entry_range.split("-")
entry_min = float(parts[0].strip())
entry_max = float(parts[1].strip())
except:
pass
# 解析止盈止损
take_profit_str = final_decision.get('take_profit', '')
take_profit = None
if take_profit_str:
try:
numbers = re.findall(r'\d+\.?\d*', str(take_profit_str))
if numbers:
take_profit = float(numbers[0])
except:
pass
stop_loss_str = final_decision.get('stop_loss', '')
stop_loss = None
if stop_loss_str:
try:
numbers = re.findall(r'\d+\.?\d*', str(stop_loss_str))
if numbers:
stop_loss = float(numbers[0])
except:
pass
# 调用监测管理器添加
from monitor_db import monitor_db
try:
# 准备进场区间数据
entry_range_dict = {}
if entry_min and entry_max:
entry_range_dict = {"min": entry_min, "max": entry_max}
# 添加到监测列表
monitor_db.add_monitored_stock(
symbol=symbol,
name=name,
rating=rating,
entry_range=entry_range_dict if entry_range_dict else None,
take_profit=take_profit,
stop_loss=stop_loss,
notes=f"主力选股批量分析 | {rating}"
)
st.success(f"{symbol} - {name} 已加入监测列表")
except Exception as e:
st.error(f"❌ 添加失败: {str(e)}")
# 失败的股票
failed_results = [r for r in results if not r['success']]
if failed_results:
st.markdown("---")
st.markdown(f"### ❌ 分析失败的股票 ({len(failed_results)}只)")
failed_data = []
for result in failed_results:
failed_data.append({
'股票代码': result.get('symbol', ''),
'失败原因': result.get('error', '未知错误')
})
df_failed = pd.DataFrame(failed_data)
st.dataframe(df_failed, use_container_width=True)