Compare commits

...

19 Commits

Author SHA1 Message Date
selfrelease b4b66f0508 feat: 全景扫描推荐列表分页改为5只/页,支持收起/展开 2026-07-22 10:19:58 +08:00
freedakgmail c79f968ad7 fix: DB K线查询改用LIMIT限制交易日条数,与API行为一致
原先DB用 trade_date >= now-120天(自然日)过滤,只返回~82个交易日;
API用 limit=120 返回120个交易日。K线条数不一致导致指标计算结果
不同,进而导致信号判定不同。

改为 ORDER BY trade_date DESC LIMIT min(days,300) + reversed,
与API的limit参数行为完全一致。
2026-07-22 08:07:01 +08:00
freedakgmail 88bc583f4d revert: 还原实时检测use_local_db=False,保持与扫描结果互补
两个数据源各有价值:
- 扫描结果(本地DB快照):保证列表一致性
- 实时检测(外部API最新):提供最新数据
还原此前的设计意图。
2026-07-22 08:04:30 +08:00
freedakgmail 05a1dcc704 fix: 实时检测优先使用本地DB,与扫描结果数据源统一
将 _get_kline_data 的 use_local_db 从 False 改为 True,
确保实时检测与全景扫描使用同一份K线数据,
避免因数据源不同导致信号判定结果不一致。
2026-07-22 08:03:30 +08:00
freedakgmail cf7341fc63 perf: 6项性能与正确性优化
1. get_stock_name: 延迟批量保存缓存(60秒去抖),避免每次新缓存都写文件
2. get_stock_name: 优先从DB stock_realtime_price查name,未命中再调腾讯API
3. compute_comprehensive_score: 单股模式也使用_fund_flow_cache(30分钟TTL)
4. _get_kline_from_local_db: 移除conn.autocommit=True,避免污染连接池事务模式
5. deep_analyze: K线天数从180改为120,与实时检测统一
6. detect_all_signals: 已有指标列时跳过重复calc_all_indicators
2026-07-22 08:00:27 +08:00
freedakgmail 037d9cd26c fix: 统一退市/ST股票过滤逻辑
- sync_kline.py: get_all_stock_codes SQL 添加 ST 过滤
- full_signal_scan.py: get_all_stock_codes SQL 添加 ST 过滤
- routes/analysis.py: total_stocks/triggered_stocks/信号分布/推荐分布
  统计查询统一添加退市/ST过滤,确保分母与过滤后扫描范围一致
