逻辑状态四态接进计划装配,结论随每张卡发给下游

此前四态只是一个纯函数模块,结论躺在库里没人读。这次接上:

计划装配时每只票算一次四态,三路输入分别是研报论断(已有)、产业研判(读行业观点
快照,按环节名对上今日被指向的环节)、券商行动(同财年同预测期的每股收益预测中位数
与覆盖机构数,两个等长窗口比较)。公司事件那一路无数据源,恒出缺失,但照样记名——
缺失要让人看得见系统缺的是什么,不能让人以为系统判过了。

两处取数搬进公用的地方,免得读数脚本和计划各写一套:行业观点快照的当日读取进
judgement.py,券商行动进 sources.py。读数脚本改调它们,自己那两份删掉。

修了一处日期口径错误:行业观点快照按计划日落库,行情与论断按数据日取,两者在生产里
差一天。原先用同一个日期取三样东西,结果是快照首日读到空。

另修一处:读数脚本原先只读含已启动成员的窄传导视图,而候选卡认的是完整那张,
算出来的产业研判覆盖比卡上真实看到的低。两边现在同一口径。

接口每行新增逻辑状态:状态、子因、每路的来龙去脉与截止日。不发权重、不发判决改动——
四态怎么作用于建仓通道是 PMS 那边的事,这里只提供状态与出处。内部中间量不外泄。

测试十七例,重点钉住四态不改判决:四个状态乘三个判决十二种组合逐个扫过,判决一次
都没被动过。这是收敛规则单调性的守卫。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
zlt 2026-09-04 11:46:01 +08:00
parent f53bcc9973
commit 5425dc6cf7
5 changed files with 287 additions and 84 deletions

View File

