akg-factor-bridge/sources.py

165 lines
7.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""候选卡的取数层:把各条证据线从三处库读成"按前缀码索引的字典"(全部只读)。
## 为什么要有它
候选卡的规则在 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 {}
try:
marks = ",".join(["%s"] * len(codes))
df = db.read_mysql(
"pms",
f"SELECT stock_code, trade_date, signal_type, support_level, pressure_level, "
f"raw_logic_json FROM strategy_daily_results WHERE stock_code IN ({marks})",
tuple(codes))
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