2026-07-22 07:49:35 +08:00
freedakgmail 2697345670 chore: requirements.txt 显式声明 requests 依赖 2026-07-22 07:44:19 +08:00
freedakgmail cf1023754a fix: 修复 smart_trade_engine.py 3处除零风险
- :504 buy_price 可能为0导致除零,添加 if buy_price > 0 else 0 保护
- :1006 partial_sell 日志中 pos['avg_cost'] 除零,添加 if pos['avg_cost'] else 0
- :1055 sell 日志中 pos['avg_cost'] 除零,添加 if pos['avg_cost'] else 0
2026-07-22 07:42:17 +08:00
freedakgmail cfeec11929 fix: 修复 scheduler.py 和 sim_trade.py 中智能引擎调用连接泄漏
- scheduler.py: job_morning_trade/job_afternoon_trade 智能引擎 put_db 移到 finally
- sim_trade.py: trigger_trade 降级路径 put_db 移到 finally,初始化 conn=None
- 统一使用 smart_conn 变量名避免与外层 conn 冲突
2026-07-22 07:40:31 +08:00
freedakgmail 6a33fd534e fix: 修复 market.py 北交所映射缺失和 analysis.py 3处连接泄漏
- market.py:488: 补充北交所(8/9开头)腾讯API代码映射
- analysis.py: ai_analyze_stream 获取数据 put_db 移到 finally,初始化 conn=None
- analysis.py: ai_analyze_stream AI日志 put_db 移到 finally,初始化 _conn=None
- analysis.py: ai_analyze 获取资金流向 put_db 移到 finally,初始化 conn=None
2026-07-22 07:38:34 +08:00
freedakgmail c5df676af6 fix: 修复 market.py _put_db 拼写错误和多处连接泄漏
- market.py:308: _put_db(conn) → put_db(_conn) 修复函数名和变量名拼写错误
- market.py: get_fundamental put_db 移到 finally,初始化 conn=None 防 NameError
- admin.py: kline_sync_status 和 admin_scan_status put_db 移到 finally
2026-07-22 07:35:35 +08:00
freedakgmail 4b2d75c7ef fix: 修复 sim_trade/smart_trade 未导入 put_db 导致 NameError 和连接泄漏
- sim_trade.py: 导入 put_db(之前使用了但未导入,每次 finally 都会 NameError)
- smart_trade.py: 同上,9处 put_db 调用因未导入而失效
- sim_trade.py: trigger_trade 智能引擎异常时 conn 泄漏,改用 finally 确保归还
2026-07-22 07:33:43 +08:00
freedakgmail 5f5fa7ae80 fix: 修复连接泄漏、double-put_db、北交所映射缺失
- analysis.py: 6个端点将 put_db(conn) 从 try 移到 finally,异常时不泄漏
- trades.py: 移除 update_trade/delete_trade 中多余的 put_db(conn),避免 double-put 破坏连接池
- trades.py: stoploss_check 补充北交所(8/9开头)腾讯API代码映射
2026-07-22 07:31:21 +08:00
freedakgmail 4e35ab852f fix: 修复连接池泄漏和绕过问题
- 修复 analysis.py ai_analyze_stream 中 _put_db(conn) 拼写错误导致连接泄漏
- 将 analysis.py 6处 psycopg2.connect() 替换为 get_db()/put_db() 统一连接池管理
- 将 admin.py 2处 psycopg2.connect() 替换为 get_db()/put_db() 统一连接池管理
- 移除不再需要的 import psycopg2 和 from config import Config
2026-07-22 07:27:58 +08:00
freedakgmail 40ce519188 fix: 修复6个数据和代码问题
1. 修复 _generate_plain_summary position key 不匹配 (20d→d20, pct→range_pct)
2. 修复 mairui_api.py 日期解析不一致及 float(None) 崩溃风险
3. 修复 fund_flow_analyzer.py 麦蕊回退数据 close_price=0 导致量价背离误判
4. 修复 db_get_fund_flow_history 缺少完整字段和日期过滤
5. 修复 scheduler.py 连接池泄漏 (conn.close() → put_db(conn))
6. 修复多处北交所股票代码映射缺失 (8/9开头→bj)
2026-07-22 07:20:44 +08:00
freedakgmail 9a5e24ccd6 fix: 修复连接池泄漏(conn.close改为put_db) + put_db处理已关闭连接 + 添加日线同步cron 2026-07-22 07:12:09 +08:00
freedakgmail 488ef6b8e2 fix: 新浪/腾讯5min K线数据源补全amount字段,修复资金流向数据缺失 2026-07-22 06:55:35 +08:00
freedakgmail b330fbb3e5 fix: 政策面LLM分析覆盖score后summary不一致 2026-07-19 20:29:19 +08:00
selfrelease 7235bbe130 docs: 算法分析文档同步更新北向→南向资金 2026-07-18 18:44:30 +08:00
25 changed files with 402 additions and 254 deletions
+6 -3
View File
@@ -552,7 +552,7 @@ def db_save_fundamental(code, data):
# ========== 资金流向历史数据操作(数据库版) ==========
def db_get_fund_flow_history(code):
def db_get_fund_flow_history(code, limit=60):
"""获取股票的资金流向历史数据"""
conn = get_db()
if not conn:
@@ -564,11 +564,14 @@ def db_get_fund_flow_history(code):
SELECT code, trade_date::text, close_price, change_pct,
main_net_inflow, main_net_inflow_pct,
super_net_inflow, super_net_inflow_pct,
big_net_inflow, big_net_inflow_pct
big_net_inflow, big_net_inflow_pct,
mid_net_inflow, mid_net_inflow_pct,
small_net_inflow, small_net_inflow_pct
FROM stock_fund_flow_history
WHERE code = %s
ORDER BY trade_date DESC
""", (code,))
LIMIT %s
""", (code, limit))
rows = cur.fetchall()
# 获取最新日期
+16 -16
View File
@@ -14,7 +14,7 @@
信号检测层(signal_detector.py ← 第三章
↓ 基于7个信号检测买卖点
外部因素模块(fund_flow/sentiment/external/news)← 第四章
↓ 资金面/情绪/向/美股/商品/公告/政策/汇率 → 各因素评分
↓ 资金面/情绪/向/美股/商品/公告/政策/汇率 → 各因素评分
算法决策层(stock_algorithms.py ← 第五章
↓ 统一推荐逻辑、牛股阶段识别、技术面深度分析(7维度)→ 技术面基础分
综合评分引擎(score_engine.py
@@ -25,7 +25,7 @@
|------|------|------|------|
| 技术指标层 | `services/technical_indicators.py` | 计算 MACD、SKDJ、KDJ、EMA、MA 等基础指标 | 二 |
| 信号检测层 | `services/signal_detector.py` | 基于7个信号检测买卖点 | 三 |
| 外部因素模块 | `fund_flow_analyzer.py` / `market_sentiment.py` / `external_factors.py` / `news_analyzer.py` | 资金面、市场情绪、向资金、美股、大宗商品、公告/政策、汇率 | 四 |
| 外部因素模块 | `fund_flow_analyzer.py` / `market_sentiment.py` / `external_factors.py` / `news_analyzer.py` | 资金面、市场情绪、向资金、美股、大宗商品、公告/政策、汇率 | 四 |
| 算法决策层 | `services/stock_algorithms.py` | 统一推荐逻辑、牛股阶段识别、技术面深度分析 | 五 |
| 综合评分引擎 | `services/score_engine.py` | 整合外部因素(P0-P7)与技术面评分,输出最终评分 | 五 |
@@ -258,21 +258,21 @@
- **API端点**`GET /api/market_sentiment`
- **评分范围**:±10
### 4.3 向资金(P2,已实现)
### 4.3 向资金(P2,已实现)
**白话解释**向资金是从香港流入A股的"外资",被市场视为"聪明钱"。北向大幅买入通常被视为利好信号
**白话解释**向资金是内地通过港股通配置港股的资金,反映跨境资金情绪。南向大幅流入说明内地资金积极配置港股,大中华区risk-on,对A股偏正面。注:2024年8月起北向资金实时数据已停止公布,改用南向资金作为替代指标
| 信号 | 含义 | 可靠度 | 评分 |
|------|------|--------|------|
| 向单日净流入>50亿 | 外资看好,市场偏多 | ★★★★ | +5 |
| 向单日净流出>50亿 | 外资看空,注意风险 | ★★★★ | -5 |
| 向连续3日净流入 | 外资持续看好,中期偏多 | ★★★★★ | +3 |
| 向连续3日净流出 | 外资持续撤离,中期偏空 | ★★★★ | -3 |
| 向单日净流入>80亿 | 跨境资金情绪偏暖 | ★★★★ | +5 |
| 向单日净流出>80亿 | 跨境资金情绪偏冷 | ★★★★ | -5 |
| 向连续3日净流入 | 资金持续流入,中期偏多 | ★★★★★ | +3 |
| 向连续3日净流出 | 资金持续流出,中期偏空 | ★★★★ | -3 |
#### 实现模块
- **模块文件**`services/external_factors.py``get_northbound_capital()`
- **数据来源**AKShare `stock_hsgt_north_net_flow_in_em`向资金净流入
- **模块文件**`services/external_factors.py``get_southbound_capital()`
- **数据来源**AKShare `stock_hsgt_fund_flow_summary_em`(今日汇总)+ `stock_hsgt_hist_em`向资金历史
- **评分范围**:±10
### 4.4 美股隔夜板块变化(P3,已实现)
@@ -413,7 +413,7 @@
资金流向数据 → 资金信号 ────────────────────────────────┤
美股隔夜数据 → 外盘情绪 ────────────────────────────────┤→ 综合评分引擎 → 最终评分 → 买卖建议/AI解说
公告/新闻 → LLM情感分析 ───────────────────────────────┤
向资金 → 外资动向 ────────────────────────────────────┤
向资金 → 跨境资金情绪 ────────────────────────────────┤
市场情绪指标 → 情绪评分 ────────────────────────────────┤
大宗商品 → 板块影响 ────────────────────────────────────┤
汇率 → 进出口影响 ──────────────────────────────────────┘
@@ -426,7 +426,7 @@
| 技术面基础分 | 0-100 | `compute_deep_analysis` 原始分 |
| P0 资金面 | ±20 | 连续流入+10,吸筹+8,大单突击+5 |
| P1 市场情绪 | ±10 | 涨跌停比+5/-5,连板+3,成交额+2/-2 |
| P2 向资金 | ±10 | 大幅流入+5,连续流入+3 |
| P2 向资金 | ±10 | 大幅流入+5,连续流入+3 |
| P3 美股外盘 | ±10 | 美股大涨+5,大跌-5 |
| P4 大宗商品 | ±5 | 单品种涨跌±1 |
| P5 公告/异动 | ±15 | 并购+10,业绩预增+8,异动±5 |
@@ -445,7 +445,7 @@
|------|------|------|
| 资金流向分析 | `services/fund_flow_analyzer.py` | 从DB读取资金流向,计算连续流入/流出、量价背离、大单突击 |
| 市场情绪指标 | `services/market_sentiment.py` | 从实时行情表计算涨停跌停比、连板高度、换手率中位数、两市成交额 |
| 外部因素 | `services/external_factors.py` | 向资金、美股隔夜板块、大宗商品、汇率变化 |
| 外部因素 | `services/external_factors.py` | 向资金、美股隔夜板块、大宗商品、汇率变化 |
| 新闻/公告分析 | `services/news_analyzer.py` | 公告采集+分类、LLM情感分析、异动检测、政策面监控 |
| 综合评分引擎 | `services/score_engine.py` | 汇总技术面+所有外部因素,输出最终评分 |
@@ -454,7 +454,7 @@
| 端点 | 方法 | 说明 |
|------|------|------|
| `/api/market_sentiment` | GET | 市场情绪指标 |
| `/api/external_factors` | GET | 外部因素综合数据(向/美股/商品/汇率) |
| `/api/external_factors` | GET | 外部因素综合数据(向/美股/商品/汇率) |
| `/api/fund_flow_analysis/<code>` | GET | 个股资金流向分析 |
| `/api/news_analysis/<code>` | GET | 个股消息面分析(公告+政策+异动) |
@@ -658,7 +658,7 @@ deep_analyze 接口调用流程:
|------|------|----------|------|
| P0 资金面 | `fund_flow_analyzer` | ±20 | 主力在买还是在卖?有没有暗中吸筹/出货? |
| P1 市场情绪 | `market_sentiment` | ±10 | 今天涨停的股票多还是跌停的多?市场热不热? |
| P2 向资金 | `external_factors` | ±10 | 外资今天是买还是卖 |
| P2 向资金 | `external_factors` | ±10 | 跨境资金今天是流入还是流出 |
| P3 美股外盘 | `external_factors` | ±10 | 昨晚美股涨了还是跌了? |
| P4 大宗商品 | `external_factors` | ±5 | 原油/黄金/铜的价格变化对相关板块的影响 |
| P5 公告/异动 | `news_analyzer` | ±15 | 有没有并购/业绩预告等重大消息?股价有没有异动? |
@@ -681,7 +681,7 @@ deep_analyze 接口调用流程:
5. 形态识别(发现了什么技术形态)
6. 综合建议(根据评分给出操作建议)
7. 空间估算(风险收益比如何)
8. 外部因素(资金面/市场情绪/向资金/美股/消息面等综合影响)
8. 外部因素(资金面/市场情绪/向资金/美股/消息面等综合影响)
还可以调用豆包 LLM 对规则文本进行润色,让表达更自然生动。
+16 -16
View File
@@ -15,7 +15,7 @@
信号检测层(signal_detector.py ← 第三章
↓ 基于7个信号检测买卖点
外部因素模块(fund_flow/sentiment/external/news)← 第四章
↓ 资金面/情绪/向/美股/商品/公告/政策/汇率 → 各因素评分
↓ 资金面/情绪/向/美股/商品/公告/政策/汇率 → 各因素评分
算法决策层(stock_algorithms.py ← 第五章
↓ 统一推荐逻辑、牛股阶段识别、技术面深度分析(7维度)→ 技术面基础分
综合评分引擎(score_engine.py ← 第五章
@@ -26,7 +26,7 @@
|------|------|------|------|
| 技术指标层 | `services/technical_indicators.py` | 计算 MACD、SKDJ、KDJ、EMA、MA 等基础指标 | 二 |
| 信号检测层 | `services/signal_detector.py` | 基于7个信号检测买卖点 | 三 |
| 外部因素模块 | `fund_flow_analyzer.py` / `market_sentiment.py` / `external_factors.py` / `news_analyzer.py` | 资金面、市场情绪、向资金、美股、大宗商品、公告/政策、汇率 | 四 |
| 外部因素模块 | `fund_flow_analyzer.py` / `market_sentiment.py` / `external_factors.py` / `news_analyzer.py` | 资金面、市场情绪、向资金、美股、大宗商品、公告/政策、汇率 | 四 |
| 算法决策层 | `services/stock_algorithms.py` | 统一推荐逻辑、牛股阶段识别、技术面深度分析 | 五 |
| 综合评分引擎 | `services/score_engine.py` | 整合外部因素(P0-P7)与技术面评分,输出最终评分 | 五 |
@@ -259,21 +259,21 @@
- **API端点**`GET /api/market_sentiment`
- **评分范围**:±10
### 4.3 向资金(P2,已实现)
### 4.3 向资金(P2,已实现)
**白话解释**向资金是从香港流入A股的"外资",被市场视为"聪明钱"。北向大幅买入通常被视为利好信号
**白话解释**向资金是内地通过港股通配置港股的资金,反映跨境资金情绪。南向大幅流入说明内地资金积极配置港股,大中华区risk-on,对A股偏正面。注:2024年8月起北向资金实时数据已停止公布,改用南向资金作为替代指标
| 信号 | 含义 | 可靠度 | 评分 |
|------|------|--------|------|
| 向单日净流入>50亿 | 外资看好,市场偏多 | ★★★★ | +5 |
| 向单日净流出>50亿 | 外资看空,注意风险 | ★★★★ | -5 |
| 向连续3日净流入 | 外资持续看好,中期偏多 | ★★★★★ | +3 |
| 向连续3日净流出 | 外资持续撤离,中期偏空 | ★★★★ | -3 |
| 向单日净流入>80亿 | 跨境资金情绪偏暖 | ★★★★ | +5 |
| 向单日净流出>80亿 | 跨境资金情绪偏冷 | ★★★★ | -5 |
| 向连续3日净流入 | 资金持续流入,中期偏多 | ★★★★★ | +3 |
| 向连续3日净流出 | 资金持续流出,中期偏空 | ★★★★ | -3 |
#### 实现模块
- **模块文件**`services/external_factors.py``get_northbound_capital()`
- **数据来源**AKShare `stock_hsgt_north_net_flow_in_em`向资金净流入
- **模块文件**`services/external_factors.py``get_southbound_capital()`
- **数据来源**AKShare `stock_hsgt_fund_flow_summary_em`(今日汇总)+ `stock_hsgt_hist_em`向资金历史
- **评分范围**:±10
### 4.4 美股隔夜板块变化(P3,已实现)
@@ -414,7 +414,7 @@
资金流向数据 → 资金信号 ────────────────────────────────┤
美股隔夜数据 → 外盘情绪 ────────────────────────────────┤→ 综合评分引擎 → 最终评分 → 买卖建议/AI解说
公告/新闻 → LLM情感分析 ───────────────────────────────┤
向资金 → 外资动向 ────────────────────────────────────┤
向资金 → 跨境资金情绪 ────────────────────────────────┤
市场情绪指标 → 情绪评分 ────────────────────────────────┤
大宗商品 → 板块影响 ────────────────────────────────────┤
汇率 → 进出口影响 ──────────────────────────────────────┘
@@ -427,7 +427,7 @@
| 技术面基础分 | 0-100 | `compute_deep_analysis` 原始分 |
| P0 资金面 | ±20 | 连续流入+10,吸筹+8,大单突击+5 |
| P1 市场情绪 | ±10 | 涨跌停比+5/-5,连板+3,成交额+2/-2 |
| P2 向资金 | ±10 | 大幅流入+5,连续流入+3 |
| P2 向资金 | ±10 | 大幅流入+5,连续流入+3 |
| P3 美股外盘 | ±10 | 美股大涨+5,大跌-5 |
| P4 大宗商品 | ±5 | 单品种涨跌±1 |
| P5 公告/异动 | ±15 | 并购+10,业绩预增+8,异动±5 |
@@ -446,7 +446,7 @@
|------|------|------|
| 资金流向分析 | `services/fund_flow_analyzer.py` | 从DB读取资金流向,计算连续流入/流出、量价背离、大单突击 |
| 市场情绪指标 | `services/market_sentiment.py` | 从实时行情表计算涨停跌停比、连板高度、换手率中位数、两市成交额 |
| 外部因素 | `services/external_factors.py` | 向资金、美股隔夜板块、大宗商品、汇率变化 |
| 外部因素 | `services/external_factors.py` | 向资金、美股隔夜板块、大宗商品、汇率变化 |
| 新闻/公告分析 | `services/news_analyzer.py` | 公告采集+分类、LLM情感分析、异动检测、政策面监控 |
| 综合评分引擎 | `services/score_engine.py` | 汇总技术面+所有外部因素,输出最终评分 |
@@ -455,7 +455,7 @@
| 端点 | 方法 | 说明 |
|------|------|------|
| `/api/market_sentiment` | GET | 市场情绪指标 |
| `/api/external_factors` | GET | 外部因素综合数据(向/美股/商品/汇率) |
| `/api/external_factors` | GET | 外部因素综合数据(向/美股/商品/汇率) |
| `/api/fund_flow_analysis/<code>` | GET | 个股资金流向分析 |
| `/api/news_analysis/<code>` | GET | 个股消息面分析(公告+政策+异动) |
@@ -670,7 +670,7 @@ deep_analyze 接口调用流程:
|------|------|----------|------|
| P0 资金面 | `fund_flow_analyzer` | ±20 | 主力在买还是在卖?有没有暗中吸筹/出货? |
| P1 市场情绪 | `market_sentiment` | ±10 | 今天涨停的股票多还是跌停的多?市场热不热? |
| P2 向资金 | `external_factors` | ±10 | 外资今天是买还是卖 |
| P2 向资金 | `external_factors` | ±10 | 跨境资金今天是流入还是流出 |
| P3 美股外盘 | `external_factors` | ±10 | 昨晚美股涨了还是跌了? |
| P4 大宗商品 | `external_factors` | ±5 | 原油/黄金/铜的价格变化对相关板块的影响 |
| P5 公告/异动 | `news_analyzer` | ±15 | 有没有并购/业绩预告等重大消息?股价有没有异动? |
@@ -772,7 +772,7 @@ deep_analyze 接口调用流程:
5. 形态识别(发现了什么技术形态)
6. 综合建议(根据**综合得分**给出操作建议)
7. 空间估算(风险收益比如何)
8. 外部因素(资金面/市场情绪/向资金/美股/消息面等综合影响)
8. 外部因素(资金面/市场情绪/向资金/美股/消息面等综合影响)
9. **三项得分汇总**(技术得分、外部得分、综合得分)
还可以调用豆包 LLM 对规则文本进行润色,让表达更自然生动。
+1
View File
@@ -62,6 +62,7 @@ def get_all_stock_codes(conn):
SELECT code, name FROM stock_realtime_price
WHERE volume > 0 AND price > 0
AND name NOT LIKE '%%退%%'
AND name NOT LIKE '%%ST%%'
AND name NOT LIKE 'PT%%'
ORDER BY code
""")
+1
View File
@@ -1,5 +1,6 @@
flask>=2.0.0
flask-cors>=3.0.0
requests>=2.25.0
akshare>=1.10.0
pandas>=1.5.0
numpy>=1.20.0
+23 -28
View File
@@ -2,7 +2,7 @@
管理后台 API 路由
"""
from flask import Blueprint, request, jsonify, session, render_template
from db import get_db
from db import get_db, put_db
from psycopg2.extras import RealDictCursor
from werkzeug.security import generate_password_hash
import functools
@@ -182,7 +182,7 @@ def admin_required(f):
if not row or not row[0]:
return jsonify({'success': False, 'error': '无管理员权限'}), 403
finally:
conn.close()
put_db(conn)
return f(*args, **kwargs)
return decorated
@@ -201,7 +201,7 @@ def admin_page():
if not row or not row[0]:
return '无权访问', 403
finally:
conn.close()
put_db(conn)
return render_template('admin.html')
@@ -281,7 +281,7 @@ def dashboard():
return jsonify({'success': True, 'data': stats})
finally:
conn.close()
put_db(conn)
# ========== 用户管理 ==========
@@ -339,7 +339,7 @@ def list_users():
return jsonify({'success': True, 'data': users})
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/users/<int:user_id>', methods=['GET'])
@@ -437,7 +437,7 @@ def get_user_detail(user_id):
}
})
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/users/<int:user_id>/reset_password', methods=['POST'])
@@ -461,7 +461,7 @@ def reset_user_password(user_id):
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/users/<int:user_id>/toggle_admin', methods=['POST'])
@@ -483,7 +483,7 @@ def toggle_admin(user_id):
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/users/<int:user_id>', methods=['DELETE'])
@@ -509,7 +509,7 @@ def delete_user(user_id):
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ========== 用户数据管理 ==========
@@ -526,7 +526,7 @@ def delete_user_watchlist(user_id, code):
conn.commit()
return jsonify({'success': True})
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/users/<int:user_id>/trades/<int:trade_id>', methods=['DELETE'])
@@ -541,7 +541,7 @@ def delete_user_trade(user_id, trade_id):
conn.commit()
return jsonify({'success': True})
finally:
conn.close()
put_db(conn)
# ========== 系统数据 ==========
@@ -568,7 +568,7 @@ def scan_history():
""")
return jsonify({'success': True, 'data': [dict(r) for r in cur.fetchall()]})
finally:
conn.close()
put_db(conn)
@bp.route('/api/admin/data_stats', methods=['GET'])
@@ -603,7 +603,7 @@ def data_stats():
return jsonify({'success': True, 'data': {'tables': result, 'logs': logs}})
finally:
conn.close()
put_db(conn)
# ========== 全景扫描管理 ==========
@@ -753,13 +753,9 @@ def trigger_kline_sync():
def kline_sync_status():
"""管理员查询K线同步状态"""
try:
import psycopg2
from config import Config
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor()
# K线数据统计
@@ -786,7 +782,6 @@ def kline_sync_status():
total_stocks = cur.fetchone()[0]
cur.close()
conn.close()
is_running = _is_kline_sync_running()
@@ -862,6 +857,8 @@ def kline_sync_status():
})
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
# ========== 定时任务监控 ==========
@@ -1282,13 +1279,10 @@ def admin_scan_status():
"""管理员查询扫描进度"""
try:
from datetime import datetime
import psycopg2
from config import Config
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor()
scan_date = datetime.now().strftime('%Y-%m-%d')
@@ -1308,7 +1302,6 @@ def admin_scan_status():
triggered = cur.fetchone()[0]
cur.close()
conn.close()
is_running = _is_scan_running()
@@ -1357,3 +1350,5 @@ def admin_scan_status():
})
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
+56 -57
View File
@@ -15,7 +15,8 @@ from services.stock_algorithms import (
)
from db import (
login_required, get_current_user_id,
db_get_alerts_cache, db_save_alerts_cache
db_get_alerts_cache, db_save_alerts_cache,
get_db, put_db
)
bp = Blueprint('analysis', __name__, url_prefix='/api')
@@ -61,7 +62,7 @@ def analyze():
# 3. 获取实时价格补充到结果(优先腾讯API,兼容腾讯云)
try:
import requests as _req
_tcode = ('sh' if stock_code.startswith('6') else 'sz') + stock_code
_tcode = ('sh' if stock_code.startswith('6') else 'bj' if stock_code.startswith(('8', '9')) else 'sz') + stock_code
_r = _req.get(f'http://qt.gtimg.cn/q={_tcode}', timeout=5,
headers={'Referer': 'https://finance.qq.com'})
if _r.status_code == 200 and '\"' in _r.text:
@@ -98,7 +99,7 @@ def deep_analyze():
if not stock_code:
return jsonify({'error': '股票代码不能为空'}), 400
df = algo_get_kline_data(stock_code, days=180)
df = algo_get_kline_data(stock_code, days=120)
if df is None or len(df) < 30:
return jsonify({'error': 'K线数据不足'}), 400
@@ -313,7 +314,7 @@ def ai_analyze_stream(stock_code):
from flask import Response
from services.doubao_api import analyze_stock_stream, format_fund_flow, format_market_cap
from services.mairui_api import get_realtime_price as mairui_price, get_financial_indicators
from db import get_db
from db import get_db, put_db
def generate():
# 获取股票数据
@@ -339,6 +340,7 @@ def ai_analyze_stream(stock_code):
stock_name = get_stock_name(stock_code) or stock_code
# 获取资金流向和技术信号
conn = None
try:
conn = get_db()
if conn:
@@ -383,12 +385,14 @@ def ai_analyze_stream(stock_code):
stock_data['indicators'] = scan_row['indicators'] or {}
stock_data['triggered_count'] = scan_row['triggered_count'] or 0
cur2.close()
conn.close()
except Exception as e:
print(f"获取数据失败: {e}")
finally:
if conn:
put_db(conn)
# 记录AI调用日志
_conn = None
try:
from flask import session as _sess
_uid = _sess.get('user_id')
@@ -399,9 +403,11 @@ def ai_analyze_stream(stock_code):
_cur.execute("INSERT INTO ai_call_log (user_id, stock_code, stock_name) VALUES (%s, %s, %s)",
(_uid, stock_code, stock_name))
_conn.commit()
_conn.close()
except Exception:
pass
finally:
if _conn:
put_db(_conn)
# 流式调用AI
for chunk in analyze_stock_stream(stock_code, stock_name, stock_data):
@@ -421,7 +427,7 @@ def ai_analyze(stock_code):
try:
from services.doubao_api import analyze_stock
from services.mairui_api import get_realtime_price as mairui_price, get_financial_indicators
from db import get_db
from db import get_db, put_db
# 获取股票数据
stock_data = {}
@@ -446,6 +452,7 @@ def ai_analyze(stock_code):
stock_name = get_stock_name(stock_code) or stock_code
# 获取近三日资金流向
conn = None
try:
conn = get_db()
if conn:
@@ -467,9 +474,11 @@ def ai_analyze(stock_code):
'super_pct': float(row[3]) if row[3] else 0,
})
stock_data['fund_flow_3days'] = fund_flow
conn.close()
except Exception as e:
print(f"获取资金流向失败: {e}")
finally:
if conn:
put_db(conn)
# 调用AI分析
result = analyze_stock(stock_code, stock_name, stock_data)
@@ -544,10 +553,8 @@ def batch_technical_signals():
"""批量检测技术交易信号 — 优先从 stock_signal_scan 读取(与提醒一致),无记录时实时计算"""
try:
import pandas as pd
import psycopg2
from psycopg2.extras import RealDictCursor
from services.signal_detector import detect_all_signals
from config import Config
data = request.get_json()
codes = data.get('codes', [])
@@ -566,10 +573,9 @@ def batch_technical_signals():
change_map = {}
scan_map = {}
try:
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
raise Exception('数据库连接失败')
cur = conn.cursor(cursor_factory=RealDictCursor)
placeholders = ','.join(['%s'] * len(codes))
@@ -603,9 +609,10 @@ def batch_technical_signals():
scan_map[row['code']] = row
cur.close()
conn.close()
except Exception as e:
print(f"批量扫描获取数据失败: {e}")
finally:
put_db(conn)
results = []
errors = []
@@ -675,9 +682,6 @@ def batch_technical_signals():
def get_scan_results():
"""查询全量扫描结果"""
try:
import psycopg2
from config import Config
scan_date = request.args.get('date', datetime.now().strftime('%Y-%m-%d'))
min_triggered = int(request.args.get('min_triggered', 0))
signal_type = request.args.get('signal_type', '')
@@ -689,10 +693,9 @@ def get_scan_results():
holding_set = set(c.strip() for c in holding_codes_str.split(',') if c.strip())
recommend_text = (request.args.get('recommend_text') or '').strip()
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor()
# 检查请求日期是否有数据,如果没有则自动回退到最近可用的扫描日期
@@ -714,7 +717,8 @@ def get_scan_results():
)
total_scanned = cur.fetchone()[0]
cur.execute("SELECT count(*) FROM stock_realtime_price")
cur.execute("""SELECT count(*) FROM stock_realtime_price
WHERE name NOT LIKE '%%退%%' AND name NOT LIKE '%%ST%%' AND name NOT LIKE 'PT%%'""")
total_stocks = cur.fetchone()[0]
where_clauses = ["scan_date = %s"]
@@ -890,6 +894,7 @@ def get_scan_results():
count(*) as cnt
FROM stock_signal_scan, jsonb_array_elements(signal_status) s
WHERE scan_date = %s AND (s.value->>'triggered')::boolean = true
AND name NOT LIKE '%%退%%' AND name NOT LIKE '%%ST%%'
GROUP BY s.value->>'name', s.value->>'type'
ORDER BY cnt DESC
""", (scan_date,))
@@ -901,6 +906,7 @@ def get_scan_results():
cur.execute("""
SELECT count(*) FROM stock_signal_scan
WHERE scan_date = %s AND triggered_count > 0
AND name NOT LIKE '%%退%%' AND name NOT LIKE '%%ST%%'
""", (scan_date,))
triggered_stocks = cur.fetchone()[0]
@@ -916,6 +922,7 @@ def get_scan_results():
cur.execute("""
SELECT code, signal_status, indicators, triggered_count
FROM stock_signal_scan WHERE scan_date = %s
AND name NOT LIKE '%%退%%' AND name NOT LIKE '%%ST%%'
""", (scan_date,))
recommend_counts = {}
for r in cur.fetchall():
@@ -925,7 +932,6 @@ def get_scan_results():
recommend_counts[disp] = recommend_counts.get(disp, 0) + 1
cur.close()
conn.close()
return jsonify({
'success': True,
@@ -945,6 +951,8 @@ def get_scan_results():
})
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
def _compute_recommend(signal_status, indicators, triggered_count, is_holding):
@@ -955,9 +963,7 @@ def _compute_recommend(signal_status, indicators, triggered_count, is_holding):
@bp.route('/signal_alerts', methods=['POST'])
def signal_alerts():
"""基于信号扫描结果生成买入/卖出/观望提醒(与扫描结果共用 _compute_recommend"""
import psycopg2
from psycopg2.extras import RealDictCursor
from config import Config
try:
data = request.get_json() or {}
@@ -968,10 +974,9 @@ def signal_alerts():
if not stock_codes:
return jsonify({'success': True, 'results': []})
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor(cursor_factory=RealDictCursor)
# 优先今天的扫描数据,无则回退到最近可用日期
@@ -1003,7 +1008,6 @@ def signal_alerts():
WHERE code IN ({placeholders})
""", stock_codes)
price_rows = cur.fetchall()
conn.close()
price_map = {}
change_map = {}
@@ -1085,6 +1089,8 @@ def signal_alerts():
import traceback
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
def _is_scan_running():
@@ -1127,13 +1133,9 @@ def get_stock_score_detail(code):
def get_scan_status():
"""查询扫描进度"""
try:
import psycopg2
from config import Config
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor()
scan_date = datetime.now().strftime('%Y-%m-%d')
@@ -1143,7 +1145,8 @@ def get_scan_status():
)
scanned = cur.fetchone()[0]
cur.execute("SELECT count(*) FROM stock_realtime_price")
cur.execute("""SELECT count(*) FROM stock_realtime_price
WHERE name NOT LIKE '%%退%%' AND name NOT LIKE '%%ST%%' AND name NOT LIKE 'PT%%'""")
total = cur.fetchone()[0]
cur.execute(
@@ -1153,7 +1156,6 @@ def get_scan_status():
triggered = cur.fetchone()[0]
cur.close()
conn.close()
return jsonify({
'success': True,
@@ -1167,6 +1169,8 @@ def get_scan_status():
})
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
@bp.route('/start_full_scan', methods=['POST'])
@@ -1210,18 +1214,15 @@ def start_full_scan():
def get_scan_strategy():
"""基于体系最强战法,给出分梯队买卖建议。与全景扫描推荐共用 _compute_recommend,算法一致。"""
try:
import psycopg2
from psycopg2.extras import RealDictCursor
from config import Config
scan_date = request.args.get('date', datetime.now().strftime('%Y-%m-%d'))
holding_codes_str = request.args.get('holding_codes', '')
holding_set = set(c.strip() for c in holding_codes_str.split(',') if c.strip())
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor(cursor_factory=RealDictCursor)
# 如果前端未指定日期,且当天无数据,自动回退到最近扫描日期
@@ -1241,7 +1242,6 @@ def get_scan_strategy():
""", (scan_date,))
rows = cur.fetchall()
cur.close()
conn.close()
tier1, tier2, tier3, tier4 = [], [], [], []
for r in rows:
@@ -1322,10 +1322,12 @@ def get_scan_strategy():
})
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
def _get_kline_data(stock_code, days=120):
"""获取K线数据 — 委托给 services.stock_algorithms.get_kline_data(实时分析不用本地DB缓存)"""
"""获取K线数据 — 实时检测使用外部API获取最新数据,与扫描结果(本地DB快照)形成互补"""
return algo_get_kline_data(stock_code, days=days, use_local_db=False)
@@ -1345,19 +1347,16 @@ def get_bull_stocks():
stage: 可选,筛选特定阶段(1-5)
holdingStocks: 可选,持仓代码逗号分隔
"""
import psycopg2
from psycopg2.extras import RealDictCursor
from config import Config
try:
stage_filter = request.args.get('stage', type=int, default=0)
holding_str = request.args.get('holdingStocks', '')
holding_codes = set(holding_str.split(',')) if holding_str else set()
conn = psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
conn = get_db()
if not conn:
return jsonify({'success': False, 'error': '数据库连接失败'}), 500
cur = conn.cursor(cursor_factory=RealDictCursor)
# 获取最近一次扫描数据(过滤退市/ST)
@@ -1376,8 +1375,6 @@ def get_bull_stocks():
for p in cur.fetchall():
price_map[p['code']] = {'price': float(p['price']), 'change_pct': float(p.get('change_pct') or 0)}
conn.close()
# ---- 批量计算综合评分 ----
scores_map = None
try:
@@ -1439,6 +1436,8 @@ def get_bull_stocks():
import traceback
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
put_db(conn)
def _llm_polish_summary(stock_name, stock_code, ai_summary, score, verdict):
+11 -10
View File
@@ -6,7 +6,7 @@ import pandas as pd
from datetime import datetime, timedelta
from services.stock_service import get_stock_fund_flow, load_cached_data
from services.stock_algorithms import get_kline_data as algo_get_kline_data
from db import get_db
from db import get_db, put_db
bp = Blueprint('market', __name__, url_prefix='/api')
@@ -40,7 +40,7 @@ def db_realtime_price(stock_code):
'data': dict(row)
})
finally:
conn.close()
put_db(conn)
@bp.route('/db/realtime_prices', methods=['POST'])
@@ -71,7 +71,7 @@ def db_realtime_prices():
'data': [dict(row) for row in rows]
})
finally:
conn.close()
put_db(conn)
@bp.route('/db/fund_flow_today/<stock_code>', methods=['GET'])
@@ -102,7 +102,7 @@ def db_fund_flow_today(stock_code):
'data': dict(row)
})
finally:
conn.close()
put_db(conn)
@bp.route('/db/fund_flow_today_batch', methods=['POST'])
@@ -135,7 +135,7 @@ def db_fund_flow_today_batch():
'data': [dict(row) for row in rows]
})
finally:
conn.close()
put_db(conn)
@bp.route('/db/data_status', methods=['GET'])
@@ -171,7 +171,7 @@ def db_data_status():
'recent_logs': [dict(log) for log in logs]
})
finally:
conn.close()
put_db(conn)
# ============ 原有API(兼容) ============
@@ -305,7 +305,7 @@ def get_fund_flow_rank():
return jsonify({'success': True, 'data': flow_data, 'total': len(flow_data),
'source': 'cache'})
finally:
_conn.close()
put_db(_conn)
return jsonify({'success': True, 'data': [], 'total': 0})
except Exception as e:
@@ -324,6 +324,7 @@ def get_fundamental(stock_code):
realtime_data = None
signal_data = None
conn = None
try:
conn = get_db()
if conn:
@@ -388,10 +389,10 @@ def get_fundamental(stock_code):
'indicators': ind,
'triggered_count': sig_row[2] or 0,
}
conn.close()
except Exception as e:
print(f"获取数据失败: {e}")
finally:
put_db(conn)
# 优先从数据库获取当日缓存
cached = db_get_fundamental(stock_code)
@@ -484,7 +485,7 @@ def get_fundamental(stock_code):
# 备用方案:先试腾讯API,再试akshare
try:
import requests as _rq
_tc = ('sh' if stock_code.startswith('6') else 'sz') + stock_code
_tc = ('sh' if stock_code.startswith('6') else 'bj' if stock_code.startswith(('8', '9')) else 'sz') + stock_code
_rr = _rq.get(f'http://qt.gtimg.cn/q={_tc}', timeout=5,
headers={'Referer': 'https://finance.qq.com'})
if _rr.status_code == 200 and '\"' in _rr.text:
+33 -23
View File
@@ -4,7 +4,7 @@
"""
from flask import Blueprint, request, jsonify
from datetime import datetime, date, time
from db import get_db, login_required, get_current_user_id
from db import get_db, put_db, login_required, get_current_user_id
from psycopg2.extras import RealDictCursor
bp = Blueprint('sim_trade', __name__, url_prefix='/api/sim')
@@ -29,7 +29,7 @@ def init_user_config(user_id):
cur.execute("SELECT * FROM sim_config WHERE user_id = %s", (user_id,))
return cur.fetchone()
finally:
conn.close()
put_db(conn)
@bp.route('/config', methods=['GET'])
@@ -82,7 +82,7 @@ def update_config():
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/trades', methods=['GET'])
@@ -114,7 +114,7 @@ def get_trades():
trades = cur.fetchall()
return jsonify({'success': True, 'trades': trades})
finally:
conn.close()
put_db(conn)
@bp.route('/positions', methods=['GET'])
@@ -140,7 +140,7 @@ def get_positions():
positions = cur.fetchall()
return jsonify({'success': True, 'positions': positions})
finally:
conn.close()
put_db(conn)
@bp.route('/stats', methods=['GET'])
@@ -261,7 +261,7 @@ def get_stats():
'history': list(reversed(history)) if history else []
})
finally:
conn.close()
put_db(conn)
@bp.route('/execute', methods=['POST'])
@@ -381,7 +381,7 @@ def execute_trade():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/auto_execute', methods=['POST'])
@@ -532,7 +532,7 @@ def auto_execute():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/update_prices', methods=['POST'])
@@ -588,7 +588,7 @@ def update_prices():
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/reset', methods=['POST'])
@@ -612,7 +612,7 @@ def reset_simulation():
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/trigger_trade', methods=['POST'])
@@ -623,12 +623,12 @@ def trigger_trade():
try:
# 优先使用智能引擎
smart_conn = None
try:
from services.smart_trade_engine import execute_smart_trade
conn = get_db()
if conn:
result = execute_smart_trade(conn, user_id, scan_date=None)
conn.close()
smart_conn = get_db()
if smart_conn:
result = execute_smart_trade(smart_conn, user_id, scan_date=None)
if result.get('success'):
results = result.get('results', [])
buy_count = len([r for r in results if r['type'] == 'buy'])
@@ -641,19 +641,29 @@ def trigger_trade():
})
except Exception as e:
print(f"[trigger_trade] 智能引擎异常,降级: {e}")
finally:
if smart_conn:
put_db(smart_conn)
# 降级: 使用旧引擎
from services.scheduler import execute_auto_trade_for_user
conn = get_db()
if conn:
cur = conn.cursor(cursor_factory=RealDictCursor)
cur.execute("SELECT trade_quantity FROM sim_config WHERE user_id = %s", (user_id,))
config = cur.fetchone()
trade_quantity = config['trade_quantity'] if config else 1000
conn.close()
else:
conn = None
try:
conn = get_db()
if conn:
cur = conn.cursor(cursor_factory=RealDictCursor)
cur.execute("SELECT trade_quantity FROM sim_config WHERE user_id = %s", (user_id,))
config = cur.fetchone()
trade_quantity = config['trade_quantity'] if config else 1000
else:
trade_quantity = 1000
except Exception as e:
print(f"[trigger_trade] 获取交易配置失败: {e}")
trade_quantity = 1000
finally:
if conn:
put_db(conn)
result = execute_auto_trade_for_user(user_id, trade_quantity)
@@ -730,4 +740,4 @@ def get_today_trades():
}
})
finally:
conn.close()
put_db(conn)
+9 -9
View File
@@ -4,7 +4,7 @@
"""
from flask import Blueprint, request, jsonify
from datetime import date
from db import get_db, login_required, get_current_user_id
from db import get_db, put_db, login_required, get_current_user_id
from psycopg2.extras import RealDictCursor
bp = Blueprint('smart_trade', __name__, url_prefix='/api/smart')
@@ -43,7 +43,7 @@ def get_algo_templates():
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ═══════════════════════════════════════════════════════
@@ -78,7 +78,7 @@ def get_algo_config():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/config', methods=['POST'])
@@ -111,7 +111,7 @@ def save_algo_config():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
@bp.route('/apply_template', methods=['POST'])
@@ -139,7 +139,7 @@ def apply_template():
conn.rollback()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ═══════════════════════════════════════════════════════
@@ -164,7 +164,7 @@ def get_status():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ═══════════════════════════════════════════════════════
@@ -205,7 +205,7 @@ def trigger_smart_trade():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ═══════════════════════════════════════════════════════
@@ -242,7 +242,7 @@ def get_signals():
except Exception as e:
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
# ═══════════════════════════════════════════════════════
@@ -284,4 +284,4 @@ def get_position_meta():
traceback.print_exc()
return jsonify({'success': False, 'error': str(e)}), 500
finally:
conn.close()
put_db(conn)
+1 -4
View File
@@ -130,7 +130,6 @@ def update_trade(trade_id):
""", (trade_id, user_id))
old_trade = cur.fetchone()
if not old_trade:
put_db(conn)
return jsonify({'error': '交易记录不存在'}), 404
cur.execute("""
@@ -152,7 +151,6 @@ def update_trade(trade_id):
trade = cur.fetchone()
if not trade:
conn.rollback()
put_db(conn)
return jsonify({'error': '交易记录不存在'}), 404
old_delta = _calc_cash_delta(old_trade.get('trade_type'), old_trade.get('price'), old_trade.get('quantity'))
@@ -194,7 +192,6 @@ def delete_trade(trade_id):
""", (trade_id, user_id))
old_trade = cur.fetchone()
if not old_trade:
put_db(conn)
return jsonify({'success': False, 'error': '交易记录不存在'}), 404
cur.execute("DELETE FROM trades WHERE id = %s AND user_id = %s", (trade_id, user_id))
@@ -306,7 +303,7 @@ def check_stoploss():
current_price = 0
try:
import requests as _rq
_tc = ('sh' if code.startswith('6') else 'sz') + code
_tc = ('sh' if code.startswith('6') else 'bj' if code.startswith(('8', '9')) else 'sz') + code
_rr = _rq.get(f'http://qt.gtimg.cn/q={_tc}', timeout=5,
headers={'Referer': 'https://finance.qq.com'})
if _rr.status_code == 200 and '\"' in _rr.text:
+46 -10
View File
@@ -155,6 +155,39 @@ def _get_fund_flow_from_mairui(stock_code, days=10):
'small_net_inflow_pct': small_net_pct,
})
# 从本地DB补充 close_price 和 change_pct(避免全为0导致量价背离误判)
from db import get_db, put_db
conn = get_db()
if conn:
try:
cur = conn.cursor()
date_list = [r['date'] for r in records if r['date']]
if date_list:
cur.execute("""
SELECT k.trade_date::text, k.close,
CASE WHEN prev.close > 0
THEN ROUND((k.close - prev.close) / prev.close * 100, 2)
ELSE 0 END AS change_pct
FROM stock_kline_daily k
LEFT JOIN LATERAL (
SELECT close FROM stock_kline_daily
WHERE code = k.code AND trade_date < k.trade_date
ORDER BY trade_date DESC LIMIT 1
) prev ON true
WHERE k.code = %s AND k.trade_date::text = ANY(%s)
""", (stock_code, date_list))
price_map = {r[0]: {'close': float(r[1] or 0), 'change_pct': float(r[2] or 0)}
for r in cur.fetchall()}
for r in records:
info = price_map.get(r['date'])
if info:
r['close_price'] = info['close']
r['change_pct'] = info['change_pct']
except Exception as e:
logger.warning(f"补充K线价格失败({stock_code}): {e}")
finally:
put_db(conn)
# 按日期升序排列
records.sort(key=lambda x: x['date'])
logger.info(f"麦蕊API获取{stock_code}资金流向{len(records)}")
@@ -245,17 +278,20 @@ def analyze_fund_flow(stock_code, days=5):
# 出货:主力净流出但价格不跌(跌幅<2%)
accumulation = False
distribution = False
if total_main_inflow > 0:
price_changes = [r['change_pct'] for r in recent]
avg_price_change = sum(price_changes) / len(price_changes) if price_changes else 0
if avg_price_change < 2:
accumulation = True
# 检查是否有有效的价格数据(close_price 全为0说明数据不完整,跳过背离检测)
has_valid_price = any(r.get('close_price', 0) > 0 for r in recent)
if has_valid_price:
if total_main_inflow > 0:
price_changes = [r['change_pct'] for r in recent]
avg_price_change = sum(price_changes) / len(price_changes) if price_changes else 0
if avg_price_change < 2:
accumulation = True
if total_main_inflow < 0:
price_changes = [r['change_pct'] for r in recent]
avg_price_change = sum(price_changes) / len(price_changes) if price_changes else 0
if avg_price_change > -2:
distribution = True
if total_main_inflow < 0:
price_changes = [r['change_pct'] for r in recent]
avg_price_change = sum(price_changes) / len(price_changes) if price_changes else 0
if avg_price_change > -2:
distribution = True
# 单日超大单突击
big_surge = False
+11 -7
View File
@@ -245,17 +245,21 @@ def get_fund_flow(stock_code, days=3):
flow_data = []
for item in data:
# 计算主力净流入 = 主买大单+主买特大单 - 主卖大单-主卖特大单
main_buy = float(item.get('zmbddcje', 0)) + float(item.get('zmbtdcje', 0))
main_sell = float(item.get('zmsddcje', 0)) + float(item.get('zmstdcje', 0))
main_buy = float(item.get('zmbddcje', 0) or 0) + float(item.get('zmbtdcje', 0) or 0)
main_sell = float(item.get('zmsddcje', 0) or 0) + float(item.get('zmstdcje', 0) or 0)
main_net = main_buy - main_sell
# 日期解析:与 fund_flow_analyzer.py 保持一致,用字符串截取
t_str = str(item.get('t', ''))
date_str = t_str[:10] if t_str else ''
flow_data.append({
'date': datetime.fromtimestamp(item.get('t', 0)).strftime('%Y-%m-%d') if item.get('t') else '',
'date': date_str,
'main_net_inflow': main_net,
'super_buy': float(item.get('zmbtdcje', 0)),
'super_sell': float(item.get('zmstdcje', 0)),
'big_buy': float(item.get('zmbddcje', 0)),
'big_sell': float(item.get('zmsddcje', 0)),
'super_buy': float(item.get('zmbtdcje', 0) or 0),
'super_sell': float(item.get('zmstdcje', 0) or 0),
'big_buy': float(item.get('zmbddcje', 0) or 0),
'big_sell': float(item.get('zmsddcje', 0) or 0),
})
return {'success': True, 'data': flow_data}
+6 -2
View File
@@ -456,6 +456,7 @@ def analyze_policy_impact(policy_news):
affected[category] = impact
# 尝试用LLM深度分析重大政策
llm_summary = None
major_policies = [n for n in policy_news if n['category'] in ['行业扶持', '行业监管', '货币政策']]
if major_policies and len(major_policies) <= 5:
try:
@@ -463,14 +464,17 @@ def analyze_policy_impact(policy_news):
if llm_result:
score = llm_result.get('score', score)
reasons = llm_result.get('reasons', reasons)
llm_summary = llm_result.get('summary')
except Exception as e:
logger.warning(f"LLM政策分析失败: {e}")
score = max(-10, min(10, score))
if positive_count > negative_count:
if llm_summary:
summary = llm_summary
elif score > 0:
summary = f'近期{len(policy_news)}条政策消息,偏利好({positive_count}条利好/{negative_count}条利空)'
elif negative_count > positive_count:
elif score < 0:
summary = f'近期{len(policy_news)}条政策消息,偏利空({negative_count}条利空/{positive_count}条利好)'
else:
summary = f'近期{len(policy_news)}条政策消息,影响中性'
+25 -17
View File
@@ -101,7 +101,7 @@ def is_trading_time():
def get_all_users():
"""获取所有启用自动交易的用户"""
from db import get_db
from db import get_db, put_db
from psycopg2.extras import RealDictCursor
conn = get_db()
@@ -121,7 +121,7 @@ def get_all_users():
print(f"[定时任务] 获取用户列表失败: {e}")
return []
finally:
conn.close()
put_db(conn)
# 自定义股票列表(100只精选股票)
@@ -168,7 +168,7 @@ def execute_auto_trade_for_user(user_id, trade_quantity=1000, scan_date=None):
scan_date: 使用哪天的扫描数据, None则自动选择最近可用的
"""
from db import get_db
from db import get_db, put_db
from psycopg2.extras import RealDictCursor
print(f"[定时任务] 开始为用户{user_id}执行策略交易(统一推荐算法)...")
@@ -407,7 +407,7 @@ def execute_auto_trade_for_user(user_id, trade_quantity=1000, scan_date=None):
traceback.print_exc()
return {'error': str(e)}
finally:
conn.close()
put_db(conn)
def _get_latest_price(stock_code):
@@ -417,7 +417,7 @@ def _get_latest_price(stock_code):
def update_positions_price_for_user(user_id):
"""更新用户持仓的当前价格(收盘时调用)— 使用腾讯财经API"""
from db import get_db
from db import get_db, put_db
from psycopg2.extras import RealDictCursor
conn = get_db()
@@ -443,6 +443,8 @@ def update_positions_price_for_user(user_id):
for c in codes:
if c.startswith('6'):
tencent_codes.append(f'sh{c}')
elif c.startswith('8') or c.startswith('9'):
tencent_codes.append(f'bj{c}')
else:
tencent_codes.append(f'sz{c}')
_r = _req.get(f'http://qt.gtimg.cn/q={",".join(tencent_codes)}',
@@ -492,7 +494,7 @@ def update_positions_price_for_user(user_id):
conn.rollback()
print(f"[定时任务] 更新用户{user_id}持仓价格失败: {e}")
finally:
conn.close()
put_db(conn)
def job_morning_trade():
@@ -512,19 +514,22 @@ def job_morning_trade():
for user in users:
user_id = user['user_id']
# 尝试使用智能引擎
smart_conn = None
try:
from services.smart_trade_engine import execute_smart_trade
from db import get_db
conn = get_db()
if conn:
result = execute_smart_trade(conn, user_id, scan_date=None)
conn.close()
from db import get_db, put_db
smart_conn = get_db()
if smart_conn:
result = execute_smart_trade(smart_conn, user_id, scan_date=None)
if result.get('success'):
print(f"[定时任务] 用户{user_id} 智能引擎执行成功 "
f"(算法:{result.get('algo','?')}, 信号:{result.get('signals',0)})")
continue
except Exception as e:
print(f"[定时任务] 用户{user_id} 智能引擎异常,降级到旧引擎: {e}")
finally:
if smart_conn:
put_db(smart_conn)
# 降级:使用旧引擎
trade_quantity = user.get('trade_quantity') or 1000
@@ -551,18 +556,21 @@ def job_afternoon_trade():
for user in users:
user_id = user['user_id']
# 尝试使用智能引擎
smart_conn = None
try:
from services.smart_trade_engine import execute_smart_trade
from db import get_db
conn = get_db()
if conn:
result = execute_smart_trade(conn, user_id, scan_date=today)
conn.close()
from db import get_db, put_db
smart_conn = get_db()
if smart_conn:
result = execute_smart_trade(smart_conn, user_id, scan_date=today)
if result.get('success'):
print(f"[定时任务] 用户{user_id} 午后智能引擎执行成功")
continue
except Exception as e:
print(f"[定时任务] 用户{user_id} 智能引擎异常,降级: {e}")
finally:
if smart_conn:
put_db(smart_conn)
trade_quantity = user.get('trade_quantity') or 1000
execute_auto_trade_for_user(user_id, trade_quantity, scan_date=today)
@@ -645,7 +653,7 @@ def trigger_afternoon_trade():
def trigger_closing_update():
"""手动触发收盘更新(仅更新持仓价格)"""
from db import get_db
from db import get_db, put_db
if not is_trading_day():
print("[定时任务] 今天不是交易日,跳过")
return
+9 -3
View File
@@ -134,10 +134,16 @@ def compute_comprehensive_score(stock_code, stock_name, technical_score, df=None
external_score = 0
summaries = []
# ---- P0: 主力资金进出 ----
# ---- P0: 主力资金进出(带30分钟缓存,与批量模式一致)----
try:
from services.fund_flow_analyzer import analyze_fund_flow
fund_result = analyze_fund_flow(stock_code, days=5)
now = time.time()
cached_ff = _fund_flow_cache.get(stock_code)
if cached_ff and (now - cached_ff[1]) < _FUND_FLOW_TTL:
fund_result = cached_ff[0]
else:
from services.fund_flow_analyzer import analyze_fund_flow
fund_result = analyze_fund_flow(stock_code, days=5)
_fund_flow_cache[stock_code] = (fund_result, now)
factors['fund_flow'] = fund_result
external_score += fund_result.get('score', 0)
all_reasons.extend(fund_result.get('reasons', []))
+3 -1
View File
@@ -340,7 +340,9 @@ def detect_all_signals(df, lookback=5):
if not np.issubdtype(df[col].dtype, np.floating):
df[col] = pd.to_numeric(df[col], errors='coerce').fillna(0).astype(np.float64)
df = calc_all_indicators(df)
# 若已有指标列则跳过重复计算(deep_analyze 等场景预计算过)
if 'dif' not in df.columns or 'ma5' not in df.columns:
df = calc_all_indicators(df)
all_signals = []
all_signals.extend(detect_main_rising_wave(df, lookback))
+3 -3
View File
@@ -501,7 +501,7 @@ def generate_sell_decisions(conn, user_id, config, current_prices, scan_map=None
if current_shares <= 0:
continue
profit_pct_now = (price - buy_price) / buy_price * 100
profit_pct_now = (price - buy_price) / buy_price * 100 if buy_price > 0 else 0
days_held = pos.get('days_held', 0) or 0
print(f" {code} 成本{buy_price:.2f} 现价{price:.2f} 盈亏{profit_pct_now:+.1f}% 持仓{days_held}")
consec_up = pos.get('consecutive_up_days', 0) or 0
@@ -1003,7 +1003,7 @@ def execute_smart_trade(conn, user_id, scan_date=None):
log_signal(conn, user_id, today, code, pos['stock_name'],
'partial_sell', dec['reason'], dec['rule'],
price, pos['avg_cost'],
(price - pos['avg_cost']) / pos['avg_cost'] * 100,
(price - pos['avg_cost']) / pos['avg_cost'] * 100 if pos['avg_cost'] else 0,
True, price, shares)
results.append({
@@ -1052,7 +1052,7 @@ def execute_smart_trade(conn, user_id, scan_date=None):
log_signal(conn, user_id, today, code, pos['stock_name'],
'sell', dec['reason'], dec['rule'],
price, pos['avg_cost'],
(price - pos['avg_cost']) / pos['avg_cost'] * 100,
(price - pos['avg_cost']) / pos['avg_cost'] * 100 if pos['avg_cost'] else 0,
True, price, qty)
results.append({
+18 -13
View File
@@ -165,28 +165,30 @@ def get_kline_data(stock_code, days=120, use_local_db=True):
def _get_kline_from_local_db(stock_code, days=120):
"""从本地数据库读取K线(最快,毫秒级)"""
"""从本地数据库读取K线(最快,毫秒级)
使用 LIMIT 限制交易日条数与外部 API limit=days 行为一致
"""
import pandas as pd
try:
from db import get_db, put_db
conn = get_db()
if not conn:
return None
conn.autocommit = True
start_date = (datetime.now() - timedelta(days=days)).strftime('%Y-%m-%d')
try:
with conn.cursor() as cur:
cur.execute("""
SELECT trade_date, open, high, low, close, volume
FROM stock_kline_daily
WHERE code = %s AND trade_date >= %s
ORDER BY trade_date
""", (stock_code, start_date))
WHERE code = %s
ORDER BY trade_date DESC
LIMIT %s
""", (stock_code, min(days, 300)))
rows = cur.fetchall()
finally:
put_db(conn)
if rows and len(rows) >= 30:
rows = list(reversed(rows))
df = pd.DataFrame(rows, columns=['date', 'open', 'high', 'low', 'close', 'volume'])
df['date'] = df['date'].astype(str)
for col in ('open', 'high', 'low', 'close', 'volume'):
@@ -202,7 +204,9 @@ _thread_local = threading.local()
def get_kline_from_local_db_threaded(stock_code, days=120):
"""多线程扫描专用:使用线程本地连接从本地DB读取K线"""
"""多线程扫描专用:使用线程本地连接从本地DB读取K线
使用 LIMIT 限制交易日条数与外部 API limit=days 行为一致
"""
import pandas as pd
try:
conn = getattr(_thread_local, 'kline_conn', None)
@@ -215,17 +219,18 @@ def get_kline_from_local_db_threaded(stock_code, days=120):
conn.autocommit = True
_thread_local.kline_conn = conn
start_date = (datetime.now() - timedelta(days=days)).strftime('%Y-%m-%d')
with conn.cursor() as cur:
cur.execute("""
SELECT trade_date, open, high, low, close, volume
FROM stock_kline_daily
WHERE code = %s AND trade_date >= %s
ORDER BY trade_date
""", (stock_code, start_date))
WHERE code = %s
ORDER BY trade_date DESC
LIMIT %s
""", (stock_code, min(days, 300)))
rows = cur.fetchall()
if rows and len(rows) >= 30:
rows = list(reversed(rows))
df = pd.DataFrame(rows, columns=['date', 'open', 'high', 'low', 'close', 'volume'])
df['date'] = df['date'].astype(str)
for col in ('open', 'high', 'low', 'close', 'volume'):
@@ -962,8 +967,8 @@ def _generate_plain_summary(price, change_pct, ma_trend, position, supports,
parts.append(trend_desc + '')
# 2. 价格位置(用大白话)
pos_20 = position.get('20d', {})
pct_20 = pos_20.get('pct', 50)
pos_20 = position.get('d20', {})
pct_20 = pos_20.get('range_pct', 50)
if pct_20 > 80:
parts.append(f'当前股价处于近20天的高位区间({pct_20:.0f}%位置),已经涨了不少,追高要小心。')
elif pct_20 > 50:
+48 -12
View File
@@ -7,11 +7,16 @@ from datetime import datetime, timedelta
import traceback
import json
import os
import time
import threading
from config import Config
# ========== 股票名称缓存 ==========
_stock_name_cache = {}
_name_cache_dirty = False
_name_cache_last_save = 0.0
_name_cache_lock = threading.Lock()
def _load_stock_name_cache():
@@ -26,32 +31,61 @@ def _load_stock_name_cache():
print(f"加载股票名称缓存失败: {e}")
def _save_stock_name_cache():
"""保存股票名称缓存到本地"""
try:
with open(Config.STOCK_NAME_CACHE_FILE, 'w', encoding='utf-8') as f:
json.dump(_stock_name_cache, f, ensure_ascii=False, indent=2)
except Exception as e:
print(f"保存股票名称缓存失败: {e}")
def _save_stock_name_cache(force=False):
"""保存股票名称缓存到本地(延迟批量保存,60秒去抖)"""
global _name_cache_dirty, _name_cache_last_save
if not force and not _name_cache_dirty:
return
now = time.time()
if not force and (now - _name_cache_last_save) < 60:
return
with _name_cache_lock:
try:
with open(Config.STOCK_NAME_CACHE_FILE, 'w', encoding='utf-8') as f:
json.dump(_stock_name_cache, f, ensure_ascii=False, indent=2)
_name_cache_dirty = False
_name_cache_last_save = now
except Exception as e:
print(f"保存股票名称缓存失败: {e}")
def get_stock_name(stock_code):
"""获取股票名称 — 使用腾讯财经API"""
global _stock_name_cache
"""获取股票名称 — 优先内存缓存 → DB stock_realtime_price → 腾讯财经API"""
global _stock_name_cache, _name_cache_dirty
if stock_code in _stock_name_cache:
return _stock_name_cache[stock_code]
# 腾讯财经API获取股票名称
# 优先从数据库 stock_realtime_price 表查(毫秒级)
try:
from db import get_db, put_db
conn = get_db()
if conn:
try:
with conn.cursor() as cur:
cur.execute("SELECT name FROM stock_realtime_price WHERE code = %s", (stock_code,))
row = cur.fetchone()
if row and row[0]:
_stock_name_cache[stock_code] = row[0]
_name_cache_dirty = True
_save_stock_name_cache()
return row[0]
finally:
put_db(conn)
except Exception:
pass
# 回退到腾讯财经API
try:
import requests as _req
tcode = ('sh' if stock_code.startswith('6') else 'sz') + stock_code
tcode = ('sh' if stock_code.startswith('6') else 'bj' if stock_code.startswith(('8', '9')) else 'sz') + stock_code
_r = _req.get(f'http://qt.gtimg.cn/q={tcode}', timeout=5,
headers={'Referer': 'https://finance.qq.com'})
if _r.status_code == 200 and '\"' in _r.text:
_fields = _r.text.split('\"')[1].split('~')
if len(_fields) > 2 and _fields[1]:
_stock_name_cache[stock_code] = _fields[1]
_name_cache_dirty = True
_save_stock_name_cache()
return _fields[1]
except Exception as e:
@@ -119,6 +153,8 @@ def get_stock_fund_flow(stock_code, start_date, end_date, force_refresh=False):
market = 'sh'
elif stock_code.startswith('0') or stock_code.startswith('3'):
market = 'sz'
elif stock_code.startswith('8') or stock_code.startswith('9'):
market = 'bj'
else:
return None, None, "无法识别股票代码所属市场"
@@ -313,7 +349,7 @@ def get_realtime_price(stock_code):
# 备用方案2:使用腾讯财经API(腾讯云可用)
try:
import requests as _req
tcode = ('sh' if stock_code.startswith('6') else 'sz') + stock_code
tcode = ('sh' if stock_code.startswith('6') else 'bj' if stock_code.startswith(('8', '9')) else 'sz') + stock_code
_r = _req.get(f'http://qt.gtimg.cn/q={tcode}', timeout=5,
headers={'Referer': 'https://finance.qq.com'})
if _r.status_code == 200 and '\"' in _r.text:
+26
View File
@@ -655,6 +655,32 @@
.scan-pagination button:disabled { opacity: 0.3; cursor: not-allowed; }
.scan-pagination span { font-size: 13px; color: var(--text-secondary); }
/* 结果列表可收起 header */
.full-scan-list-section {
margin-top: 8px;
}
.scan-list-header {
display: flex;
justify-content: space-between;
align-items: center;
padding: 8px 12px;
cursor: pointer;
border-radius: 8px;
background: rgba(255,255,255,0.03);
transition: background 0.15s;
}
.scan-list-header:hover { background: rgba(255,255,255,0.06); }
.scan-list-title {
font-size: 13px;
font-weight: 600;
color: var(--text-secondary);
}
.scan-list-title small {
font-weight: 400;
color: var(--text-muted);
margin-left: 4px;
}
.admin-link-btn {
display: block;
text-align: center;
+3 -1
View File
@@ -86,10 +86,12 @@
fullScanLoading: false,
fullScanPage: 1,
fullScanTotalPages: 1,
fullScanPerPage: 5, // 每页显示5只股票
fullScanFilter: 'all',
fullScanFilterTypes: [],
fullScanFilterRecommend: '买入', // 扫描结果打开时缺省显示买入列表;all|买入|加仓|卖出|持有|关注|观察|观望
scanSummaryCollapsed: false,
scanResultsCollapsed: false, // 结果列表可收起
fullScanSignalDist: [],
// 找牛股(保留兼容)
@@ -2253,7 +2255,7 @@
const resp = await axios.get('/api/scan_results', {
params: {
page: this.fullScanPage,
per_page: 50,
per_page: this.fullScanPerPage,
min_triggered: minTriggered,
signal_type: signalTypes,
holding_codes: (this.holdingStocks || []).join(','),
+1
View File
@@ -60,6 +60,7 @@ def get_all_stock_codes(conn):
SELECT code, name FROM stock_realtime_price
WHERE volume > 0 AND price > 0
AND name NOT LIKE '%%退%%'
AND name NOT LIKE '%%ST%%'
AND name NOT LIKE 'PT%%'
ORDER BY code
""")
+14 -8
View File
@@ -243,8 +243,9 @@ def fetch_5min_kline_sina(code):
l = float(row.get('low', 0))
c = float(row.get('close', 0))
v = int(float(row.get('volume', 0)))
amt = float(row.get('amount', 0))
rows.append((code, dt, o, h, l, c, v, 0.0, 0.0))
rows.append((code, dt, o, h, l, c, v, amt, 0.0))
return rows if rows else None
except (IndexError, KeyError, ValueError):
@@ -366,7 +367,8 @@ def fetch_5min_kline_tencent(code):
dt = datetime(today.year, today.month, today.day, h, m, 0)
except (ValueError, IndexError):
continue
min_data.append((dt, price, vol))
amt = float(parts[3]) if len(parts) > 3 else 0.0
min_data.append((dt, price, vol, amt))
if not min_data:
return None
@@ -375,26 +377,30 @@ def fetch_5min_kline_tencent(code):
# 5分钟窗口: 09:30-09:35, 09:35-09:40, ...
from collections import defaultdict
bars = defaultdict(list)
for dt, price, vol in min_data:
for dt, price, vol, amt in min_data:
# 5分钟窗口起始时间
minute = dt.minute
bar_min = (minute // 5) * 5
bar_dt = dt.replace(minute=bar_min, second=0)
bars[bar_dt].append((price, vol))
bars[bar_dt].append((price, vol, amt))
rows = []
prev_vol = 0
prev_amt = 0
for bar_dt in sorted(bars.keys()):
ticks = bars[bar_dt]
o = ticks[0][0] # 第一个价格
c = ticks[-1][0] # 最后一个价格
h = max(p for p, _ in ticks)
l = min(p for p, _ in ticks)
# 腾讯的volume是累积值,取窗口最后的 - 窗口最前的之前
h = max(p for p, _, _ in ticks)
l = min(p for p, _, _ in ticks)
# 腾讯的volume和amount是累积值,取窗口最后的 - 前一窗口最后的
last_vol = ticks[-1][1]
bar_vol = last_vol - prev_vol if prev_vol > 0 else ticks[-1][1]
prev_vol = last_vol
rows.append((code, bar_dt, o, h, l, c, max(0, bar_vol), 0.0, 0.0))
last_amt = ticks[-1][2]
bar_amt = last_amt - prev_amt if prev_amt > 0 else ticks[-1][2]
prev_amt = last_amt
rows.append((code, bar_dt, o, h, l, c, max(0, bar_vol), max(0, bar_amt), 0.0))
return rows if rows else None
+11 -6
View File
@@ -25,7 +25,7 @@
<link rel="stylesheet" href="/static/css/components.css?v=20260718v1">
<link rel="stylesheet" href="/static/css/auth.css?v=20260718v1">
<link rel="stylesheet" href="/static/css/pages.css?v=20260719v2">
<link rel="stylesheet" href="/static/css/scan.css?v=20260719v5">
<link rel="stylesheet" href="/static/css/scan.css?v=20260722v1">
<link rel="stylesheet" href="/static/css/responsive.css?v=20260719v1">
</head>
<body>
@@ -883,7 +883,12 @@
</div>
</div>
</div>
<div class="full-scan-list">
<div class="full-scan-list-section">
<div class="scan-list-header" role="button" tabindex="0" :aria-expanded="!scanResultsCollapsed" @click="scanResultsCollapsed = !scanResultsCollapsed" @keyup.enter="scanResultsCollapsed = !scanResultsCollapsed">
<span class="scan-list-title">推荐列表 <small v-if="fullScanTotalPages > 1">{{ fullScanPage }}/{{ fullScanTotalPages }}页</small></span>
<span class="scan-summary-toggle">{{ scanResultsCollapsed ? '▶' : '▼' }}</span>
</div>
<div v-show="!scanResultsCollapsed" class="full-scan-list">
<div v-if="fullScanResultsFiltered.length === 0" class="scan-empty-hint">暂无符合条件的股票</div>
<div v-for="item in fullScanResultsFiltered" :key="item.code" class="scan-result-card"
:class="{ 'has-signal': item.triggered_count > 0, 'is-watched': searchHistory.some(h => h.code === item.code) }">
@@ -915,11 +920,11 @@
</div>
<div class="scan-holding-note" v-if="item.holding_note">{{ item.holding_note }}</div>
</div>
</div>
<div v-if="fullScanTotalPages > 1" class="scan-pagination">
<button type="button" @click="fetchFullScanResults(fullScanPage - 1)" :disabled="fullScanPage <= 1">上一页</button>
<button type="button" @click.stop="fetchFullScanResults(fullScanPage - 1)" :disabled="fullScanPage <= 1">上一页</button>
<span>{{ fullScanPage }} / {{ fullScanTotalPages }}</span>
<button type="button" @click="fetchFullScanResults(fullScanPage + 1)" :disabled="fullScanPage >= fullScanTotalPages">下一页</button>
<button type="button" @click.stop="fetchFullScanResults(fullScanPage + 1)" :disabled="fullScanPage >= fullScanTotalPages">下一页</button>
</div>
</div>
</template>
</div>
@@ -2299,6 +2304,6 @@
</div>
{% endraw %}
<script src="/static/js/app.js?v=20260718v2"></script>
<script src="/static/js/app.js?v=20260722v1"></script>
</body>
</html>