2026-09-02 16:39:52 +08:00
|
|
|
|
"""候选卡的取数层:把各条证据线从三处库读成"按前缀码索引的字典"(全部只读)。
|
|
|
|
|
|
|
|
|
|
|
|
## 为什么要有它
|
|
|
|
|
|
|
|
|
|
|
|
候选卡的规则在 card.py,是纯函数;证据从哪来、怎么对齐代码格式、缺了怎么办,
|
|
|
|
|
|
全部收在这里,plan.py 只做装配。三处来源:
|
|
|
|
|
|
|
|
|
|
|
|
基座 PG v_factor_transmission_moved 已启动成员及其所在环节的传导证据(只认 Segment 目标)
|
|
|
|
|
|
v_factor_stock_daily 数据日涨幅、主力净额异常值、热度变化
|
|
|
|
|
|
153 代理 strategy_daily_results 决策系统昨夜结论:信号、支撑压力位、吸筹块
|
|
|
|
|
|
(吸筹的评分、状态、评分日三项一次读齐;基座落库的吸筹版本没有评分日,所以不从基座取)
|
|
|
|
|
|
平台 MySQL gp_day_data 交易日历(算评分日龄用;取一只长期存在的票的日期序列)
|
|
|
|
|
|
|
|
|
|
|
|
代码格式:基座是点后缀式 600000.SH,决策系统与桥是前缀式 SH600000,进出都过 common.to_prefix。
|
|
|
|
|
|
读失败的语义:每一路读不到都返回空字典并打印一行原因,候选卡按"缺失"处理(进关注或仅展示),
|
|
|
|
|
|
不让计划断产——与 pool.py 的安全边界一致。
|
|
|
|
|
|
"""
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
|
|
import datetime as dt
|
|
|
|
|
|
import json
|
|
|
|
|
|
|
|
|
|
|
|
import pandas as pd
|
|
|
|
|
|
|
|
|
|
|
|
import common
|
|
|
|
|
|
import db
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def moved_members(ds: str) -> dict[str, dict]:
|
|
|
|
|
|
"""数据日 ds 被传导指向的环节里,已启动(不在未动名单)的成员。
|
|
|
|
|
|
同一票挂在多个被指向环节上时,取源数最多、其次链符最高的那条作卡上的证据。"""
|
|
|
|
|
|
try:
|
|
|
|
|
|
df = db.read_pg(
|
|
|
|
|
|
"SELECT ts_code, target, n_sources, chain_fit, members_total, moved, "
|
|
|
|
|
|
"moved_ratio, mkt_trade_date FROM v_factor_transmission_moved "
|
|
|
|
|
|
"WHERE scan_date = %s", (ds,))
|
|
|
|
|
|
except Exception as e: # noqa: BLE001
|
|
|
|
|
|
print(f" (已动成员视图读取失败,候选卡的传导门槛整体缺席: {e!r})")
|
|
|
|
|
|
return {}
|
|
|
|
|
|
if df.empty:
|
|
|
|
|
|
return {}
|
|
|
|
|
|
df["k"] = df["ts_code"].map(lambda s: common.to_prefix(str(s).strip()))
|
|
|
|
|
|
df["n_sources"] = pd.to_numeric(df["n_sources"], errors="coerce").fillna(0)
|
|
|
|
|
|
df["chain_fit"] = pd.to_numeric(df["chain_fit"], errors="coerce").fillna(0)
|
|
|
|
|
|
df = df.sort_values(["n_sources", "chain_fit"], ascending=False).drop_duplicates("k")
|
|
|
|
|
|
out = {}
|
|
|
|
|
|
for r in df.itertuples():
|
|
|
|
|
|
out[r.k] = {"theme": str(r.target), "n_sources": int(r.n_sources),
|
|
|
|
|
|
"chain_fit": float(r.chain_fit),
|
|
|
|
|
|
"members_total": None if pd.isna(r.members_total) else int(r.members_total),
|
|
|
|
|
|
"moved": None if pd.isna(r.moved) else int(r.moved),
|
|
|
|
|
|
"mkt_trade_date": None if pd.isna(r.mkt_trade_date) else str(r.mkt_trade_date)}
|
|
|
|
|
|
return out
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def stock_daily(ds: str) -> dict[str, dict]:
|
|
|
|
|
|
"""数据日 ds 的个股行情三列:涨幅(百分数)、主力净额异常值、热度五日变化。"""
|
|
|
|
|
|
try:
|
|
|
|
|
|
df = db.read_pg(
|
|
|
|
|
|
"SELECT ts_code, pct_change, net_z, heat_chg FROM v_factor_stock_daily "
|
|
|
|
|
|
"WHERE trade_date = %s", (ds,))
|
|
|
|
|
|
except Exception as e: # noqa: BLE001
|
|
|
|
|
|
print(f" (个股日行情视图读取失败,候选卡的已启动门槛整体缺席: {e!r})")
|
|
|
|
|
|
return {}
|
|
|
|
|
|
out = {}
|
|
|
|
|
|
for r in df.itertuples():
|
|
|
|
|
|
k = common.to_prefix(str(r.ts_code).strip())
|
|
|
|
|
|
out[k] = {"pct0": _f(r.pct_change), "net_z": _f(r.net_z), "heat_chg": _f(r.heat_chg)}
|
|
|
|
|
|
return out
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def trading_days(upto: str, back_days: int = 90) -> list[str]:
|
|
|
|
|
|
"""交易日历:取一只长期存在的票在 gp_day_data 里的日期序列(表有一千四百万行,
|
|
|
|
|
|
全表 DISTINCT 太慢;按 symbol 走索引)。返回升序 ISO 日期串。"""
|
|
|
|
|
|
start = (dt.date.fromisoformat(upto) - dt.timedelta(days=back_days)).isoformat()
|
|
|
|
|
|
try:
|
|
|
|
|
|
df = db.read_mysql(
|
|
|
|
|
|
"factor", "SELECT DISTINCT DATE(`timestamp`) d FROM gp_day_data "
|
|
|
|
|
|
"WHERE symbol = %s AND `timestamp` >= %s AND `timestamp` <= %s "
|
|
|
|
|
|
"ORDER BY d", ("SH600519", start, upto))
|
|
|
|
|
|
return [pd.Timestamp(x).date().isoformat() for x in df["d"].tolist()]
|
|
|
|
|
|
except Exception as e: # noqa: BLE001
|
|
|
|
|
|
print(f" (交易日历读取失败,评分日龄按自然日×5/7 近似: {e!r})")
|
|
|
|
|
|
return []
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def night_conclusions(codes, ds: str) -> dict[str, dict]:
|
|
|
|
|
|
"""决策系统昨夜结论(每票最新一行):信号、支撑压力位、吸筹块与评分日龄。
|
|
|
|
|
|
|
|
|
|
|
|
评分日龄 = 结论行的 trade_date 到数据日 ds 之间的交易日数(含头不含尾)。
|
|
|
|
|
|
该表会被盘中补扫就地改写、没有落库时刻列,日龄以行的 trade_date 为准,
|
|
|
|
|
|
这是已知局限(方案 2.2 第三项)。"""
|
|
|
|
|
|
codes = sorted({common.to_prefix(str(c).strip()) for c in codes if c})
|
|
|
|
|
|
if not codes:
|
|
|
|
|
|
return {}
|
2026-09-02 16:53:13 +08:00
|
|
|
|
# 只取不晚于数据日 ds 的结论行:生产上 ds 就是数据日,与"取最新一行"等价;
|
|
|
|
|
|
# 复盘重建历史日时这一条防前视(否则历史日会读到今天的吸筹状态)。
|
|
|
|
|
|
# trade_date 是整数 YYYYMMDD(数据源盘点 §1a),直接比大小。
|
|
|
|
|
|
ds_int = int(ds.replace("-", "")) if ds else 99999999
|
2026-09-02 16:39:52 +08:00
|
|
|
|
try:
|
|
|
|
|
|
marks = ",".join(["%s"] * len(codes))
|
|
|
|
|
|
df = db.read_mysql(
|
|
|
|
|
|
"pms",
|
|
|
|
|
|
f"SELECT stock_code, trade_date, signal_type, support_level, pressure_level, "
|
2026-09-02 16:53:13 +08:00
|
|
|
|
f"raw_logic_json FROM strategy_daily_results "
|
|
|
|
|
|
f"WHERE stock_code IN ({marks}) AND trade_date <= %s",
|
|
|
|
|
|
tuple(codes) + (ds_int,))
|
2026-09-02 16:39:52 +08:00
|
|
|
|
except Exception as e: # noqa: BLE001
|
|
|
|
|
|
print(f" (决策系统结论表读取失败,吸筹确认线与坏信号风险整体缺席: {e!r})")
|
|
|
|
|
|
return {}
|
|
|
|
|
|
if df.empty:
|
|
|
|
|
|
return {}
|
|
|
|
|
|
df = df.sort_values("trade_date").drop_duplicates("stock_code", keep="last")
|
|
|
|
|
|
cal = trading_days(ds)
|
|
|
|
|
|
cal_index = {d: i for i, d in enumerate(cal)}
|
|
|
|
|
|
out = {}
|
|
|
|
|
|
for r in df.itertuples():
|
|
|
|
|
|
k = common.to_prefix(str(r.stock_code).strip())
|
|
|
|
|
|
tdate = _ymd(r.trade_date)
|
|
|
|
|
|
age = _age(tdate, ds, cal_index)
|
|
|
|
|
|
ff = {}
|
|
|
|
|
|
try:
|
|
|
|
|
|
raw = json.loads(r.raw_logic_json) if isinstance(r.raw_logic_json, str) else (r.raw_logic_json or {})
|
|
|
|
|
|
ff = (raw or {}).get("fund_flow") or {}
|
|
|
|
|
|
except Exception: # noqa: BLE001 —— 坏 JSON 当无吸筹块
|
|
|
|
|
|
ff = {}
|
|
|
|
|
|
out[k] = {
|
|
|
|
|
|
"signal": str(r.signal_type or "").strip().upper() or None,
|
|
|
|
|
|
"support": _f(r.support_level), "pressure": _f(r.pressure_level),
|
|
|
|
|
|
"conclusion_date": tdate,
|
|
|
|
|
|
"accum_state": str(ff.get("state") or "") or None,
|
|
|
|
|
|
"accum_score": _f(ff.get("score")),
|
|
|
|
|
|
"accum_structure": ff.get("structure"), "accum_pos_tag": ff.get("pos_tag"),
|
|
|
|
|
|
"accum_age": age,
|
|
|
|
|
|
}
|
|
|
|
|
|
return out
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _ymd(v) -> str | None:
|
|
|
|
|
|
"""strategy_daily_results.trade_date 是整数 YYYYMMDD(也可能是日期),统一成 ISO 串。"""
|
|
|
|
|
|
if v is None or (isinstance(v, float) and pd.isna(v)):
|
|
|
|
|
|
return None
|
|
|
|
|
|
s = str(v).strip()
|
|
|
|
|
|
if len(s) == 8 and s.isdigit():
|
|
|
|
|
|
return f"{s[:4]}-{s[4:6]}-{s[6:]}"
|
|
|
|
|
|
try:
|
|
|
|
|
|
return pd.Timestamp(s).date().isoformat()
|
|
|
|
|
|
except Exception: # noqa: BLE001
|
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _age(tdate: str | None, ds: str, cal_index: dict) -> int | None:
|
|
|
|
|
|
if not tdate:
|
|
|
|
|
|
return None
|
|
|
|
|
|
if tdate in cal_index and ds in cal_index:
|
|
|
|
|
|
return cal_index[ds] - cal_index[tdate]
|
|
|
|
|
|
try: # 日历缺失或日期在日历之外:自然日 × 5/7 近似
|
|
|
|
|
|
nat = (dt.date.fromisoformat(ds) - dt.date.fromisoformat(tdate)).days
|
|
|
|
|
|
return max(0, round(nat * 5 / 7))
|
|
|
|
|
|
except ValueError:
|
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _f(v):
|
|
|
|
|
|
try:
|
|
|
|
|
|
x = float(v)
|
|
|
|
|
|
except (TypeError, ValueError):
|
|
|
|
|
|
return None
|
|
|
|
|
|
return None if x != x else x
|