@ -142,6 +142,50 @@ def build_rows(day: str, fetched: list, prev: dict, now: str | None = None) -> l
return out
def snapshot_of(day: str, read_mysql=None) -> dict:
"""这个计划日**当天**那一版快照,按簇键索引;表没建或当天没写过返回空字典。
load_previous 的区别那个读的是严格早于计划日的上一版用来算迁移这个读的是
当天这一版用来判断"今天这个环节的行业观点是什么"用错了会在快照首日读到空
空值统一归一成 None数据库驱动把空值读成不是 None 的东西时例如 pandas NaN
直接往下传会让"没有上一版倾向"看着像有值判据那边只认 None
"""
table = config.JUDGEMENT_SNAPSHOT_TABLE
reader = read_mysql or db.read_mysql
try:
df = reader("factor", f"SELECT * FROM {table} WHERE plan_date = %s", (day,))
except Exception as e: # noqa: BLE001
print(f" (行业观点快照表读取失败,产业研判这一路整体缺席: {e!r}")
return {}
rows = df.to_dict("records") if hasattr(df, "to_dict") else list(df or [])
out = {}
for r in rows:
key = str(r.get("cluster_key") or "").strip()
if key:
out[key] = {k: _blank_to_none(v) for k, v in r.items()}
return out
def _blank_to_none(v):
if v is None:
return None
if isinstance(v, float) and v != v: # NaN 只跟自己不相等
return None
return v
def by_segment_name(snaps: dict) -> dict:
"""把快照按环节名索引,供"这只票所在的环节有没有行业观点"这类查询用。
环节名为空的簇例如产业主题级的那些不进这个索引"""
out = {}
for r in (snaps or {}).values():
name = str(r.get("segment_name") or "").strip()
if name:
out[name] = r
return out
def load_previous(day: str, read_mysql=None) -> dict:
"""本计划日之前、回看窗口之内,每个簇最近一次快照的行,按簇键索引。

View File

@ -33,88 +33,13 @@ import config
import db
import judgement
import logic_state as ls
import plan
import sources
# 丙路两个窗口各自的长度(自然日)。等长是硬要求,见模块说明第一条。
BROKER_WINDOW_DAYS = 45
def broker_paths(codes, ds: str) -> dict:
"""按票算丙路信号。返回前缀码到 logic_state.signal 的字典(算不出的票不进字典)。"""
end = dt.date.fromisoformat(ds)
mid = end - dt.timedelta(days=BROKER_WINDOW_DAYS)
start = end - dt.timedelta(days=BROKER_WINDOW_DAYS * 2)
# noqa: SLF001 —— _to_dot 是同仓自用的代码格式转换
dotted = sorted({sources._to_dot(c) for c in codes if c}) # noqa: SLF001
if not dotted:
return {}
out: dict[str, dict] = {}
# 一次拉两个窗口的全部行,按票在内存里分窗——逐票查库要发几千次请求。
marks = ",".join(["%s"] * len(dotted))
try:
df = db.read_mysql(
"factor",
f"SELECT ts_code, report_date, quarter, org_name, eps FROM gp_report_rc "
f"WHERE ts_code IN ({marks}) AND report_date > %s AND report_date <= %s "
f"AND eps IS NOT NULL AND quarter IS NOT NULL",
tuple(dotted) + (start.isoformat(), end.isoformat()))
except Exception as e: # noqa: BLE001
print(f" (券商研报明细表读取失败,丙路整体缺席: {e!r}")
return {}
# 分票、分窗、分财年地堆起来:{票: {财年: {"now": {机构: (日期, 每股收益)}, "prev": ...}}}
box: dict = defaultdict(lambda: defaultdict(lambda: {"now": {}, "prev": {}}))
for r in df.itertuples():
d = sources._ymd(r.report_date) # noqa: SLF001 —— 同仓自用
if not d:
continue
win = "now" if d > mid.isoformat() else "prev"
k = common.to_prefix(str(r.ts_code).strip())
q = str(r.quarter).strip()
org = str(r.org_name or "").strip() or "未署名"
slot = box[k][q][win]
# 同一家机构在窗口里发了多篇,只留最近一篇(模块说明第三条)。
if org not in slot or d > slot[org][0]:
slot[org] = (d, float(r.eps))
for k, by_q in box.items():
# 两个窗口都有料的财年里,取行数最多的那个作可比口径(模块说明第二条)。
usable = [(q, v) for q, v in by_q.items() if v["now"] and v["prev"]]
if not usable:
continue
q, v = max(usable, key=lambda kv: len(kv[1]["now"]) + len(kv[1]["prev"]))
now = {"eps": st.median([x[1] for x in v["now"].values()]), "firms": len(v["now"])}
prev = {"eps": st.median([x[1] for x in v["prev"].values()]), "firms": len(v["prev"])}
s = ls.from_broker(now, prev, as_of=ds)
s["refs"] = [{**(s["refs"][0] if s["refs"] else {}), "quarter": q}]
out[k] = s
return out
def snapshot_of(day: str) -> dict:
"""行业观点快照表里这个计划日的那一版,按簇键索引;表没建或没有当日行返回空字典。"""
try:
df = db.read_mysql(
"factor",
f"SELECT * FROM {config.JUDGEMENT_SNAPSHOT_TABLE} WHERE plan_date = %s", (day,))
except Exception as e: # noqa: BLE001
print(f" (行业观点快照表读取失败,乙路整体缺席: {e!r}")
return {}
if df.empty:
return {}
# pandas 把空值读成 NaN而 NaN 是真值——直接往下传会让"没有上一版倾向"看着像有值。
# 全部归一成 None判据那边只认 None。
import math
def _n(v):
if v is None or (isinstance(v, float) and math.isnan(v)):
return None
return v
return {str(r["cluster_key"]): {k: _n(v) for k, v in r.items()}
for _, r in df.iterrows()}
def main(ds: str | None = None, plan_day: str | None = None) -> None:
"""ds 是数据日行情与论断按它取plan_day 是计划日(行业观点快照按它取)。
@ -136,20 +61,16 @@ def main(ds: str | None = None, plan_day: str | None = None) -> None:
print(f"当日有行情的票 {len(codes)}\n")
claims = sources.logic_claims(codes, ds, per_stock=200)
brokers = broker_paths(codes, ds)
brokers = sources.broker_actions(codes, ds)
# 要的是这个计划日当天那一版快照(带迁移与陈旧两列),不是它之前的那一版——
# load_previous 读的是严格早于计划日的,用它会在快照首日读到空。当天没有就退回上一版。
snaps = snapshot_of(plan_day) or judgement.load_previous(plan_day)
snaps = judgement.snapshot_of(plan_day) or judgement.load_previous(plan_day)
print(f"甲路取到 {len(claims)} 只票的论断;丙路算得出 {len(brokers)} 只票;"
f"乙路快照 {len(snaps)} 个簇\n")
# 乙路按环节名对上主题:候选卡按环节,产业研判按主题聚簇,两者不在一个命名空间,
# 这里只做同名匹配,对不上的票乙路就是缺失。这一路的天花板本来就低(实测 1.4%)。
by_seg = {}
for r in snaps.values():
name = str(r.get("segment_name") or r.get("subject_name") or "").strip()
if name:
by_seg[name] = r
by_seg = judgement.by_segment_name(snaps)
seg_of = {}
# 用完整的传导视图,不是只含已启动成员的那张。候选卡判"所在环节被指向"时认的就是
# 完整这张plan.py 的 evd["pointed"] 是两张视图取或),读数脚本必须跟它一致,

