akg-factor-bridge/logic_state_daily.py

283 lines
14 KiB
Python
Raw Permalink Normal View History

"""逐票日频逻辑状态表:每个数据日为档位表里的每只票记一行原始态与落定态,攒版本史并做抗抖动。
## 为什么要有它2026-09-07下一阶段方案第三件桥侧前置
逻辑状态四态logic_state.py在计划装配里已经随每张卡产出但有两件事它自己做不到
抗抖动settle 那个函数写好了可它要看"昨天落定的是什么、近几天原始态是什么"
这些只有存下来才有不存四态一天一变PMS 按它分流就会跟着一天一变
给持仓票用四态唯一真正有用的地方是持仓而持仓票在候选筛选第一步就被整行剔掉
/plan 里根本没有它接口 /logic_state 要能回答"这只在持的票今天证据还在不在"
就得有一张按票按日可查的表
## 落点、表名、键
写在平台因子库 t_factor_akg_* 的同一个 MySQL与行业观点快照 t_akg_judgement_snapshot
同库同理由桥对数据基座只有只读账号这张表又是选股系统自己的派生记录表名默认
t_akg_logic_state_daily不带 t_factor_ 前缀免得平台的因子清单把它当因子表收进去
键是数据日前缀码数据日就是计划文件名里那个日期plan_<ds> t_factor_akg_*
trade_date 同口径不是行业观点快照的 plan_date那个是写入当天两张表的日期口径不同
读的时候别混计划在 D+1 凌晨构建数据日是 D这张表记的是 D
## 只在早上那一次生成里写
写入点只有 plan.generaterun.py plan调度中心的默认三步之一/plan 每次实时重算会读这张表
做抗抖动但不写实时重算一天几十次写进去会把"当天落定态"改成"最后一次有人查时的状态"
派生数据就不可回溯了重跑同一天幂等先删该日行再整批插
## 落定的口径settle_one
上一次落定态 = 表里这只票严格早于数据日回看窗口之内最近一行的 state
近几日原始态 = 那几行的 raw_state 加上今天的原始态最新的在最后
落定规则在 logic_state.settle进入逻辑存疑即刻成立退出存疑要连续几个计划日不再存疑
其余迁移要连续几个计划日同向天数由 config 一次定死不按复盘读数回调
表里没有这只票第一次见到或超出回看窗口落定态就是原始态没有历史就没有抖动可抗
## 持仓票不在档位表里时的局限(写进设计不回避)
这张表只记档位表里的票一只在持的票掉出档位表之后每天没有新行接口现算时能用的历史只有
它掉出去之前那几行掉出去超过回看窗口落定态就等于原始态要让持仓票也逐日攒历史
得让桥知道 PMS 的持仓那是跨层的事本阶段不做先在接口返回里用 source=computed 标明
离线单测见 test_logic_state_daily.py不连库
"""
from __future__ import annotations
import datetime as dt
import json
import config
import db
import judgement
import logic_state
import sources
# 落库列的顺序(插入语句按此顺序拼参数)。
COLUMNS = ("trade_date", "code", "raw_state", "raw_why", "state", "settle_note", "prev_state",
"as_of", "usable", "missing", "paths", "reasons",
"verdict", "card_rank", "plan_version", "snapshot_at")
# 主键(数据日,前缀码):一天一票最多一行,重跑覆盖;另建(前缀码,数据日)索引,按票回看走它。
_CREATE_TABLE = """
CREATE TABLE IF NOT EXISTS {t} (
trade_date DATE NOT NULL,
code VARCHAR(16) NOT NULL,
raw_state VARCHAR(16),
raw_why VARCHAR(16),
state VARCHAR(16),
settle_note VARCHAR(128),
prev_state VARCHAR(16),
as_of DATE,
usable VARCHAR(64),
missing VARCHAR(64),
paths TEXT,
reasons TEXT,
verdict VARCHAR(16),
card_rank INT,
plan_version VARCHAR(32),
snapshot_at DATETIME,
PRIMARY KEY (trade_date, code),
KEY idx_code (code, trade_date)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
"""
# 四路各自的分隔usable / missing 两列存"研报论断、券商行动"这种顿号串,读回来再拆。
_SEP = ""
# ============================================================================
# 抗抖动:一只票的原始态配上历史,得到落定态
# ============================================================================
def settle_one(code: str, raw: dict, hist: list | None, *, confirm_days: int | None = None,
exit_days: int | None = None) -> dict:
"""纯函数:把一只票今天的原始合成结果,配上它在本表里的近日行,落定成今天的状态。
raw logic_state.compose 的返回state 是原始态hist 是这只票严格早于今天
按数据日升序的近日行history 的返回每行至少有 raw_state state 两列
没有历史传 None 或空列表那时落定态就是原始态
返回一个新字典state 改成落定态原始态挪到 raw_state另带 settle_note prev_state
其余键whypathsusablemissingreasonsas_of原样保留不改传入的 raw
"""
confirm = int(config.LOGIC_SETTLE_CONFIRM_DAYS if confirm_days is None else confirm_days)
exit_n = int(config.LOGIC_SETTLE_EXIT_DAYS if exit_days is None else exit_days)
rows = [r for r in (hist or []) if isinstance(r, dict)]
raw_state = raw.get("state") if isinstance(raw, dict) else None
prev_state = rows[-1].get("state") if rows else None
need = max(confirm, exit_n) - 1
recent = [r.get("raw_state") for r in (rows[-need:] if need > 0 else [])] + [raw_state]
settled, note = logic_state.settle(prev_state, raw_state, recent,
confirm_days=confirm, exit_doubt_days=exit_n)
out = dict(raw or {})
out.update(raw_state=raw_state, state=settled, settle_note=note, prev_state=prev_state)
return out
# ============================================================================
# 读:历史行(做抗抖动)、当日行(给接口)
# ============================================================================
def _in_clause(codes) -> tuple[str, tuple]:
codes = [str(c).strip() for c in (codes or []) if str(c).strip()]
return ",".join(["%s"] * len(codes)), tuple(codes)
def history(ds: str, codes=None, read_mysql=None) -> dict:
"""每只票严格早于数据日、回看窗口之内的行,按前缀码索引、行按数据日升序。
codes 不传就取全部票计划装配一轮要看一千多只一次查询取回不逐票查传了就只取
这几只接口用只取抗抖动要读的几列表还没建或读不到返回空字典并打印一行原因
那样全部票按第一次见到处理落定态等于原始态计划照出不断产
"""
table = config.LOGIC_STATE_TABLE
reader = read_mysql or db.read_mysql
try:
start = (dt.date.fromisoformat(ds)
- dt.timedelta(days=config.LOGIC_STATE_LOOKBACK_DAYS)).isoformat()
except ValueError:
print(f" (数据日 {ds!r} 不是合法日期,逻辑状态不做抗抖动)")
return {}
sql = (f"SELECT trade_date, code, raw_state, state FROM {table} "
f"WHERE trade_date >= %s AND trade_date < %s")
params: tuple = (start, ds)
if codes is not None:
marks, vals = _in_clause(codes)
if not vals:
return {}
sql += f" AND code IN ({marks})"
params = params + vals
try:
rows = sources._records(reader("factor", sql + " ORDER BY code, trade_date", params)) # noqa: SLF001
except Exception as e: # noqa: BLE001 —— 头一次跑时表还不存在,属正常
print(f" {table} 读不到历史行,逻辑状态本次不做抗抖动,落定态等于原始态: {e!r}")
return {}
out: dict = {}
for r in rows:
k = str(r.get("code") or "").strip()
if k:
out.setdefault(k, []).append({
"trade_date": sources._ymd(r.get("trade_date")), # noqa: SLF001
"raw_state": judgement._blank_to_none(r.get("raw_state")), # noqa: SLF001
"state": judgement._blank_to_none(r.get("state"))}) # noqa: SLF001
return out
def lookup(codes, ds: str, read_mysql=None) -> dict:
"""这几只票在数据日当天的行,按前缀码索引;没写过、表没建、读不到都返回空字典。"""
table = config.LOGIC_STATE_TABLE
reader = read_mysql or db.read_mysql
marks, vals = _in_clause(codes)
if not vals:
return {}
try:
rows = sources._records(reader( # noqa: SLF001
"factor", f"SELECT * FROM {table} WHERE trade_date = %s AND code IN ({marks})",
(ds,) + vals))
except Exception as e: # noqa: BLE001
print(f" {table} 当日行读取失败: {e!r}")
return {}
out = {}
for r in rows:
k = str(r.get("code") or "").strip()
if k:
out[k] = {c: judgement._blank_to_none(v) for c, v in dict(r).items()} # noqa: SLF001
return out
def row_to_out(row: dict) -> dict:
"""把表里的一行还原成发给下游的形状(与 plan._state_out 同形,多带 raw_state 等三键)。"""
def _split(v):
s = str(v or "").strip()
return [x for x in s.split(_SEP) if x] if s else []
def _loads(v, default):
if v is None or v == "":
return default
try:
return json.loads(v) if isinstance(v, str) else v
except (TypeError, ValueError):
return default
return {"state": row.get("state"), "raw_state": row.get("raw_state"), "why": row.get("raw_why"),
"settle_note": row.get("settle_note"), "prev_state": row.get("prev_state"),
"as_of": sources._ymd(row.get("as_of")), # noqa: SLF001
"usable": _split(row.get("usable")), "missing": _split(row.get("missing")),
"reasons": _loads(row.get("reasons"), []), "paths": _loads(row.get("paths"), [])}
# ============================================================================
# 写:只在 plan.generate 里调
# ============================================================================
def build_rows(ds: str, plan_rows: list, plan_version: str | None = None,
now: str | None = None) -> list:
"""纯函数把计划快照的行plan.collect 里 _row 的输出,带 logic_state拼成可落库的行。
没有 logic_state 的行档位表里有票但没装配出卡跳过同一只票出现多次只留最后一条"""
stamp = now or dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
by_code = {}
for r in plan_rows or []:
st = r.get("logic_state") if isinstance(r, dict) else None
code = str((r or {}).get("code") or "").strip()
if not code or not isinstance(st, dict):
continue
by_code[code] = (r, st)
out = []
for code in sorted(by_code):
r, st = by_code[code]
out.append({
"trade_date": ds, "code": code[:16],
"raw_state": st.get("raw_state") or st.get("state"), "raw_why": st.get("why"),
"state": st.get("state"), "settle_note": (st.get("settle_note") or "")[:128] or None,
"prev_state": st.get("prev_state"), "as_of": st.get("as_of"),
"usable": _SEP.join(st.get("usable") or [])[:64] or None,
"missing": _SEP.join(st.get("missing") or [])[:64] or None,
"paths": json.dumps(st.get("paths") or [], ensure_ascii=False, default=str),
"reasons": json.dumps((st.get("reasons") or [])[:6], ensure_ascii=False, default=str),
"verdict": r.get("verdict"),
"card_rank": r.get("card_rank"),
"plan_version": (plan_version or "")[:32] or None,
"snapshot_at": stamp,
})
return out
def save(ds: str, rows: list, conn_factory=None) -> None:
"""幂等落库:建表(已存在就跳过)、删该日行、整批插,删与插同一事务。与 judgement.save 同范式。
rows 为空也照样删该日行那表示这一天一张卡都没装配出来不该留上一次重跑的残行"""
table = config.LOGIC_STATE_TABLE
factory = conn_factory or db.factor_conn
cols = ",".join(COLUMNS)
marks = ",".join(["%s"] * len(COLUMNS))
payload = [tuple(r.get(c) for c in COLUMNS) for r in rows]
with factory() as conn:
with conn.cursor() as cur:
cur.execute(_CREATE_TABLE.format(t=table))
conn.commit()
with conn.cursor() as cur:
cur.execute(f"DELETE FROM {table} WHERE trade_date = %s", (ds,))
if payload:
cur.executemany(f"INSERT INTO {table} ({cols}) VALUES ({marks})", payload)
conn.commit()
def persist(ds: str, plan_rows: list, plan_version: str | None = None, write=None) -> dict:
"""plan.generate 的一步从快照行拼表行、幂等写该日、打印读数。write 可注入(离线单测)。"""
rows = build_rows(ds, plan_rows, plan_version)
(write or save)(ds, rows)
by_state: dict = {}
for r in rows:
by_state[r["state"] or ""] = by_state.get(r["state"] or "", 0) + 1
moved = [r for r in rows if r["prev_state"] and r["prev_state"] != r["state"]]
held = [r for r in rows if r["raw_state"] != r["state"]]
print(f"逐票逻辑状态 {ds}:写入 {len(rows)}"
f"{''.join(f'{k} {v}' for k, v in sorted(by_state.items())) or ''}"
f",落定态迁移 {len(moved)} 只,被抗抖动按住 {len(held)} 只 -> {config.LOGIC_STATE_TABLE}")
for r in moved[:10]:
print(f" 迁移:{r['code']} {r['prev_state']} -> {r['state']}{r['settle_note']}")
return {"date": ds, "rows": len(rows), "by_state": by_state, "moved": len(moved),
"held": len(held), "table": config.LOGIC_STATE_TABLE}