264 lines
13 KiB
Python
264 lines
13 KiB
Python
"""日频行业观点快照:每个计划日把数据基座最新一版的产业研判与环节评析原样抄一行,只写不判。
|
||
|
||
## 为什么要有它
|
||
|
||
产业研判的结论存在数据基座的评析表里,那张表按主题唯一:每次重新评估就整行覆盖,并把生成
|
||
时间改成当前。所以库里永远只有最新一版,查不到"这个观点上一次是什么样、什么时候变的"。
|
||
而逻辑状态四态里"采信倾向从偏多变成偏空"这条判据依赖的正是这种变化,变化必须由读的一方自己
|
||
建立版本史。不先攒,这条判据永远看不到方向变过。
|
||
|
||
本模块就是攒版本史的那一步(2026-09-03 方案第四之五之三节,落地顺序第一步):每个计划日为
|
||
每个主题或环节写一行,只抄不判,不产生任何判决、不改任何排序、不拦任何票。
|
||
|
||
## 落点与表名
|
||
|
||
写在平台因子库(写 t_factor_akg_* 的同一个 MySQL,config.factor_mysql)。理由两条:桥对数据
|
||
基座 PostgreSQL 只有只读账号(见 config 模块说明),写不进去;这张表又是选股系统自己的派生
|
||
记录,不是数据基座的事实,本来也不该回写基座。表名默认 t_akg_judgement_snapshot,不带
|
||
t_factor_ 前缀,免得平台的因子清单把它当成一张因子表收进去。
|
||
|
||
## 幂等
|
||
|
||
同一计划日重跑先删当日行再整批插,不追加重复行。删与插在同一个事务里(一天几十行,量小,
|
||
不必像因子表那样分块提交)。上一版的比对只看"早于本计划日"的行,所以重跑得到的迁移与陈旧
|
||
天数与首次跑完全一样。
|
||
|
||
## 迁移与陈旧的判定口径(写进列里,但本模块不判任何状态)
|
||
|
||
只有材料指纹(input_version)变化才算一次迁移:migrated 记 1,陈旧天数归零。
|
||
指纹没变而日期变老只累加陈旧天数:migrated 记 0,stale_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 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
|