akg-factor-bridge/judgement.py

308 lines
15 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 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",
f"SELECT {', '.join(_PREV_COLS)} 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] = dict(r)
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) -> dict:
"""一个计划日的完整一步:取最新一版结论、配上一版比对、幂等写当日行、打印读数。
day 不传取今天。补写一个早于今天的计划日要当心:视图给的永远是此刻的最新一版,
把它记在过去某一天名下,等于把今天的观点写成那天的观察结果,版本史会失真——
所以这种情况下会多打印一行提醒,但仍照写(重跑与补跑都幂等)。
三个函数都可注入(离线单测):取数、读上一版、落库。"""
day = day or dt.date.today().isoformat()
if day != dt.date.today().isoformat():
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