akg-factor-bridge/judgement.py

355 lines
18 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.

"""日频行业观点快照:每个计划日把数据基座最新一版的产业研判与环节评析原样抄一行,只写不判。
## 为什么要有它
产业研判的结论存在数据基座的评析表里,那张表按主题唯一:每次重新评估就整行覆盖,并把生成
时间改成当前。所以库里永远只有最新一版,查不到"这个观点上一次是什么样、什么时候变的"
而逻辑状态四态里"采信倾向从偏多变成偏空"这条判据依赖的正是这种变化,变化必须由读的一方自己
建立版本史。不先攒,这条判据永远看不到方向变过。
本模块就是攒版本史的那一步2026-09-03 方案第四之五之三节,落地顺序第一步):每个计划日为
每个主题或环节写一行,只抄不判,不产生任何判决、不改任何排序、不拦任何票。
## 落点与表名
写在平台因子库(写 t_factor_akg_* 的同一个 MySQLconfig.factor_mysql。理由两条桥对数据
基座 PostgreSQL 只有只读账号(见 config 模块说明),写不进去;这张表又是选股系统自己的派生
记录,不是数据基座的事实,本来也不该回写基座。表名默认 t_akg_judgement_snapshot不带
t_factor_ 前缀,免得平台的因子清单把它当成一张因子表收进去。
## 幂等
同一计划日重跑先删当日行再整批插,不追加重复行。删与插在同一个事务里(一天几十行,量小,
不必像因子表那样分块提交)。上一版的比对只看"早于本计划日"的行,所以重跑得到的迁移与陈旧
天数与首次跑完全一样。
## 迁移与陈旧的判定口径(写进列里,但本模块不判任何状态)
只有材料指纹input_version变化才算一次迁移migrated 记 1陈旧天数归零。
指纹没变而日期变老只累加陈旧天数migrated 记 0stale_days 加上距上一个计划日的自然日数。
本表里第一次见到这个主题无从比对migrated 留空,陈旧天数从零起算。
指纹本身缺失视图给了空值时不判迁移migrated 留空,陈旧天数照常累加——陈旧只依赖时间。
采信倾向的上一版值单独存一列 leaning_prev"从偏多变成偏空"这类问题不必回表自己拼。
往后要算逻辑状态四态,入口是本模块的 history 与 migrations 两个函数:它们只回答"什么时候
变过、从什么变成什么",不回答"该判成哪个状态"——四态合成是后续步骤(落地顺序第四步)的事,
现在不判。
"""
from __future__ import annotations
import datetime as dt
import config
import db
import sources
# 落库列的顺序插入语句按此顺序拼参数。plan_date 是计划日snapshot_at 是写入时刻。
COLUMNS = ("plan_date", "cluster_key", "scope", "subject_name", "segment_name",
"leaning", "leaning_prev", "n_bull", "n_bear", "n_flags",
"verified", "n_verify_problems", "n_materials", "input_version",
"migrated", "stale_days", "review_date", "reviewed_at",
"model", "rounds", "review_id", "snapshot_at")
# 建表语句与服务器上手工执行的那一份逐字相同docs 里给的就是这一段)。
# 主键取(计划日,簇键):一个计划日一个簇最多一行,重跑覆盖;另建(簇键,计划日)索引,
# 按主题回看版本史时走它。
_CREATE_TABLE = """
CREATE TABLE IF NOT EXISTS {t} (
plan_date DATE NOT NULL,
cluster_key VARCHAR(191) NOT NULL,
scope VARCHAR(16),
subject_name VARCHAR(128),
segment_name VARCHAR(128),
leaning VARCHAR(32),
leaning_prev VARCHAR(32),
n_bull INT,
n_bear INT,
n_flags INT,
verified TINYINT,
n_verify_problems INT,
n_materials INT,
input_version VARCHAR(80),
migrated TINYINT,
stale_days INT,
review_date DATE,
reviewed_at DATETIME,
model VARCHAR(64),
rounds INT,
review_id VARCHAR(64),
snapshot_at DATETIME,
PRIMARY KEY (plan_date, cluster_key),
KEY idx_cluster (cluster_key, plan_date)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
"""
# 回看上一版时读的列(不必把全部列拉回来)。
_PREV_COLS = ("plan_date", "cluster_key", "leaning", "input_version", "stale_days")
def _day_gap(earlier: str | None, later: str | None) -> int:
"""两个计划日之间的自然日数;认不出日期时按一天算(陈旧只会少算不会多算)。"""
try:
return max(0, (dt.date.fromisoformat(later) - dt.date.fromisoformat(earlier)).days)
except (TypeError, ValueError):
return 1
def build_rows(day: str, fetched: list, prev: dict, now: str | None = None) -> list:
"""纯函数:把取回来的最新一版结论,配上一版的比对结果,拼成可落库的行。
day 是计划日ISO 日期串fetched 是 sources.judgement_rows 的返回prev 是
簇键到上一版行的字典load_previous 的返回now 是写入时刻(不传取当前时刻)。
迁移与陈旧两列的口径见模块说明。同一簇键在 fetched 里出现多次时只留最后一条。"""
stamp = now or dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
by_key = {}
for r in fetched or []:
key = (r.get("cluster_key") or "").strip()
if not key:
continue
by_key[key] = r
out = []
for key in sorted(by_key):
r = by_key[key]
p = prev.get(key) if prev else None
fp = r.get("input_version")
if not p:
migrated, stale, leaning_prev = None, 0, None
else:
inc = _day_gap(sources._ymd(p.get("plan_date")), day) # noqa: SLF001 —— 同仓自用
base = int(p.get("stale_days") or 0)
prev_fp = (p.get("input_version") or None)
leaning_prev = p.get("leaning") or None
if fp is None or prev_fp is None:
migrated, stale = None, base + inc # 指纹缺失无从比对,但日子照样变老
elif fp == prev_fp:
migrated, stale = 0, base + inc
else:
migrated, stale = 1, 0
out.append({
"plan_date": day, "cluster_key": key[:191], "scope": r.get("scope"),
"subject_name": (r.get("subject_name") or "")[:128] or None,
"segment_name": (r.get("segment_name") or "")[:128] or None,
"leaning": r.get("leaning"), "leaning_prev": leaning_prev,
"n_bull": r.get("n_bull"), "n_bear": r.get("n_bear"), "n_flags": r.get("n_flags"),
"verified": None if r.get("verified") is None else int(bool(r.get("verified"))),
"n_verify_problems": r.get("n_verify_problems"),
"n_materials": r.get("n_materials"), "input_version": fp,
"migrated": migrated, "stale_days": stale,
"review_date": r.get("review_date"), "reviewed_at": r.get("reviewed_at"),
"model": r.get("model"), "rounds": r.get("rounds"),
"review_id": r.get("review_id"), "snapshot_at": stamp,
})
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 recent_rows(day: str, days: int = 5, read_mysql=None) -> dict:
"""每个簇最近 days 个计划日的快照行,按簇键索引、行按计划日升序,只取严格早于 day 的。
给四态乙路判"硬触发要不要维持"2026-09-07迁移只在指纹变化那天记 migrated=1
次日归零,光看当天那一行硬触发只能维持一个计划日。一次查询取回全部簇的近几行,
不逐簇查——计划装配一轮要看几百只票,逐簇查是几百次往返。
读不到返回空字典,那样乙路退回"只看当天一行",不断产。
"""
table = config.JUDGEMENT_SNAPSHOT_TABLE
reader = read_mysql or db.read_mysql
try:
start = (dt.date.fromisoformat(day) - dt.timedelta(days=int(days) * 2 + 7)).isoformat()
except ValueError:
return {}
try:
rows = sources._records(reader( # noqa: SLF001 —— 同仓自用
"factor",
f"SELECT * FROM {table} WHERE plan_date >= %s AND plan_date < %s "
f"ORDER BY cluster_key, plan_date", (start, day)))
except Exception as e: # noqa: BLE001
print(f" {table} 读不到近日行,乙路只看当天一行: {e!r}")
return {}
out: dict = {}
for r in rows:
key = str(r.get("cluster_key") or "").strip()
if key:
out.setdefault(key, []).append({k: _blank_to_none(v) for k, v in dict(r).items()})
return {k: v[-int(days):] for k, v in out.items()}
def load_previous(day: str, read_mysql=None) -> dict:
"""本计划日之前、回看窗口之内,每个簇最近一次快照的行,按簇键索引。
表还没建、或读不到,返回空字典并打印一行原因:那样全部主题都按"第一次见到"处理,
快照照写只是这一天算不出迁移。read_mysql 可注入(离线单测)。"""
table = config.JUDGEMENT_SNAPSHOT_TABLE
reader = read_mysql or db.read_mysql
try:
start = (dt.date.fromisoformat(day)
- dt.timedelta(days=config.JUDGEMENT_PREV_LOOKBACK_DAYS)).isoformat()
except ValueError:
print(f" (计划日 {day!r} 不是合法日期,跳过上一版比对)")
return {}
try:
rows = sources._records(reader( # noqa: SLF001 —— 同仓自用
"factor",
# 取全部列而不是只取算迁移用的那五列:这个函数 2026-09-07 起也给计划装配用,
# 它要按环节名索引、再读采信倾向与自我校验——只取五列时环节名是空的,
# 产业研判这一路就整个缺席(周一计划里可用数为零,就是这个原因)。
f"SELECT * FROM {table} "
f"WHERE plan_date >= %s AND plan_date < %s ORDER BY plan_date", (start, day)))
except Exception as e: # noqa: BLE001 —— 头一次跑时表还不存在,属正常
print(f" {table} 读不到历史行,本次全部按第一次见到处理: {e!r}")
return {}
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 dict(r).items()}
return out
def save(day: str, rows: list, conn_factory=None) -> None:
"""幂等落库:建表(已存在就跳过)、删当日行、整批插当日行,删与插同一事务。
rows 为空也照样删当日行:那表示这一天视图里没有可抄的结论,当日不该留下上一次重跑的残行。
conn_factory 可注入(离线单测),默认走 db.factor_conn。"""
table = config.JUDGEMENT_SNAPSHOT_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 plan_date = %s", (day,))
if payload:
cur.executemany(f"INSERT INTO {table} ({cols}) VALUES ({marks})", payload)
conn.commit()
def snapshot(day: str | None = None, fetch=None, load_prev=None, write=None,
calendar=None) -> dict:
"""一个计划日的完整一步:取最新一版结论、配上一版比对、幂等写当日行、打印读数。
day 不传取今天。补写一个早于今天的计划日要当心:视图给的永远是此刻的最新一版,
把它记在过去某一天名下,等于把今天的观点写成那天的观察结果,版本史会失真——
所以这种情况下会多打印一行提醒,但仍照写(重跑与补跑都幂等)。
三个函数都可注入离线单测取数、读上一版、落库。calendar 注入交易日历(升序 ISO 串)。
非交易日不落行2026-09-07 审查修)。快照按工作日 20:40 投递、plan_date 取当天,
法定假日(工作日)那行 plan_date 大于节前的数据日、小于节后首个数据日;计划读的是
plan_date 不晚于数据日的最新一行,节后计划读到的是与假日行比对得到的 migrated=0
指纹迁移若首次被假日快照记录,任何计划都读不到那个 migrated=1乙路的硬触发就丢了。
国庆那一周连续五个工作日假日,几乎必然命中。所以:今天不是交易日就不写,直接返回。"""
today = dt.date.today().isoformat()
if day is None:
cal = calendar if calendar is not None else sources.trading_days(today, back_days=21)
if cal and today not in cal:
print(f" 行业观点快照:{today} 不是交易日,不落行(最近交易日 {cal[-1]}")
return {"date": today, "rows": 0, "skipped": "非交易日", "last_trading_day": cal[-1],
"table": config.JUDGEMENT_SNAPSHOT_TABLE}
day = today
if day != today:
print(f" ⚠️ 计划日 {day} 不是今天:抄下来的是此刻的最新一版结论,"
f"记在过去的日期名下会让版本史失真,确认是有意为之再用。")
rows_in = (fetch or sources.judgement_rows)()
prev = (load_prev or load_previous)(day)
rows = build_rows(day, rows_in, prev)
(write or save)(day, rows)
moved = [r for r in rows if r["migrated"] == 1]
first = [r for r in rows if r["migrated"] is None]
stale_max = max([r["stale_days"] for r in rows], default=0)
by_scope = {}
for r in rows:
by_scope[r["scope"] or "未标"] = by_scope.get(r["scope"] or "未标", 0) + 1
print(f"行业观点快照 {day}:写入 {len(rows)}"
f"{''.join(f'{k} {v}' for k, v in sorted(by_scope.items())) or ''}"
f",指纹变化 {len(moved)} 个,第一次见到 {len(first)} 个,"
f"最长陈旧 {stale_max} 天 -> {config.JUDGEMENT_SNAPSHOT_TABLE}")
for r in moved:
print(f" 指纹变化:{r['subject_name']}(采信倾向 {r['leaning_prev']} -> {r['leaning']}"
f"生成日 {r['review_date']}")
return {"date": day, "rows": len(rows), "migrated": len(moved), "first_seen": len(first),
"stale_max": stale_max, "by_scope": by_scope,
"table": config.JUDGEMENT_SNAPSHOT_TABLE}
# ============================================================================
# 版本史的读取入口。给后面的逻辑状态合成用,现在只回答"变过没有、什么时候变的"。
# ============================================================================
def history(cluster_key: str, upto: str | None = None, read_mysql=None) -> list:
"""一个簇在本表里的全部版本按计划日升序。upto 给了就只要不晚于它的行。"""
table = config.JUDGEMENT_SNAPSHOT_TABLE
reader = read_mysql or db.read_mysql
sql = (f"SELECT plan_date, leaning, leaning_prev, input_version, migrated, stale_days, "
f"review_date, verified, n_bull, n_bear FROM {table} WHERE cluster_key = %s")
params = (cluster_key,)
if upto:
sql += " AND plan_date <= %s"
params = (cluster_key, upto)
try:
rows = sources._records(reader("factor", sql + " ORDER BY plan_date", params)) # noqa: SLF001
except Exception as e: # noqa: BLE001
print(f" {table} 版本史读取失败: {e!r}")
return []
return [dict(r) for r in rows]
def migrations(rows: list) -> list:
"""纯函数:从一个簇的版本史里摘出被判为迁移的那些计划日。
判定口径就是落库时那一条,这里只是把已经算好的列摘出来,不重新判、也不据此判任何状态。
每条给出:发生在哪个计划日、采信倾向从什么变成什么、当时的材料指纹与生成日期。
采信倾向没变的迁移也照样列出来(材料换了但结论一样,本身就是一条信息)。"""
out = []
for r in rows or []:
if r.get("migrated") != 1:
continue
out.append({"plan_date": sources._ymd(r.get("plan_date")), # noqa: SLF001
"leaning_from": r.get("leaning_prev"), "leaning_to": r.get("leaning"),
"input_version": r.get("input_version"),
"review_date": sources._ymd(r.get("review_date"))}) # noqa: SLF001
return out