78
plan.py
View File

@ -27,6 +27,8 @@ import pandas as pd
import card
import common
import config
import judgement
import logic_state
import db
import sources
import version
@ -110,6 +112,67 @@ def _val(series: pd.Series, k: str):
return None if v is None or pd.isna(v) else float(v)
def _state_out(st) -> dict | None:
"""逻辑状态发给下游的形状:只留人要读的和判决要用的,不发内部中间量。
每路都带自己的截止日与一句话说明四态的合成规则第一条就是每路必须带截止日
下游要能看出"这一路的证据是哪天的"否则分不清"没有证据""证据很旧"
"""
if not isinstance(st, dict):
return None
return {"state": st.get("state"), "why": st.get("why"), "as_of": st.get("as_of"),
"usable": st.get("usable") or [], "missing": st.get("missing") or [],
"reasons": (st.get("reasons") or [])[:6],
"paths": [{"path": p.get("path"), "signal": p.get("signal"),
"as_of": p.get("as_of"), "why": p.get("why")}
for p in (st.get("paths") or []) if isinstance(p, dict)]}
def _next_day(ds: str) -> str:
"""数据日的次日,也就是计划日。行业观点快照按计划日落库,行情与论断按数据日取,
两者在生产里本来就差一天计划日凌晨构建用的是上一个交易日的数据"""
try:
return (dt.date.fromisoformat(ds) + dt.timedelta(days=1)).isoformat()
except ValueError:
return ds
def _segments_of(ds: str) -> dict:
"""每只票当日所在的全部被指向环节,按前缀码索引。
读的是完整的传导视图不是只含已启动成员的那张候选卡判"所在环节被指向"
认的就是完整这张这里必须跟它一致否则算出来的产业研判覆盖比卡上真实看到的低
读不到返回空字典产业研判这一路整体缺席不让计划断产
"""
try:
d = db.read_pg("SELECT ts_code, target FROM v_factor_transmission "
"WHERE scan_date = %s", (ds,))
except Exception as e: # noqa: BLE001
print(f" (传导视图读取失败,产业研判这一路整体缺席: {e!r}")
return {}
out: dict = {}
for r in d.itertuples():
out.setdefault(common.to_prefix(str(r.ts_code).strip()), []).append(str(r.target))
return out
def _logic_state_of(k: str, evd: dict, seg_of: dict, seg_view: dict,
broker: dict, ds: str) -> dict:
"""一只票的逻辑状态四态。四路各自归一,再按合成规则合成。
乙路的取法一只票可能挂在多个被指向的环节上取第一个有行业观点的那个
取不到就是缺失卡上会写明缺的是哪一路缺失既不算负面也不算正面证据
但必须让人看得见系统缺的是什么不能让人以为系统判过了
"""
a = logic_state.from_claims(evd.get("logic"), ds, stale_days=config.LOGIC_STALE_DAYS)
row = next((seg_view[t] for t in seg_of.get(k, []) if t in seg_view), None)
b = logic_state.from_judgement(row)
c = broker.get(k) or logic_state.signal(
logic_state.PATH_BROKER, logic_state.SIG_NONE,
why="两个等长窗口里算不出可比的每股收益预测")
return logic_state.compose([a, b, c, logic_state.from_events()])
def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series,
mkt_days: set, risk: set | None = None) -> tuple[dict, list]:
"""候选卡装配2026-09-02 方案第 2.2 节第三项):对档位表里的全部票(主榜与观察档,
@ -126,6 +189,14 @@ def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series,
night = sources.night_conclusions(codes, ds)
# 因果论断2026-09-03数据基座抽取的论断挂在卡上作证据线只展示不进判决视图未建时为空。
logic = sources.logic_claims(codes, ds)
# 逻辑状态四态的三路输入丁路公司事件无数据源logic_state 那边恒出缺失)。
# 这三路都是"研究证据还在不在"的跟踪,与候选卡的三门槛判决是正交的两维:
# 判决回答今天要不要买,四态回答支撑它的研究证据还在不在。收敛规则在
# logic_state.apply_to_card是单调的——强化只能提前卡内序、永远不升判决。
plan_day = _next_day(ds)
seg_view = judgement.by_segment_name(judgement.snapshot_of(plan_day))
broker = sources.broker_actions(codes, ds)
seg_of = _segments_of(ds)
if risk is None: # collect 会传入读过一次的名单;单独调用时自己读
try:
risk = factors._risk_set() or set() # noqa: SLF001 —— 同仓自用
@ -146,13 +217,14 @@ def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series,
"accum_state": n.get("accum_state"), "accum_score": n.get("accum_score"),
"accum_age": n.get("accum_age"), "y_signal": n.get("signal"),
"stale_snapshot": stale, "logic": logic.get(k) or []}
state = _logic_state_of(k, evd, seg_of, seg_view, broker, ds)
j = card.judge(evd, start_pct=config.CARD_START_PCT,
accum_max_age=config.CARD_ACCUM_MAX_AGE,
neg_tol=config.UPSIDE_NEG_TOLERANCE,
logic_stale_days=config.LOGIC_STALE_DAYS,
require_started=config.CARD_REQUIRE_STARTED)
cards[k] = {
**j,
**j, "logic_state": state,
"theme": theme, "n_sources": n_sources, "chain_fit": evd["chain_fit"],
"started_source": "moved_view" if mv else None,
"logic_claims": evd["logic"],
@ -275,6 +347,10 @@ def collect(date: str | None = None, top: int = 20, obs_top: int = 10,
r.update(verdict=c["verdict"], reasons=c["reasons"], missing=c["missing"],
risk=c["risk"], card_rank=c["card_rank"],
basis=c.get("basis"), logic=c.get("logic") or [],
# 2026-09-04 新增:逻辑状态四态。只发状态、子因、每路的来龙去脉与
# 截止日,不发权重也不发判决改动——四态怎么作用于建仓通道是 PMS
# 那边的事,这里只提供状态与出处。
logic_state=_state_out(c.get("logic_state")),
card={"pct0": c.get("pct0"), "net_z": c.get("net_z"),
"heat_chg": c.get("heat_chg"), "accum": c.get("accum"),
"night": c.get("night"), "gates": c.get("gates"),

View File

@ -31,6 +31,7 @@ from __future__ import annotations
import datetime as dt
import json
import statistics
from collections import defaultdict
import pandas as pd
@ -563,3 +564,72 @@ def _latest_row(rmy, src: str, table: str) -> tuple[dict | None, str | None]:
order = f"`{date_col}`" if date_col else "1"
rows = _records(rmy(src, f"SELECT * FROM {table} ORDER BY {order} DESC LIMIT 1"))
return (rows[0] if rows else None), date_col
# 丙路(券商行动)两个窗口各自的长度,自然日。等长是硬要求:窗口不等长会让八成的票
# 假显示覆盖收缩——实测前 135 天对近 45 天时有 907 只票误报。等长本身也是抗抖动的低通。
BROKER_WINDOW_DAYS = 45
def broker_actions(codes, ds: str, *, window_days: int = BROKER_WINDOW_DAYS,
read_mysql=None) -> dict:
"""券商用行动说话这一路:同一财年同一预测期的每股收益预测中位数与覆盖机构数,
比较最近两个等长窗口返回前缀码到 logic_state.signal 的字典算不出的票不进字典
三条口径必须照做否则读数是错的
两个窗口等长见上面那条常量的说明
同一财年才可比按预测期字段精确匹配跨财年比较没有意义
同一家机构在窗口里可能发多篇先按机构取最近一篇再算中位数
否则发得勤的机构会被重复计入
看的是券商的行动不是言辞券商极少明说不看好某个行业所以等不到它开口
只能看预测在不在下修覆盖在不在收缩
"""
import logic_state as ls
reader = read_mysql or db.read_mysql
end = dt.date.fromisoformat(ds)
mid = end - dt.timedelta(days=int(window_days))
start = end - dt.timedelta(days=int(window_days) * 2)
dotted = sorted({_to_dot(c) for c in codes if c})
if not dotted:
return {}
marks = ",".join(["%s"] * len(dotted))
try:
df = reader(
"factor",
f"SELECT ts_code, report_date, quarter, org_name, eps FROM gp_report_rc "
f"WHERE ts_code IN ({marks}) AND report_date > %s AND report_date <= %s "
f"AND eps IS NOT NULL AND quarter IS NOT NULL",
tuple(dotted) + (start.isoformat(), end.isoformat()))
except Exception as e: # noqa: BLE001
print(f" (券商研报明细表读取失败,券商行动这一路整体缺席: {e!r}")
return {}
rows = df.itertuples() if hasattr(df, "itertuples") else []
box: dict = defaultdict(lambda: defaultdict(lambda: {"now": {}, "prev": {}}))
for r in rows:
d = _ymd(r.report_date)
if not d:
continue
win = "now" if d > mid.isoformat() else "prev"
k = common.to_prefix(str(r.ts_code).strip())
org = str(r.org_name or "").strip() or "未署名"
slot = box[k][str(r.quarter).strip()][win]
if org not in slot or d > slot[org][0]: # 同机构多篇只留最近一篇
slot[org] = (d, float(r.eps))
out = {}
for k, by_q in box.items():
usable = [(q, v) for q, v in by_q.items() if v["now"] and v["prev"]]
if not usable:
continue
q, v = max(usable, key=lambda kv: len(kv[1]["now"]) + len(kv[1]["prev"]))
now = {"eps": statistics.median([x[1] for x in v["now"].values()]),
"firms": len(v["now"])}
prev = {"eps": statistics.median([x[1] for x in v["prev"].values()]),
"firms": len(v["prev"])}
sig = ls.from_broker(now, prev, as_of=ds)
if sig["refs"]:
sig["refs"][0]["quarter"] = q
out[k] = sig
return out

92
test_plan_logic_state.py Normal file
View File

@ -0,0 +1,92 @@
"""计划装配接入逻辑状态四态的离线单测(不连库)。
钉住三件事
四态随每张卡一起产出并按约定的形状发给下游每路带截止日与一句话说明
四态**不改判决**候选卡的三门槛判决与逻辑四态是正交的两维收敛规则是单调的
逻辑强化只能提前卡内序永远不能把判决往上升一档逻辑存疑只能改分流通道
不能把仅展示变成可执行
取数任一路读不到都不让计划断产那一路记缺失
跑法python3 test_plan_logic_state.py pytest test_plan_logic_state.py
"""
import logic_state as ls
import plan
DS = "2026-09-03"
def t(name, cond):
assert cond, name
print(" ok", name)
def claim(date, direction="利好"):
return {"disclosure_date": date, "direction": direction, "mechanism": "机制",
"doc_title": "研报", "claim_id": "c" + date.replace("-", "")}
def jrow(leaning="偏多", verified=1, migrated=0, stale=3, seg="固态电解质"):
return {"leaning": leaning, "leaning_prev": None, "migrated": migrated,
"stale_days": stale, "verified": verified, "review_date": "2026-09-03",
"n_materials": 43, "subject_name": seg, "segment_name": seg,
"n_bull": 3, "n_bear": 3}
def main():
print("四态随卡产出")
evd = {"logic": [claim("2026-08-25")]}
seg_of = {"SH600000": ["固态电解质"]}
seg_view = {"固态电解质": jrow()}
broker = {"SH600000": ls.signal(ls.PATH_BROKER, ls.SIG_FLAT, as_of=DS,
coverage=6, why="预测中位数变化 +1%,在阈值之内")}
st = plan._logic_state_of("SH600000", evd, seg_of, seg_view, broker, DS) # noqa: SLF001
t("三路都在且无负面 -> 逻辑成立", st["state"] == ls.STATE_HOLD)
t("四路都记了名,缺的那一路写明是公司事件",
len(st["paths"]) == 4 and ls.PATH_EVENT in st["missing"])
print("取数缺席时不断产")
st = plan._logic_state_of("SH600000", {}, {}, {}, {}, DS) # noqa: SLF001
t("四路全缺 -> 无法判断加证据不足,不抛异常",
st["state"] == ls.STATE_UNKNOWN and st["why"] == ls.WHY_THIN)
t("缺失名单写全了四路", len(st["missing"]) == 4)
st = plan._logic_state_of("SH999999", evd, seg_of, seg_view, broker, DS) # noqa: SLF001
t("这只票不在任何被指向环节上 -> 产业研判缺席,其余照算",
ls.PATH_JUDGE in st["missing"] and ls.PATH_CLAIM in st["usable"])
print("一票挂多个环节")
st = plan._logic_state_of( # noqa: SLF001
"SH600000", evd, {"SH600000": ["没评过的环节", "固态电解质"]}, seg_view, broker, DS)
t("取第一个有行业观点的那个环节", ls.PATH_JUDGE in st["usable"])
print("发给下游的形状")
out = plan._state_out(st) # noqa: SLF001
t("状态、子因、截止日都在", set(out) >= {"state", "why", "as_of", "usable", "missing",
"reasons", "paths"})
t("每一路都带自己的截止日与一句话说明",
all(set(p) == {"path", "signal", "as_of", "why"} for p in out["paths"]))
t("不发内部中间量(硬触发标记、出处原值不外泄)",
all("hard" not in p and "refs" not in p for p in out["paths"]))
t("没有状态时发 None不硬拼一个空壳", plan._state_out(None) is None) # noqa: SLF001
print("四态不改判决(收敛规则单调)")
for state in (ls.STATE_STRONG, ls.STATE_HOLD, ls.STATE_UNKNOWN, ls.STATE_DOUBT):
for verdict in ("候选", "关注", "仅展示"):
r = ls.apply_to_card(verdict, state)
assert r["verdict"] == verdict, (verdict, state, r)
t("十二种组合逐个扫过,判决一次都没被四态改动", True)
t("逻辑强化最多提前卡内序", ls.apply_to_card("关注", ls.STATE_STRONG)["rank_bonus"] == 1)
t("逻辑存疑不把仅展示变成可执行",
not ls.apply_to_card("仅展示", ls.STATE_DOUBT)["force_confirm"])
t("候选加逻辑存疑是唯一需要新语义的一格:强制人工确认",
ls.apply_to_card("候选", ls.STATE_DOUBT)["force_confirm"])
print("计划日与数据日差一天")
t("次日推算正确", plan._next_day("2026-09-03") == "2026-09-04") # noqa: SLF001
t("认不出的日期原样返回,不抛异常", plan._next_day("不是日期") == "不是日期") # noqa: SLF001
print("ALL OK — 四态随卡产出 / 缺席不断产 / 下游形状 / 判决不被改动 / 计划日推算 全部通过")
if __name__ == "__main__":
main()