refactor: BSP检测改为新笔驱动,不再逐tick对比BSP列表
- 用 last_bi_id (start_klc.start_time) 跟踪最后一笔 - 新笔确认时检查上一笔终点是否为 BSP → 推送 - 中枢更新同样在新笔产生时触发 - 移除时间过滤、BSP列表diff、持久化去重等冗余逻辑 - 无新笔时快速跳过,tick从40s降到15s Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
cf9097a540
commit
78d02cf2ef
+40
-37
@@ -5,8 +5,8 @@ main.py - 缠论买卖点监控主程序。
|
|||||||
每整分钟:
|
每整分钟:
|
||||||
1. 从 data_provider 拉取所有币对最新 1m K 线
|
1. 从 data_provider 拉取所有币对最新 1m K 线
|
||||||
2. 每个币对独立跑缠论管线
|
2. 每个币对独立跑缠论管线
|
||||||
3. 检测新出现的买卖点(BSP)+ 中枢特征变化
|
3. 检测新笔 → 检查上一笔终点是否 BSP → 推送
|
||||||
4. 推送到 Telegram
|
4. 检测中枢特征变化 → 推送
|
||||||
"""
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
@@ -15,6 +15,7 @@ import os
|
|||||||
import time
|
import time
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from datetime import datetime, timezone, timedelta
|
from datetime import datetime, timezone, timedelta
|
||||||
|
from typing import Optional
|
||||||
|
|
||||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
|
||||||
@@ -39,12 +40,16 @@ def _short(symbol: str) -> str:
|
|||||||
return symbol.split(":")[0].replace("/", "")
|
return symbol.split(":")[0].replace("/", "")
|
||||||
|
|
||||||
|
|
||||||
|
def _bi_id(bi) -> tuple:
|
||||||
|
"""笔的稳定标识,基于首K线时间戳。"""
|
||||||
|
return (bi.start_klc.start_time,)
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class SymbolState:
|
class SymbolState:
|
||||||
symbol: str
|
symbol: str
|
||||||
pivot_monitor: ChanPivotMonitor = field(default_factory=ChanPivotMonitor)
|
pivot_monitor: ChanPivotMonitor = field(default_factory=ChanPivotMonitor)
|
||||||
known_bsp_keys: set = field(default_factory=set)
|
last_bi_id: Optional[tuple] = None # 上次检查过的最后一笔 ID
|
||||||
last_df_ts: object = None
|
|
||||||
first_run: bool = True
|
first_run: bool = True
|
||||||
|
|
||||||
|
|
||||||
@@ -80,11 +85,6 @@ class BSPMonitor:
|
|||||||
logger.warning(f"[{name}] DataFrame 为空")
|
logger.warning(f"[{name}] DataFrame 为空")
|
||||||
return
|
return
|
||||||
|
|
||||||
latest_ts = df.iloc[-1]["timestamp"]
|
|
||||||
if st.last_df_ts and latest_ts <= st.last_df_ts:
|
|
||||||
return # 无新K线
|
|
||||||
st.last_df_ts = latest_ts
|
|
||||||
|
|
||||||
# 2. 运行缠论管线
|
# 2. 运行缠论管线
|
||||||
try:
|
try:
|
||||||
engine = ChanEngine(df)
|
engine = ChanEngine(df)
|
||||||
@@ -92,23 +92,46 @@ class BSPMonitor:
|
|||||||
logger.error(f"[{name}] 缠论计算失败: {e}", exc_info=True)
|
logger.error(f"[{name}] 缠论计算失败: {e}", exc_info=True)
|
||||||
return
|
return
|
||||||
|
|
||||||
# 3. 中枢特征监控 + 推送
|
# 3. 获取已确认的笔
|
||||||
pivot_state = st.pivot_monitor.update(engine.bi_zs_list)
|
confirmed = [b for b in engine.bi_list if b.is_sure]
|
||||||
|
if len(confirmed) < 2:
|
||||||
|
return
|
||||||
|
|
||||||
# 4. BSP 检测 + 推送
|
last_bi = confirmed[-1]
|
||||||
current_bsps = engine.bsp_list
|
current_bi_id = _bi_id(last_bi)
|
||||||
current_keys = {_bsp_stable_key(b, symbol) for b in current_bsps}
|
|
||||||
|
|
||||||
|
# 4. 首轮:记录状态,不推送
|
||||||
if st.first_run:
|
if st.first_run:
|
||||||
st.first_run = False
|
st.first_run = False
|
||||||
st.known_bsp_keys = current_keys
|
st.last_bi_id = current_bi_id
|
||||||
confirmed = [b for b in engine.bi_list if b.is_sure]
|
st.pivot_monitor.update(engine.bi_zs_list) # 初始化中枢状态
|
||||||
logger.info(
|
logger.info(
|
||||||
f"[{name}] 首次完成 — {len(confirmed)} 笔, "
|
f"[{name}] 首次完成 — {len(confirmed)} 笔, "
|
||||||
f"{len(current_bsps)} BSP (不推送)"
|
f"{len(engine.bsp_list)} BSP (不推送)"
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# 5. 检测新笔
|
||||||
|
if current_bi_id == st.last_bi_id:
|
||||||
|
return # 无新笔,跳过
|
||||||
|
|
||||||
|
st.last_bi_id = current_bi_id
|
||||||
|
logger.info(f"[{name}] 新笔确认 — #{len(confirmed)} "
|
||||||
|
f"{'⬆️' if last_bi.dir.value == 1 else '⬇️'} "
|
||||||
|
f"高度: ${last_bi.height:.2f}")
|
||||||
|
|
||||||
|
# 6. 检查上一笔终点是否为 BSP
|
||||||
|
prev_bi = confirmed[-2]
|
||||||
|
bsp = engine.get_bsp_for_bi(prev_bi)
|
||||||
|
if bsp:
|
||||||
|
key = f"{symbol}_{bsp.type}_{bsp.klc.end_time}"
|
||||||
|
msg = engine.format_bsp_detail(bsp, symbol)
|
||||||
|
msg = _escape_html(msg)
|
||||||
|
if send_bsp_alert(msg, bsp_key=key):
|
||||||
|
logger.info(f"[{name}] ✅ BSP: {key}")
|
||||||
|
|
||||||
|
# 7. 中枢特征更新
|
||||||
|
pivot_state = st.pivot_monitor.update(engine.bi_zs_list)
|
||||||
if pivot_state:
|
if pivot_state:
|
||||||
logger.info(
|
logger.info(
|
||||||
f"[{name}] 中枢更新: bi_count={pivot_state['bi_count']} "
|
f"[{name}] 中枢更新: bi_count={pivot_state['bi_count']} "
|
||||||
@@ -117,22 +140,6 @@ class BSPMonitor:
|
|||||||
)
|
)
|
||||||
self._push_pivot(pivot_state, name)
|
self._push_pivot(pivot_state, name)
|
||||||
|
|
||||||
new_keys = current_keys - st.known_bsp_keys
|
|
||||||
st.known_bsp_keys |= current_keys
|
|
||||||
|
|
||||||
if not new_keys:
|
|
||||||
return
|
|
||||||
|
|
||||||
logger.info(f"[{name}] {len(new_keys)} 个新 BSP")
|
|
||||||
for bsp in current_bsps:
|
|
||||||
key = _bsp_stable_key(bsp, symbol)
|
|
||||||
if key not in new_keys:
|
|
||||||
continue
|
|
||||||
msg = engine.format_bsp_detail(bsp, symbol)
|
|
||||||
msg = _escape_html(msg)
|
|
||||||
if send_bsp_alert(msg, bsp_key=key):
|
|
||||||
logger.info(f"[{name}] ✅ BSP: {key}")
|
|
||||||
|
|
||||||
def _push_pivot(self, state: dict, name: str):
|
def _push_pivot(self, state: dict, name: str):
|
||||||
if state["is_sure"]:
|
if state["is_sure"]:
|
||||||
phase = "✅ 已确认"
|
phase = "✅ 已确认"
|
||||||
@@ -197,10 +204,6 @@ class BSPMonitor:
|
|||||||
await asyncio.sleep(5)
|
await asyncio.sleep(5)
|
||||||
|
|
||||||
|
|
||||||
def _bsp_stable_key(bsp, symbol: str) -> str:
|
|
||||||
return f"{symbol}_{bsp.type}_{bsp.klc.end_time}"
|
|
||||||
|
|
||||||
|
|
||||||
def _escape_html(msg: str) -> str:
|
def _escape_html(msg: str) -> str:
|
||||||
"""HTML 转义,保留已有的 <b>/<code> 标签。"""
|
"""HTML 转义,保留已有的 <b>/<code> 标签。"""
|
||||||
msg = msg.replace("&", "&")
|
msg = msg.replace("&", "&")
|
||||||
|
|||||||
Reference in New Issue
Block a user