diff --git a/README.md b/README.md index b676627..9ae9d92 100644 --- a/README.md +++ b/README.md @@ -46,6 +46,24 @@ PMS 的参考位、择时执行区间、研判上下文由此可用。入池=当 掉榜未恶化留池观察;无持仓、不在计划且形态恶化 → 移入回收站 stock_recycle_bin; 持仓永不出池。规则、时间线、部署与判收见 `docs/选股计划入池_对接说明.md`。 +## 日频行业观点快照(2026-09-03,只写不判) + +`python run.py judgement-snapshot [--date D]` 把数据基座当前这一版产业研判与环节评析 +(视图 `v_factor_judgement`)抄进选股系统自己的表 `t_akg_judgement_snapshot`,每个计划日 +每个主题或环节一行。为什么要抄:那张评析表按主题唯一,重评时整行覆盖并把生成时间改成当前, +库里永远只有最新一版,查不到"这个观点上一次是什么样、什么时候变的";而逻辑状态四态里 +"采信倾向从偏多变成偏空"这条判据依赖的正是这种变化,版本史只能由读的一方自己攒。 + +这一步只写不判:不产生判决、不改排序、不拦任何票。迁移与陈旧的口径写在行上——只有材料指纹 +(`input_version`)变化才算一次迁移(`migrated=1`、`stale_days` 归零),指纹没变而日期变老 +只累加陈旧天数(`migrated=0`);本表第一次见到的主题 `migrated` 留空。落点选平台因子库而不是 +数据基座,因为桥对基座只有只读账号,且这是选股系统的派生记录不是基座事实;表名不带 +`t_factor_` 前缀,免得平台的因子清单把它当因子表收进去。细节见 `judgement.py` 模块说明。 + +表由代码首次运行时自动建(`CREATE TABLE IF NOT EXISTS`)。若代理不允许自动建表, +在平台因子库(`FACTOR_MYSQL_*` 指向的那个库,即写 `t_factor_akg_*` 的同一个)手工执行 +`judgement.py` 里 `_CREATE_TABLE` 的那段语句,两处逐字相同。 + ## 用法(全程 Docker,不在宿主机直跑) ```bash @@ -117,6 +135,9 @@ docker compose exec akg-factor-bridge python run.py freeze --date 2026-07-24 | `WRITE_CHUNK_ROWS` | `50000` | 因子表分块提交,防单事务过大 | | `FROZEN_ROOT` | `/app/data/frozen` | 冻结落点 | | `UPSTREAM_MEMBER_CAP` / `UPSTREAM_QUIET_CAP` | `30` / `12` | 上游截断体检阈值,命中只告警不拦截 | +| `JUDGEMENT_SNAPSHOT_TABLE` | `t_akg_judgement_snapshot` | 日频行业观点快照落哪张表(平台因子库) | +| `JUDGEMENT_SCOPES` | `industry,segment` | 抄哪几类评析簇:产业研判与环节评析。个股与概念评析默认不抄 | +| `JUDGEMENT_PREV_LOOKBACK_DAYS` | `60` | 比对上一版时往前看几个自然日,超窗按第一次见到处理 | ## 网络(Docker) diff --git a/config.py b/config.py index 0dc2fc3..53a78cc 100644 --- a/config.py +++ b/config.py @@ -228,3 +228,17 @@ MARKET_MYSQL_SOURCE = os.environ.get("MARKET_MYSQL_SOURCE", "price").strip().low # 每只票从数据基座的因果论断视图 v_factor_logic 取最近披露日的最多几条论断挂在卡上。 # 只展示、不作门槛、不进判决;设 0 表示不读该视图(视图未建时也可用它关掉那一行告警)。 LOGIC_CLAIMS_PER_STOCK = int(os.environ.get("LOGIC_CLAIMS_PER_STOCK", "3")) + +# --- 日频行业观点快照(2026-09-03 方案第四之五之三节,落地顺序第一步;模块见 judgement.py)------ +# 数据基座的研判结论按簇键唯一,重评时整行覆盖,库里永远只有最新一版,查不到"上一次是什么样、 +# 什么时候变的"。选股系统每个计划日把最新一版抄一行存下来,自己攒版本史,只写不判。 +# 落点是平台因子库(写因子表的同一个 MySQL):桥对数据基座只有只读账号,这张表又是选股系统的 +# 派生记录不是基座事实。表名不带 t_factor_ 前缀,免得平台的因子清单把它当成一张因子表。 +JUDGEMENT_SNAPSHOT_TABLE = os.environ.get("JUDGEMENT_SNAPSHOT_TABLE", "t_akg_judgement_snapshot") +# 抄哪几类评析簇。默认产业研判与环节评析两类:产业研判是主题级的慢信号(四态里的乙路), +# 环节评析将来接进来时版本史已经在攒了。个股评析与概念评析暂不抄,要抄把它们加进这个值。 +JUDGEMENT_SCOPES = {s.strip() for s in + os.environ.get("JUDGEMENT_SCOPES", "industry,segment").split(",") if s.strip()} +# 比对上一版时往前看几个自然日。取到窗口内最近一个计划日的行作为上一版;超过这个窗口没写过 +# 快照的主题,本次按"第一次见到"处理,迁移记为空、陈旧天数从零重新起算。 +JUDGEMENT_PREV_LOOKBACK_DAYS = int(os.environ.get("JUDGEMENT_PREV_LOOKBACK_DAYS", "60")) diff --git a/judgement.py b/judgement.py new file mode 100644 index 0000000..b2c7c02 --- /dev/null +++ b/judgement.py @@ -0,0 +1,263 @@ +"""日频行业观点快照:每个计划日把数据基座最新一版的产业研判与环节评析原样抄一行,只写不判。 + +## 为什么要有它 + +产业研判的结论存在数据基座的评析表里,那张表按主题唯一:每次重新评估就整行覆盖,并把生成 +时间改成当前。所以库里永远只有最新一版,查不到"这个观点上一次是什么样、什么时候变的"。 +而逻辑状态四态里"采信倾向从偏多变成偏空"这条判据依赖的正是这种变化,变化必须由读的一方自己 +建立版本史。不先攒,这条判据永远看不到方向变过。 + +本模块就是攒版本史的那一步(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 diff --git a/judgement_coverage.py b/judgement_coverage.py new file mode 100644 index 0000000..dc44f6a --- /dev/null +++ b/judgement_coverage.py @@ -0,0 +1,645 @@ +"""产业研判与因果论断的覆盖率读数(只读,不写任何库、不落任何文件)。 + +## 为什么要有它 + +《主观量化系统方案_2026-09-03》第四之五之三节把逻辑状态四态定了稿,其中乙路是产业研判。 +接之前必须先量出能对上多少:产业研判只有八个主题,投影表里有一千九百多个环节,而候选卡是按 +环节判的;两边的命名也不在一个体系里(产业研判按股票池主题聚簇,候选卡按环节名)。 +不先量就接,结果是大面积为空。本脚本出四组读数,末尾给一句结论:产业研判这一路今天 +能覆盖多少只票。 + +## 四组读数 + + 甲 产业研判覆盖 —— 主题簇、环节簇、概念簇、个股评析各有多少个,各自最新一条的生成日期 + 与距今天数,方向取值分布,自我校验未通过的条数;主题簇只有几个,逐条列出。 + 乙 命名对齐 —— 产业研判的主题名里有多少个与投影表里的环节名逐字相同,对不上的逐个列出 + 并给出包含关系的候选环节名(这决定要不要做主题与环节的对照表)。 + 环节簇评析的名字也照同样口径对一遍。 + 丙 候选卡侧可达 —— 当日被传导指向的环节里有多少个能对上任意一条产业研判;当日档位表里的票, + 按当日指向的环节算能对上多少只,按投影表全部所属环节算又能对上多少只。 + 丁 因果论断覆盖 —— 当日档位表的票里有多少只有因果论断,论断的披露日最新、中位、最老是哪天, + 距数据日超过三十、六十、九十、一百八十、三百六十五天的各占多少。 + +## 纪律 + +全部是 SELECT,没有任何写操作,也不生成文件,只往终端打印。 +每一路读不到都打印一行原因然后继续,任何一组失败都不影响其余组。 +本脚本不接进常规链路,也不注册到 run.py 的子命令里:它是一次性的决策读数,跑法见下。 + +## 跑法【桥机 192.168.16.155 · ~/akg-factor-bridge】 + + docker compose exec -T akg-factor-bridge python judgement_coverage.py + docker compose exec -T akg-factor-bridge python judgement_coverage.py --date 2026-09-02 + +不传日期时,数据日取平台因子表 t_factor_akg_score 的最新一天(与出计划的口径同源); +取不到时退回基座传导视图的最新扫描日。 + +## 已知口径与局限(读数时要一起看) + + 一,投影表每天 06:10 删后插,只存当前态,所以环节名与成员集合都是"今天的";跑本脚本 + 避开 06:10 那个窗口。 + 二,桥的只读视图 v_factor_segment_members 暴露的是投影后的环节名 segment_name,而传导扫描 + 与因果论断视图连接用的是图上原名 COALESCE(orig_segment_name, segment_name)。改派重排过的 + 环节这两个名字会不同。脚本会额外试读一次投影底表拿原名,读不到就只用视图名并打印一行说明。 + 三,因果论断视图里客体是环节的论断会按成员展开成多行,所以论断条数一律按 claim_id 去重计。 + 四,研判结论表以簇键唯一,重评时整行覆盖,库里只有每个簇的最新一版,没有历史版本, + 所以"距今天数"读的是最新一版的生成日期,不是这个簇第一次生成的日期。 +""" +from __future__ import annotations # 注解不在定义时求值:开发机的 Python 3.9 也能导入本模块 + +import argparse +import datetime as dt +import unicodedata + +import common +import db + +# 论断披露日距数据日的分桶边界(天)。这里只是把分布摊开给人看,不是判定阈值, +# 四态合成真正要用的陈旧天数在方案第四之五之三节里一次定死,不由本脚本决定。 +_AGE_BUCKETS = (30, 60, 90, 180, 365) + + +# ---------------------------------------------------------------- 打印与表格 +def say(msg: str = "") -> None: + print(msg) + + +def sec(title: str) -> None: + print(f"\n{'=' * 76}\n{title}\n{'=' * 76}") + + +def _w(s) -> int: + """终端显示宽度:中日韩宽字符按两列算,其余按一列算。""" + return sum(2 if unicodedata.east_asian_width(ch) in ("W", "F") else 1 + for ch in str(s)) + + +def _pad(s, width: int, right: bool = False) -> str: + gap = max(0, width - _w(s)) + return (" " * gap + str(s)) if right else (str(s) + " " * gap) + + +def table(headers, rows, right_cols=()) -> None: + """打印一张定宽文字表。rows 是二维序列,None 显示为空。""" + body = [["" if c is None else str(c) for c in r] for r in rows] + if not body: + say("(无行)") + return + n = len(headers) + widths = [max([_w(headers[i])] + [_w(r[i]) for r in body]) for i in range(n)] + say(" ".join(_pad(headers[i], widths[i], i in right_cols) for i in range(n))) + say(" ".join("─" * widths[i] for i in range(n))) + for r in body: + say(" ".join(_pad(r[i], widths[i], i in right_cols) for i in range(n))) + + +def _pct(part: int, whole: int) -> str: + return "—" if not whole else f"{part / whole:.1%}" + + +# ---------------------------------------------------------------- 只读取数 +def _rows(df) -> list: + """把 db.read_* 的返回统一成字典列表;空表返回空列表。""" + if df is None: + return [] + if getattr(df, "empty", False): + return [] + if hasattr(df, "to_dict"): + return list(df.to_dict("records")) + return [dict(r) for r in df] + + +def read_pg(sql: str, params=None, what: str = "") -> list: + """只读一条 PG 查询。读不到打印一行原因并返回空列表,不抛出。""" + try: + return _rows(db.read_pg(sql, params)) + except Exception as e: # noqa: BLE001 —— 单路失败不拖累其余读数 + say(f" (读不到{what},本项留空: {e!r})") + return [] + + +def read_factor(sql: str, params=None, what: str = "") -> list: + """只读一条平台因子库查询。读不到打印一行原因并返回空列表,不抛出。""" + try: + return _rows(db.read_mysql("factor", sql, params)) + except Exception as e: # noqa: BLE001 + say(f" (读不到{what},本项留空: {e!r})") + return [] + + +def _ymd(v): + """各表的日期列形态不一,统一成 ISO 日期串;认不出返回 None。""" + if v is None or (isinstance(v, float) and v != v): + return None + if isinstance(v, dt.datetime): + return v.date().isoformat() + if isinstance(v, dt.date): + return v.isoformat() + s = str(v).strip() + if len(s) == 8 and s.isdigit(): + return f"{s[:4]}-{s[4:6]}-{s[6:]}" + if len(s) >= 10 and s[4] == "-" and s[7] == "-" and s[:4].isdigit(): + return s[:10] + return None + + +def _days_between(d_from, d_to): + """两个 ISO 日期串相差多少自然日;任一为空返回 None。""" + if not d_from or not d_to: + return None + try: + return (dt.date.fromisoformat(d_to) - dt.date.fromisoformat(d_from)).days + except ValueError: + return None + + +def _s(v): + if v is None or (isinstance(v, float) and v != v): + return None + s = str(v).strip() + return s or None + + +def _median(vals: list): + xs = sorted(v for v in vals if v is not None) + return xs[len(xs) // 2] if xs else None + + +def _tri(v): + """三态布尔:真、假、说不上(空值)。不能直接拿列值判真假——pandas 读回来的布尔列是 + numpy 的布尔,与 Python 的 True / False 不是同一个对象;空值是 NaN,而 NaN 本身为真。""" + if v is None or (isinstance(v, float) and v != v): + return None + return bool(v) + + +def pick_date(arg): + """数据日:优先命令行,其次平台因子表最新一天,再次基座传导视图最新扫描日。""" + if arg: + return arg.strip() + rows = read_factor("SELECT MAX(trade_date) AS d FROM t_factor_akg_score", + what="平台因子表 t_factor_akg_score 的最新日期") + d = _ymd(rows[0]["d"]) if rows else None + if d: + return d + say(" (退回基座传导视图取最新扫描日)") + rows = read_pg("SELECT MAX(scan_date) AS d FROM v_factor_transmission", + what="传导视图 v_factor_transmission 的最新扫描日") + return _ymd(rows[0]["d"]) if rows else None + + +# ---------------------------------------------------------------- 数据装载 +def load_judgements() -> list: + """研判结论视图全表(列子集)。scope 分产业研判、个股评析、环节评析、概念评析。 + + 这里不走取数层的 sources.judgement_rows:那个函数是给日频快照用的,只取配置里那两类 + scope,并且会跳过簇键或主题名为空的行。覆盖率读数要的恰恰是"库里到底有什么", + 少数一行都会让结论偏乐观,所以本脚本自己读全量、不做任何过滤。 + """ + return read_pg( + "SELECT scope, subject_name, segment_name, ts_code, leaning, verified, " + "review_date, n_materials, n_bull, n_bear, input_version, cluster_key " + "FROM v_factor_judgement", + what="研判结论视图 v_factor_judgement(视图未建时属正常)") + + +def load_segment_names() -> tuple: + """投影表里的环节名与各自的上市成员数。 + + 返回 (视图名到成员数的字典, 图上原名集合)。原名从投影底表试读,只读账号没有底表权限时 + 返回空集合并打印一行说明——此时命名对齐只按视图名判,改派重排过的环节可能被少算。 + """ + rows = read_pg( + "SELECT segment_name, count(DISTINCT ts_code) AS n_members " + "FROM v_factor_segment_members " + "WHERE ts_code IS NOT NULL AND ts_code <> '' GROUP BY segment_name", + what="环节投影视图 v_factor_segment_members") + by_name = {} + for r in rows: + name = _s(r.get("segment_name")) + if name: + by_name[name] = int(r.get("n_members") or 0) + orig = set() + try: + raw = _rows(db.read_pg( + "SELECT DISTINCT COALESCE(orig_segment_name, segment_name) AS seg " + "FROM segment_members_projection")) + orig = {_s(r.get("seg")) for r in raw if _s(r.get("seg"))} + except Exception as e: # noqa: BLE001 —— 底表没权限是常态,视图名已经够用 + say(f" (读不到投影底表的图上原名,命名对齐只按视图里的环节名判: {e!r})") + return by_name, orig + + +def load_pointed(ds: str) -> tuple: + """当日被传导指向的环节,以及每只票挂在哪些被指向环节上。 + + 两张视图合起来才是完整的被指向成员:v_factor_transmission 摊平的是未动名单, + v_factor_transmission_moved 出的是已启动成员。返回 (环节到源数的字典, 票到环节集合的字典)。 + """ + seg_sources, by_code = {}, {} + for sql, what in ( + ("SELECT ts_code, target, n_sources FROM v_factor_transmission " + "WHERE scan_date = %s", "传导视图 v_factor_transmission 当日行"), + ("SELECT ts_code, target, n_sources FROM v_factor_transmission_moved " + "WHERE scan_date = %s", "已动成员视图 v_factor_transmission_moved 当日行")): + for r in read_pg(sql, (ds,), what=what): + seg = _s(r.get("target")) + if not seg: + continue + try: + n = int(r.get("n_sources") or 0) + except (TypeError, ValueError): + n = 0 + seg_sources[seg] = max(seg_sources.get(seg, 0), n) + code = _s(r.get("ts_code")) + if code: + by_code.setdefault(common.to_prefix(code), set()).add(seg) + return seg_sources, by_code + + +def load_panel(ds: str) -> tuple: + """当日档位表:主榜与观察档的票。返回 (主榜集合, 观察档集合, 档位表覆盖只数)。 + + 分界与 plan.py 一致:主榜分不低于 150,其余是观察档。档位表覆盖只数单独从 + t_factor_akg_gate 数,与计划里的 gate_covered 同源。 + """ + main, obs = set(), set() + for r in read_factor( + "SELECT stock_code, factor_value FROM t_factor_akg_score WHERE trade_date = %s", + (ds,), what=f"平台因子表 t_factor_akg_score 的 {ds} 截面"): + code = _s(r.get("stock_code")) + if not code: + continue + try: + v = float(r.get("factor_value")) + except (TypeError, ValueError): + continue + (main if v >= 150.0 else obs).add(common.to_prefix(code)) + gate_rows = read_factor( + "SELECT count(*) AS n FROM t_factor_akg_gate WHERE trade_date = %s", + (ds,), what=f"平台因子表 t_factor_akg_gate 的 {ds} 截面") + gate_n = int(gate_rows[0]["n"]) if gate_rows else 0 + return main, obs, gate_n + + +def load_member_segments(codes: set) -> dict: + """投影表里每只票所属的全部环节(不限于当日被指向)。只留给定的票。""" + out = {} + for r in read_pg( + "SELECT segment_name, ts_code FROM v_factor_segment_members " + "WHERE ts_code IS NOT NULL AND ts_code <> ''", + what="环节投影视图的成员关系"): + code = _s(r.get("ts_code")) + seg = _s(r.get("segment_name")) + if not code or not seg: + continue + k = common.to_prefix(code) + if k in codes: + out.setdefault(k, set()).add(seg) + return out + + +def load_logic(ds: str) -> list: + """因果论断视图按票聚合:论断条数、最新与最老披露日、经环节展开的条数、两类标记。 + + 论断条数一律按 claim_id 去重:客体是环节的论断在视图里会按成员展开成多行。 + """ + return read_pg( + "SELECT ts_code, " + "count(DISTINCT claim_id) AS n_claims, " + "max(disclosure_date) AS d_max, min(disclosure_date) AS d_min, " + "count(DISTINCT claim_id) FILTER (WHERE link_method = 'segment') AS n_via_segment, " + "count(DISTINCT claim_id) FILTER (WHERE disputed) AS n_disputed, " + "count(DISTINCT claim_id) FILTER (WHERE review_flag IS NOT NULL) AS n_flagged " + "FROM v_factor_logic WHERE disclosure_date <= %s GROUP BY ts_code", + (ds,), what="因果论断视图 v_factor_logic(视图未建时属正常)") + + +# ---------------------------------------------------------------- 甲 +_SCOPE_LABEL = {"industry": "产业研判(主题簇)", "segment": "环节评析(环节簇)", + "concept": "概念评析(概念簇)", "company": "个股评析(公司簇)", + "other": "其他簇键"} + + +def part_jia(judg: list, today: str) -> None: + sec("甲 · 产业研判覆盖:各类簇各有多少、最新一条多久没动") + if not judg: + say("研判结论视图没有出行:视图未建、无权限,或表里还没有跑完第二轮的簇。本组读数留空。") + return + groups = {} + for r in judg: + groups.setdefault(_s(r.get("scope")) or "other", []).append(r) + rows = [] + for scope in ("industry", "segment", "concept", "company", "other"): + g = groups.get(scope) + if not g: + rows.append([_SCOPE_LABEL[scope], 0, "—", "—", "—", "—"]) + continue + dates = [d for d in (_ymd(r.get("review_date")) for r in g) if d] + newest, oldest = (max(dates), min(dates)) if dates else (None, None) + gap = _days_between(newest, today) + n_unverified = sum(1 for r in g if _tri(r.get("verified")) is False) + rows.append([_SCOPE_LABEL[scope], len(g), newest or "—", oldest or "—", + "—" if gap is None else f"{gap} 天", n_unverified]) + table(["簇的类别", "簇数", "最新生成日", "最老生成日", "最新一条距今", "自我校验未通过"], + rows, right_cols=(1, 5)) + + lean = {} + for r in judg: + scope = _s(r.get("scope")) or "other" + lean.setdefault(scope, {}) + key = _s(r.get("leaning")) or "(方向为空)" + lean[scope][key] = lean[scope].get(key, 0) + 1 + say("\n方向取值分布(四态合成的乙路判据用的就是这个受控枚举):") + lr = [] + for scope in ("industry", "segment", "concept", "company", "other"): + if scope in lean: + for k, v in sorted(lean[scope].items(), key=lambda kv: -kv[1]): + lr.append([_SCOPE_LABEL[scope], k, v]) + table(["簇的类别", "方向", "簇数"], lr, right_cols=(2,)) + + ind = groups.get("industry") or [] + if ind: + say(f"\n产业研判逐条清单(共 {len(ind)} 个主题,乙路今天全部的材料就是这几行):") + ir = [] + for r in sorted(ind, key=lambda x: _ymd(x.get("review_date")) or "", reverse=True): + d = _ymd(r.get("review_date")) + gap = _days_between(d, today) + ir.append([_s(r.get("subject_name")) or "(无名)", + _s(r.get("leaning")) or "—", + r.get("n_bull"), r.get("n_bear"), r.get("n_materials"), + {True: "是", False: "否", None: "未知"}[_tri(r.get("verified"))], + d or "—", "—" if gap is None else f"{gap} 天", + (_s(r.get("input_version")) or "—")[:12]]) + table(["主题名", "方向", "多头条数", "空头条数", "材料条数", "自我校验通过", + "生成日", "距今", "材料指纹"], ir, right_cols=(2, 3, 4, 7)) + + +# ---------------------------------------------------------------- 乙 +def _substr_hints(name: str, seg_names, limit: int = 3) -> str: + """给对不上的名字找包含关系的候选环节名(一方是另一方的子串),最多给几个。""" + hits = [s for s in seg_names if s != name and (name in s or s in name)] + hits.sort(key=len) + if not hits: + return "(无包含关系的候选)" + tail = " …" if len(hits) > limit else "" + return "、".join(hits[:limit]) + tail + + +def part_yi(judg: list, seg_members: dict, seg_orig: set) -> set: + """返回能被逐字对上的环节名集合(主题名命中的加上环节簇名命中的)。""" + sec("乙 · 命名对齐:产业研判的主题名与投影表里的环节名对不对得上") + all_seg = set(seg_members) | set(seg_orig) + if not all_seg: + say("投影表里一个环节名都没读到,命名对齐无法判定。本组读数留空。") + return set() + say(f"投影表里的环节名:视图名 {len(seg_members)} 个" + + (f",加上图上原名后去重 {len(all_seg)} 个。" if seg_orig else "(图上原名未读到)。")) + if not judg: + say("研判结论视图没有出行,没有主题名可对。本组读数留空。") + return set() + + matched = set() + for scope, label in (("industry", "主题名"), ("segment", "环节簇名")): + rows = [r for r in judg if (_s(r.get("scope")) or "") == scope] + if not rows: + say(f"\n{label}:一条都没有,跳过。") + continue + out, hit = [], 0 + for r in sorted(rows, key=lambda x: _ymd(x.get("review_date")) or "", reverse=True): + name = _s(r.get("segment_name")) if scope == "segment" else None + name = name or _s(r.get("subject_name")) or "(无名)" + ok = name in all_seg + if ok: + hit += 1 + matched.add(name) + out.append([name, "是" if ok else "否", + seg_members.get(name, "—") if ok else "—", + _ymd(r.get("review_date")) or "—", + "" if ok else _substr_hints(name, all_seg)]) + say(f"\n{label} 共 {len(rows)} 个,与环节名逐字相同 {hit} 个" + f"({_pct(hit, len(rows))}),对不上 {len(rows) - hit} 个:") + table([label, "逐字命中环节名", "该环节上市成员数", "生成日", "包含关系的候选环节名"], + out, right_cols=(2,)) + say("\n读法:逐字命中列全是否,就意味着乙路按名字直连一只票也接不上," + "要接必须先做主题与环节的对照表;包含关系那一列是起草对照表时的线索,不是自动匹配的依据。") + return matched + + +# ---------------------------------------------------------------- 丙 +def part_bing(judg: list, seg_sources: dict, code_pointed: dict, + main: set, obs: set, gate_n: int, member_segs: dict) -> tuple: + """返回 (能对上研判的票集合, 按投影环节能对上的票集合)。""" + sec("丙 · 候选卡侧的可达性:当日被指向的环节与档位表的票,有多少能对上产业研判") + theme_names = {_s(r.get("subject_name")) for r in judg + if (_s(r.get("scope")) or "") == "industry"} + seg_cluster_names = {_s(r.get("segment_name")) or _s(r.get("subject_name")) + for r in judg if (_s(r.get("scope")) or "") == "segment"} + theme_names.discard(None) + seg_cluster_names.discard(None) + reach = theme_names | seg_cluster_names + say(f"能被对上的名字一共 {len(reach)} 个:同名主题 {len(theme_names)} 个," + f"环节簇 {len(seg_cluster_names)} 个。") + + if not seg_sources: + say("当日没有读到任何被传导指向的环节,本组前半留空。") + else: + by_seg = sum(1 for s in seg_sources if s in seg_cluster_names) + by_theme = sum(1 for s in seg_sources if s in theme_names) + both = sum(1 for s in seg_sources if s in reach) + say(f"\n当日被传导指向的环节 {len(seg_sources)} 个:按环节簇对上 {by_seg} 个," + f"按同名主题对上 {by_theme} 个,去重合计 {both} 个({_pct(both, len(seg_sources))})。") + top = sorted(seg_sources.items(), key=lambda kv: (-kv[1], kv[0]))[:20] + rows = [[seg, n, + "环节簇" if seg in seg_cluster_names else + ("同名主题" if seg in theme_names else "对不上")] + for seg, n in top] + say("\n被指向最强的二十个环节(源数从高到低):") + table(["被指向的环节", "源数", "对上哪一路研判"], rows, right_cols=(1,)) + + panel = main | obs + say(f"\n当日档位表:主榜 {len(main)} 只,观察档 {len(obs)} 只,合计 {len(panel)} 只;" + f"档位表覆盖 {gate_n} 只。") + if not panel: + say("档位表当日没有票,本组后半留空。") + return set(), set() + + hit_pointed = {k for k in panel if (code_pointed.get(k) or set()) & reach} + hit_member = {k for k in panel if (member_segs.get(k) or set()) & reach} + n_pointed = sum(1 for k in panel if code_pointed.get(k)) + n_member = sum(1 for k in panel if member_segs.get(k)) + rows = [ + ["当日被传导指向(卡上的主题就是它)", n_pointed, _pct(n_pointed, len(panel)), + len(hit_pointed), _pct(len(hit_pointed), len(panel))], + ["投影表里的全部所属环节(不限当日)", n_member, _pct(n_member, len(panel)), + len(hit_member), _pct(len(hit_member), len(panel))], + ] + say("\n档位表的票能不能对上研判,按两种口径各算一遍:") + table(["算所在环节的口径", "有所在环节的票数", "占档位表", "能对上研判的票数", "占档位表"], + rows, right_cols=(1, 2, 3, 4)) + say("两种口径的差别:上一行是候选卡当天真正用的(卡上的主题来自当日传导指向)," + "下一行是把票在投影表里挂着的环节全算上,是乙路理论上的天花板。") + return hit_pointed, hit_member + + +# ---------------------------------------------------------------- 丁 +def part_ding(logic: list, main: set, obs: set, ds: str) -> set: + """返回档位表里有因果论断的票集合。""" + sec("丁 · 因果论断覆盖:档位表的票有多少只有论断、论断有多旧") + panel = main | obs + if not logic: + say("因果论断视图没有出行:视图未建、无权限,或披露日不晚于数据日的论断为零。本组读数留空。") + return set() + by_code = {} + for r in logic: + code = _s(r.get("ts_code")) + if code: + by_code[common.to_prefix(code)] = r + say(f"因果论断视图里一共有 {len(by_code)} 只票带论断(披露日不晚于数据日 {ds})。") + if not panel: + say("档位表当日没有票,只能给出全库读数。") + return set() + + have = {k for k in panel if k in by_code} + have_main = {k for k in main if k in by_code} + rows = [["主榜", len(main), len(have_main), _pct(len(have_main), len(main))], + ["观察档", len(obs), len(have) - len(have_main), + _pct(len(have) - len(have_main), len(obs))], + ["合计", len(panel), len(have), _pct(len(have), len(panel))]] + table(["档位", "票数", "有因果论断", "覆盖率"], rows, right_cols=(1, 2, 3)) + if not have: + say("档位表里一只票都没有论断,披露日分布无从谈起。") + return set() + + n_claims, n_via_seg, n_disputed, n_flagged = 0, 0, 0, 0 + per_stock, newest_dates = [], [] + for k in have: + r = by_code[k] + try: + n_claims += int(r.get("n_claims") or 0) + per_stock.append(int(r.get("n_claims") or 0)) + n_via_seg += int(r.get("n_via_segment") or 0) + n_disputed += int(r.get("n_disputed") or 0) + n_flagged += int(r.get("n_flagged") or 0) + except (TypeError, ValueError): + pass + d = _ymd(r.get("d_max")) + if d: + newest_dates.append(d) + say(f"\n这 {len(have)} 只票身上一共 {n_claims} 条论断(按论断编号去重)," + f"每票中位数 {_median(per_stock)} 条,最多 {max(per_stock) if per_stock else 0} 条;" + f"其中 {n_via_seg} 条是客体为环节、按成员展开挂上来的,不是直接讲这家公司的。") + say(f"标记:处于未结的多空分歧复核 {n_disputed} 条;逻辑评析标了疑似误抽 {n_flagged} 条。" + "这两类在接进四态之前要先决定怎么处置。") + + if not newest_dates: + say("论断的披露日全部认不出来,日期分布留空。") + return have + ages = [a for a in (_days_between(d, ds) for d in newest_dates) if a is not None] + say(f"\n每只票最新一条论断的披露日:最新 {max(newest_dates)}," + f"中位 {_median(newest_dates)},最老 {min(newest_dates)}" + f"(距数据日 {ds} 分别是 {min(ages) if ages else '—'}、{_median(ages)}、" + f"{max(ages) if ages else '—'} 天)。") + if ages: + rows, prev = [], 0 + for b in _AGE_BUCKETS: + n = sum(1 for a in ages if prev < a <= b) + rows.append([f"{prev + 1} 到 {b} 天", n, _pct(n, len(ages)), + sum(1 for a in ages if a > b), + _pct(sum(1 for a in ages if a > b), len(ages))]) + prev = b + n_last = sum(1 for a in ages if a > _AGE_BUCKETS[-1]) + rows.append([f"超过 {_AGE_BUCKETS[-1]} 天", n_last, _pct(n_last, len(ages)), 0, "0.0%"]) + say("\n按距数据日的天数分桶(末两列是超过本桶上界的累计,看陈旧面有多大):") + table(["距数据日", "票数", "占比", "超过本桶上界的票数", "占比"], + rows, right_cols=(1, 2, 3, 4)) + return have + + +# ---------------------------------------------------------------- 结论 +def conclusion(main: set, obs: set, hit_pointed: set, hit_member: set, + have_logic: set, judg: list, today: str) -> None: + sec("结论") + panel = main | obs + if not panel: + say("当日档位表没有票,给不出覆盖结论;先确认那一天的因子表已经构建。") + return + ind = [r for r in judg if (_s(r.get("scope")) or "") == "industry"] + dates = [d for d in (_ymd(r.get("review_date")) for r in ind) if d] + gap = _days_between(max(dates), today) if dates else None + say(f"产业研判这一路今天只有 {len(ind)} 个主题,最新一条生成于 " + f"{max(dates) if dates else '未知'}" + f"{'' if gap is None else f',距今 {gap} 天'}。") + say(f"当日档位表 {len(panel)} 只票里,按候选卡当天用的口径(当日被传导指向的环节)" + f"能对上产业研判的有 {len(hit_pointed)} 只,占 {_pct(len(hit_pointed), len(panel))};" + f"把投影表里挂着的环节全算上也只有 {len(hit_member)} 只," + f"占 {_pct(len(hit_member), len(panel))}——这是乙路的天花板。") + diff = len(have_logic) - len(hit_pointed) + say(f"同一批票里有因果论断的是 {len(have_logic)} 只,占 {_pct(len(have_logic), len(panel))};" + + (f"甲路比乙路多覆盖 {diff} 只。" if diff > 0 else + (f"乙路比甲路多覆盖 {-diff} 只。" if diff < 0 else "两路覆盖的只数一样多。"))) + if len(hit_pointed) == 0: + say("\n判断:乙路现在按名字直连一只票都接不上,接进四态合成等于全票都落到" + "无法判断加证据不足。要么先做主题与环节的对照表再接,要么这一轮先只接甲路," + "乙路等主题数上来、对照表建好再说。") + elif len(hit_pointed) * 10 < len(panel): + say("\n判断:乙路能覆盖的不到档位表的一成,接进去的话绝大多数票仍然是证据不足。" + "可以接,但要在卡上明写缺的是哪一路,不要让读的人误以为是系统判过了。") + else: + say("\n判断:乙路的覆盖面够得上接。仍要在卡上分开标注缺失与矛盾两个子因," + "缺失不当成负面证据。") + say("\n本脚本只读不写,跑完不留任何文件。以上读数随投影表每天 06:10 重建而变," + "拍板前建议连着两个交易日各跑一次看是否稳定。") + + +# ---------------------------------------------------------------- 入口 +def run(date=None) -> int: + today = dt.date.today().isoformat() + say("产业研判与因果论断覆盖率读数(只读,不写任何库、不落任何文件)") + ds = pick_date(date) + if not ds: + say("取不到数据日:平台因子表与基座传导视图都读不到。先确认连接与当日构建,再跑本脚本。") + return 1 + say(f"数据日 {ds};脚本运行日 {today}(容器走国际时间,与北京时间可能差一天," + f"两个日期都打出来是为了让距今天数有据可查)。") + + judg = load_judgements() + seg_members, seg_orig = load_segment_names() + seg_sources, code_pointed = load_pointed(ds) + main, obs, gate_n = load_panel(ds) + member_segs = load_member_segments(main | obs) if (main or obs) else {} + logic = load_logic(ds) + + hit_pointed, hit_member, have_logic = set(), set(), set() + for name, fn in ( + ("甲", lambda: part_jia(judg, today)), + ("乙", lambda: part_yi(judg, seg_members, seg_orig)), + ("丙", lambda: part_bing(judg, seg_sources, code_pointed, main, obs, + gate_n, member_segs)), + ("丁", lambda: part_ding(logic, main, obs, ds))): + try: + r = fn() + if name == "丙" and isinstance(r, tuple): + hit_pointed, hit_member = r + elif name == "丁" and isinstance(r, set): + have_logic = r + except Exception as e: # noqa: BLE001 —— 单组失败不拖累其余读数 + say(f"\n第{name}组读数失败(其余照出): {e!r}") + try: + conclusion(main, obs, hit_pointed, hit_member, have_logic, judg, today) + except Exception as e: # noqa: BLE001 + say(f"\n结论一节失败: {e!r}") + return 0 + + +def main() -> int: + ap = argparse.ArgumentParser( + description="产业研判与因果论断的覆盖率读数(只读)") + ap.add_argument("--date", default=None, + help="数据日 YYYY-MM-DD,不传时取平台因子表最新一天") + a = ap.parse_args() + return run(a.date) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/run.py b/run.py index ac3caf2..5ba0709 100644 --- a/run.py +++ b/run.py @@ -9,6 +9,9 @@ python run.py push-pool [--date D] [--dry-run] # 计划入池:写 Mongo 股票池分组, # 供决策系统每晚推理覆盖(详见 pool.py) python run.py register # 注册全部因子到 factor_metadata + python run.py judgement-snapshot [--date D] # 日频行业观点快照:把数据基座当前这一版 + # 产业研判与环节评析抄一行存起来攒版本史, + # 只写不判(详见 judgement.py) python run.py build all --mode history --start 2024-01-01 --end 2025-12-31 python run.py build akg_heat --mode daily --date 2026-07-24 python run.py build akg_event --mode history --start 2025-01-01 --end 2026-07-24 @@ -182,6 +185,8 @@ def main(): av.add_argument("--dry-run", action="store_true", help="只列语句不执行") ra = sub.add_parser("regime-append") # 08:45 环境追加(写当日计划快照的 regime 段) ra.add_argument("--date", help="默认取 score 表最新日(与 plan 同口径)") + jg = sub.add_parser("judgement-snapshot") # 日频行业观点快照(只写不判,攒版本史) + jg.add_argument("--date", help="计划日,默认今天;补写过去的日期会让版本史失真,会有提示") f = sub.add_parser("freeze") f.add_argument("--date", help="默认今天") b = sub.add_parser("build") @@ -219,6 +224,9 @@ def main(): if not day: raise SystemExit("t_factor_akg_score 还没有数据,没有当日快照可追加。") regime.append_to_snapshot(day) + elif a.cmd == "judgement-snapshot": + import judgement + judgement.snapshot(a.date) elif a.cmd == "tracks": import tracks tracks.coverage_report() diff --git a/sources.py b/sources.py index edb0d6f..3f2b900 100644 --- a/sources.py +++ b/sources.py @@ -15,6 +15,8 @@ 平台 MySQL zs_day_data / eastmoney_rzrq_data / fear_greed_index 计划环境段的市场四项:两市成交额、融资余额、恐贪指数; 市场广度从基座 v_factor_stock_daily 当日行自算(同一节) + 基座 PG v_factor_judgement 产业研判与环节评析的最新一版结论(采信倾向、多空条数、 + 自我校验、材料指纹),每个计划日抄一份存版本史(judgement.py) 代码格式:基座是点后缀式 600000.SH,决策系统与桥是前缀式 SH600000,进出都过 common.to_prefix。 读失败的语义:每一路读不到都返回空字典并打印一行原因,候选卡按"缺失"处理(进关注或仅展示), @@ -291,6 +293,111 @@ def _s(v) -> str | None: return s or None +def _i(v) -> int | None: + x = _f(v) + return None if x is None else int(x) + + +def _b(v) -> bool | None: + """自我校验结果这类三值列:真、假、还没有结果。空值保持为空,不当成假。""" + if v is None or (isinstance(v, float) and v != v): + return None + if isinstance(v, bool): + return v + s = str(v).strip().lower() + if s in ("true", "t", "1", "yes", "y"): + return True + if s in ("false", "f", "0", "no", "n"): + return False + return None + + +def _n_items(v) -> int | None: + """JSON 数组列只留条数:列表直接数,字符串先按 JSON 解析,认不出返回空。""" + if v is None or (isinstance(v, float) and v != v): + return None + if isinstance(v, (list, tuple)): + return len(v) + if isinstance(v, str): + try: + parsed = json.loads(v) + except ValueError: + return None + return len(parsed) if isinstance(parsed, list) else None + return None + + +# ============================================================================ +# 产业研判与环节评析(数据基座 v_factor_judgement) +# 用处见 judgement.py:选股系统每个计划日自留一份快照,攒版本史。 +# ============================================================================ +_JUDGEMENT_COLS = ("scope", "subject_name", "segment_name", "cluster_key", "leaning", + "n_bull", "n_bear", "n_flags", "verified", "verify_problems", + "n_materials", "model", "rounds", "review_date", "reviewed_at", + "input_version", "review_id") + + +def judgement_rows(scopes=None, read_pg=None) -> list[dict]: + """数据基座研判结论视图里 scope 落在指定几类的全部行,一行一个评析簇。 + + scopes 不传时取 config.JUDGEMENT_SCOPES(默认产业研判 industry 与环节评析 segment)。 + 每行带:主题名或环节名、采信倾向、多头论点条数、空头论点条数、旗标条数、自我校验结果与 + 校验问题条数、材料条数、模型与轮次、生成日期与生成时刻、材料指纹、评析编号、簇键。 + + 这张视图按簇键唯一、重评时整行覆盖(视图注释里的已知局限 a),所以它永远只有最新一版; + 版本史由读的一方自己攒,见 judgement.py。读失败返回空列表并打印一行原因,不阻断上层。 + read_pg 可注入(离线单测),默认走 db.read_pg。""" + want = tuple(scopes) if scopes else tuple(sorted(config.JUDGEMENT_SCOPES)) + if not want: + return [] + reader = read_pg or db.read_pg + try: + marks = ",".join(["%s"] * len(want)) + rows = _records(reader( + f"SELECT {', '.join(_JUDGEMENT_COLS)} FROM v_factor_judgement " + f"WHERE scope IN ({marks})", want)) + except Exception as e: # noqa: BLE001 + print(f" (研判结论视图 v_factor_judgement 读取失败,本计划日没有行业观点行可写: {e!r})") + return [] + out = [] + for r in rows: + key = _s(r.get("cluster_key")) + subject = _s(r.get("subject_name")) or _s(r.get("segment_name")) + if not key or not subject: + continue # 簇键或主题名缺一样,这一行没法参与逐版本比对,跳过 + out.append({ + "scope": _s(r.get("scope")), "subject_name": subject, + "segment_name": _s(r.get("segment_name")), "cluster_key": key, + "leaning": _s(r.get("leaning")), + "n_bull": _i(r.get("n_bull")), "n_bear": _i(r.get("n_bear")), + "n_flags": _i(r.get("n_flags")), + "verified": _b(r.get("verified")), + "n_verify_problems": _n_items(r.get("verify_problems")), + "n_materials": _i(r.get("n_materials")), + "model": _s(r.get("model")), "rounds": _i(r.get("rounds")), + "review_date": _ymd(r.get("review_date")), + "reviewed_at": _dt_str(r.get("reviewed_at")), + "input_version": _s(r.get("input_version")), + "review_id": _s(r.get("review_id")), + }) + return out + + +def _dt_str(v) -> str | None: + """时刻列统一成 MySQL 认的 'YYYY-MM-DD HH:MM:SS';只有日期的补零点;认不出返回空。""" + if v is None or (isinstance(v, float) and v != v): + return None + if isinstance(v, dt.datetime): + return v.strftime("%Y-%m-%d %H:%M:%S") + if isinstance(v, dt.date): + return f"{v.isoformat()} 00:00:00" + s = str(v).strip().replace("T", " ") + if len(s) >= 19 and s[4] == "-" and s[7] == "-": + return s[:19] + day = _ymd(s) + return f"{day} 00:00:00" if day else None + + # ============================================================================ # 计划环境段的市场四项(2026-09-03 方案第 1.4 节清单里"有"与"可自算"的项) # ============================================================================ diff --git a/test_judgement_snapshot.py b/test_judgement_snapshot.py new file mode 100644 index 0000000..78632ea --- /dev/null +++ b/test_judgement_snapshot.py @@ -0,0 +1,288 @@ +"""日频行业观点快照的离线单测(不连库)。 + +覆盖:从研判结论视图正常取数并归一(只要产业研判与环节评析两类、自我校验三值、校验问题条数、 +指纹与生成日期);视图读失败返回空列表不抛错;同一计划日重跑覆盖当日行不追加重复行;材料指纹 +未变时陈旧天数累加、采信倾向的上一版值带出来;材料指纹变化时算作一次迁移且陈旧天数归零; +第一次见到的主题不判迁移;版本史里摘迁移的入口函数。 + +取数、读上一版、落库三处都接受注入的函数,落库这一路用一张内存里的假表接住真实的 SQL +(DELETE 当日行 + 批插),所以幂等测的是真语句不是桩。 + +开发机没有 pandas 与数据库驱动时,只给缺席的模块装最小桩(与 test_market_context.py 同一约定: +仅在模块缺席时装桩,不覆盖真实模块)。 + +跑法:python3 test_judgement_snapshot.py 或 pytest test_judgement_snapshot.py +""" +import datetime as dt +import sys +import types + +_STUBS = ("pandas", "psycopg", "pymysql", "dotenv") +for _n in _STUBS: + if _n not in sys.modules: + try: + __import__(_n) + except ImportError: + _m = types.ModuleType(_n) + if _n == "pandas": # db.py 的函数签名在定义时引用这两个名字 + _m.DataFrame = type("DataFrame", (), {}) + _m.Series = type("Series", (), {}) + sys.modules[_n] = _m + +import config # noqa: E402 +import judgement # noqa: E402 +import sources # noqa: E402 + + +def t(name, cond): + assert cond, name + print(" ok", name) + + +DAY = "2026-09-03" +TABLE = "t_akg_judgement_snapshot_test" + + +# ---------------------------------------------------------------- 假数据 +def _view_rows(fp_a="fp-a-1", leaning_a="偏多"): + """研判结论视图的两行产业研判、一行环节评析,外加两行该被挡掉的(个股评析、簇键为空)。""" + return [ + {"scope": "industry", "subject_name": "液冷散热", "segment_name": None, + "cluster_key": "Industry::液冷散热", "leaning": leaning_a, "n_bull": 4, "n_bear": 2, + "n_flags": 1, "verified": True, "verify_problems": [], "n_materials": 37, + "model": "m1", "rounds": 3, "review_date": dt.date(2026, 8, 19), + "reviewed_at": dt.datetime(2026, 8, 19, 10, 30, 5), "input_version": fp_a, + "review_id": "rv1"}, + {"scope": "industry", "subject_name": "固态电池", "segment_name": None, + "cluster_key": "Industry::固态电池", "leaning": "证据不足", "n_bull": "1", "n_bear": 0, + "n_flags": 0, "verified": "false", "verify_problems": '["缺少产能口径", "价格来源不一致"]', + "n_materials": 12, "model": "m1", "rounds": 3, "review_date": "2026-08-11", + "reviewed_at": "2026-08-11T09:00:00", "input_version": "fp-b-1", "review_id": "rv2"}, + {"scope": "segment", "subject_name": "正极材料", "segment_name": "正极材料", + "cluster_key": "Segment:正极材料", "leaning": "中性", "n_bull": 2, "n_bear": 2, + "n_flags": None, "verified": None, "verify_problems": None, "n_materials": 8, + "model": "m1", "rounds": 2, "review_date": "2026-08-20", "reviewed_at": None, + "input_version": None, "review_id": "rv3"}, + {"scope": "industry", "subject_name": None, "segment_name": None, + "cluster_key": "Industry::", "leaning": "偏多", "n_bull": 1, "n_bear": 0, + "n_flags": 0, "verified": True, "verify_problems": None, "n_materials": 3, + "model": "m1", "rounds": 3, "review_date": "2026-08-01", "reviewed_at": None, + "input_version": "fp-x", "review_id": "rv4"}, + ] + + +def _pg_ok(sql, params=None): + s = " ".join(sql.split()) + assert "v_factor_judgement" in s, s + assert set(params) == {"industry", "segment"}, params # 只要这两类 + return _view_rows() + + +def _boom(*a, **k): + raise OSError("connection refused") + + +class _FakeCursor: + """只认本模块会发的三种语句的假游标:建表、按计划日删、批插。""" + + def __init__(self, store): + self.store = store + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + def execute(self, sql, params=None): + s = " ".join(sql.split()) + if s.upper().startswith("CREATE TABLE"): + self.store["created"] += 1 + assert TABLE in s, s + return + if s.upper().startswith("DELETE"): + assert "WHERE plan_date = %s" in s, s + day = params[0] + self.store["deleted"].append(day) + self.store["rows"] = [r for r in self.store["rows"] if r[0] != day] + return + raise AssertionError(f"意外的语句: {s}") + + def executemany(self, sql, rows): + s = " ".join(sql.split()) + assert s.upper().startswith("INSERT INTO") and TABLE in s, s + assert s.count("%s") == len(judgement.COLUMNS), s + for r in rows: + assert len(r) == len(judgement.COLUMNS) + self.store["rows"].extend(list(rows)) + + +class _FakeConn: + def __init__(self, store): + self.store = store + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + def cursor(self): + return _FakeCursor(self.store) + + def commit(self): + self.store["commits"] += 1 + + +def _fake_store(): + return {"rows": [], "deleted": [], "created": 0, "commits": 0} + + +def _as_dicts(store): + return [dict(zip(judgement.COLUMNS, r)) for r in store["rows"]] + + +# ---------------------------------------------------------------- 用例 +def test_fetch(): + config.JUDGEMENT_SCOPES = {"industry", "segment"} + rows = sources.judgement_rows(read_pg=_pg_ok) + t("只出簇键与主题名齐全的行(空主题那行被挡掉)", len(rows) == 3) + a = next(r for r in rows if r["cluster_key"] == "Industry::液冷散热") + t("产业研判字段齐:采信倾向、多空条数、材料指纹、生成日期", + a["scope"] == "industry" and a["subject_name"] == "液冷散热" and a["leaning"] == "偏多" + and a["n_bull"] == 4 and a["n_bear"] == 2 and a["input_version"] == "fp-a-1" + and a["review_date"] == "2026-08-19" and a["reviewed_at"] == "2026-08-19 10:30:05") + t("自我校验为真、校验问题零条", a["verified"] is True and a["n_verify_problems"] == 0) + b = next(r for r in rows if r["cluster_key"] == "Industry::固态电池") + t("字符串形态的数字与真假值都归一、JSON 串的校验问题按条数", + b["n_bull"] == 1 and b["verified"] is False and b["n_verify_problems"] == 2 + and b["reviewed_at"] == "2026-08-11 09:00:00") + c = next(r for r in rows if r["scope"] == "segment") + t("环节评析带环节名;空的自我校验保持为空不当成假", + c["segment_name"] == "正极材料" and c["verified"] is None + and c["n_verify_problems"] is None and c["input_version"] is None) + t("scope 传参可覆盖旋钮", + sources.judgement_rows(scopes=("industry", "segment"), read_pg=_pg_ok) == rows) + + +def test_fetch_failure(): + t("视图读失败返回空列表不抛错", sources.judgement_rows(read_pg=_boom) == []) + config.JUDGEMENT_SCOPES = set() + t("一类都不抄时不查库", sources.judgement_rows(read_pg=_boom) == []) + config.JUDGEMENT_SCOPES = {"industry", "segment"} + + +def test_prev_read_failure(): + config.JUDGEMENT_SNAPSHOT_TABLE = TABLE + t("表还不存在时上一版为空字典、不抛错", judgement.load_previous(DAY, read_mysql=_boom) == {}) + t("计划日不合法时也只是没有上一版", + judgement.load_previous("不是日期", read_mysql=_boom) == {}) + + +def test_same_day_idempotent(): + config.JUDGEMENT_SNAPSHOT_TABLE = TABLE + store = _fake_store() + + def _write(day, rows): + judgement.save(day, rows, conn_factory=lambda: _FakeConn(store)) + + def _fetch(): + return sources.judgement_rows(read_pg=_pg_ok) + + r1 = judgement.snapshot(DAY, fetch=_fetch, load_prev=lambda d: {}, write=_write) + t("首日写三行、每个主题一行", r1["rows"] == 3 and len(store["rows"]) == 3) + t("按 scope 计数:产业研判两个、环节评析一个", + r1["by_scope"] == {"industry": 2, "segment": 1}) + judgement.snapshot(DAY, fetch=_fetch, load_prev=lambda d: {}, write=_write) + t("同一计划日重跑仍是三行,不追加重复行", len(store["rows"]) == 3) + t("重跑先删当日行(删的正是这个计划日)", store["deleted"] == [DAY, DAY]) + keys = [r["cluster_key"] for r in _as_dicts(store)] + t("三个簇键各一行、无重复", sorted(keys) == sorted(set(keys)) and len(keys) == 3) + t("每行都带计划日与写入时刻", + all(r["plan_date"] == DAY and len(r["snapshot_at"] or "") == 19 for r in _as_dicts(store))) + t("首次见到不判迁移、陈旧天数从零起算", + all(r["migrated"] is None and r["stale_days"] == 0 for r in _as_dicts(store)) + and r1["first_seen"] == 3 and r1["migrated"] == 0) + + # 视图这天一行都读不到:当日行照样先删,不留上一次重跑的残行。 + judgement.snapshot(DAY, fetch=lambda: [], load_prev=lambda d: {}, write=_write) + t("视图无行时当日行被清空,不留残行", store["rows"] == []) + + +def test_stale_days_accumulate(): + """指纹未变:陈旧天数按距上一个计划日的自然日数累加,采信倾向的上一版值带出来。""" + prev = {"Industry::液冷散热": {"plan_date": "2026-09-01", "cluster_key": "Industry::液冷散热", + "leaning": "偏多", "input_version": "fp-a-1", "stale_days": 5}} + rows = judgement.build_rows(DAY, sources.judgement_rows(read_pg=_pg_ok), prev, now="x") + a = next(r for r in rows if r["cluster_key"] == "Industry::液冷散热") + t("指纹未变:不算迁移", a["migrated"] == 0) + t("陈旧天数 = 上一版 5 天 + 两个计划日相隔 2 天 = 7", a["stale_days"] == 7) + t("上一版采信倾向带出来", a["leaning_prev"] == "偏多" and a["leaning"] == "偏多") + other = next(r for r in rows if r["cluster_key"] == "Industry::固态电池") + t("上一版里没有的主题按第一次见到", other["migrated"] is None and other["stale_days"] == 0) + + # 指纹本身缺失(环节评析那行没有指纹):不判迁移,但日子照样变老。 + prev2 = {"Segment:正极材料": {"plan_date": "2026-09-02", "leaning": "中性", + "input_version": None, "stale_days": 3}} + rows = judgement.build_rows(DAY, sources.judgement_rows(read_pg=_pg_ok), prev2, now="x") + seg = next(r for r in rows if r["cluster_key"] == "Segment:正极材料") + t("指纹缺失:不判迁移,陈旧天数照常累加", seg["migrated"] is None and seg["stale_days"] == 4) + + # 计划日之间隔了一整个周末也照样按自然日数累加。 + prev3 = {"Industry::液冷散热": {"plan_date": "2026-08-28", "leaning": "偏多", + "input_version": "fp-a-1", "stale_days": 0}} + rows = judgement.build_rows("2026-08-31", sources.judgement_rows(read_pg=_pg_ok), prev3, now="x") + a = next(r for r in rows if r["cluster_key"] == "Industry::液冷散热") + t("跨周末按自然日数累加 3 天", a["stale_days"] == 3) + + +def test_migration_on_fingerprint_change(): + """指纹变化才算一次迁移:陈旧天数归零,上一版采信倾向留在 leaning_prev 里。""" + prev = {"Industry::液冷散热": {"plan_date": "2026-09-02", "cluster_key": "Industry::液冷散热", + "leaning": "偏多", "input_version": "fp-a-1", "stale_days": 9}} + + def _fetch_changed(sql=None, params=None): + return _view_rows(fp_a="fp-a-2", leaning_a="偏空") + + rows = judgement.build_rows(DAY, sources.judgement_rows(read_pg=_fetch_changed), prev, now="x") + a = next(r for r in rows if r["cluster_key"] == "Industry::液冷散热") + t("指纹变化:算一次迁移", a["migrated"] == 1) + t("迁移时陈旧天数归零", a["stale_days"] == 0) + t("采信倾向从偏多变成偏空,两版都在行上", + a["leaning_prev"] == "偏多" and a["leaning"] == "偏空") + + # 同一份材料指纹没变、只是采信倾向的文本被重新生成时,不算迁移(口径就是只认指纹)。 + def _fetch_same_fp(sql=None, params=None): + return _view_rows(fp_a="fp-a-1", leaning_a="偏空") + + rows = judgement.build_rows(DAY, sources.judgement_rows(read_pg=_fetch_same_fp), prev, now="x") + a = next(r for r in rows if r["cluster_key"] == "Industry::液冷散热") + t("指纹没变就不算迁移,哪怕采信倾向的文本变了", + a["migrated"] == 0 and a["stale_days"] == 10) + + hist = [{"plan_date": "2026-09-01", "migrated": None, "leaning": "偏多", "leaning_prev": None, + "input_version": "fp-a-1", "review_date": "2026-08-19"}, + {"plan_date": "2026-09-02", "migrated": 0, "leaning": "偏多", "leaning_prev": "偏多", + "input_version": "fp-a-1", "review_date": "2026-08-19"}, + {"plan_date": DAY, "migrated": 1, "leaning": "偏空", "leaning_prev": "偏多", + "input_version": "fp-a-2", "review_date": "2026-09-02"}] + got = judgement.migrations(hist) + t("版本史里只摘出被判为迁移的那一天,带方向与出处", + len(got) == 1 and got[0]["plan_date"] == DAY and got[0]["leaning_from"] == "偏多" + and got[0]["leaning_to"] == "偏空" and got[0]["review_date"] == "2026-09-02") + t("空版本史不报错", judgement.migrations([]) == [] and judgement.migrations(None) == []) + + +def main(): + test_fetch() + test_fetch_failure() + test_prev_read_failure() + test_same_day_idempotent() + test_stale_days_accumulate() + test_migration_on_fingerprint_change() + print("ALL OK — 行业观点快照:取数归一 / 读失败为空 / 同日重跑幂等 / 陈旧天数累加 / " + "指纹变化算迁移 / 版本史摘迁移 全部通过") + + +if __name__ == "__main__": + main() diff --git a/xxl.py b/xxl.py index 38b47a9..b746fde 100644 --- a/xxl.py +++ b/xxl.py @@ -24,6 +24,9 @@ tasks_periodic 的 notify_xxl),对平台表现完全一致: build 因子构建(run.py build all --mode daily) plan 生成当日选股计划文件 push-pool 计划写入股票池 + 触发决策系统增量补扫 +后面两步不在默认三步里,平台各建一个任务、显式传 steps 才跑: + regime-append 把决策系统 08:40 预热的市场区制写进当天计划快照的环境段(平台 08:45 触发) + judgement-snapshot 把数据基座当前这一版产业研判与环节评析抄一行,攒版本史(只写不判) """ import datetime as dt import os @@ -46,7 +49,7 @@ HERE = os.path.dirname(os.path.abspath(__file__)) LOG_PATH = os.path.join(HERE, "data", "xxl_build.log") # 步骤 → 命令(顺序即执行顺序;--date 由触发参数统一追加) -STEP_ORDER = ("build", "plan", "push-pool", "regime-append") +STEP_ORDER = ("build", "plan", "push-pool", "regime-append", "judgement-snapshot") STEP_CMDS = { "build": ["run.py", "build", "all", "--mode", "daily"], "plan": ["run.py", "plan"], @@ -55,6 +58,10 @@ STEP_CMDS = { # 桥 07:10 出计划时拿不到;平台在 08:45 另建一个任务只跑这一步,把当日区制写进当天的 # 计划快照 regime 段。它不在默认三步里,必须显式 steps=regime-append 才跑。 "regime-append": ["run.py", "regime-append"], + # 日频行业观点快照(2026-09-03 方案第四之五之三节,落地顺序第一步):把数据基座当前这一版 + # 产业研判与环节评析抄进选股系统自己的表,攒版本史,只写不判。它与出计划互不依赖,同样 + # 不在默认三步里,必须显式 steps=judgement-snapshot 才跑。 + "judgement-snapshot": ["run.py", "judgement-snapshot"], } DEFAULT_STEPS = ("build", "plan", "push-pool") STEP_TIMEOUT_SEC = 3600 # 单步上限一小时,防呆死(正常盘前链全程分钟级) @@ -162,7 +169,9 @@ def trigger_daily_build( steps: str = Query(",".join(DEFAULT_STEPS), description="要跑哪几步,逗号分隔;默认三步 build→plan→push-pool。" "环境追加 regime-append 不在默认里,平台 08:45 的任务单独传 " - "steps=regime-append(决策系统区制快照 08:40 才预热)"), + "steps=regime-append(决策系统区制快照 08:40 才预热);" + "行业观点快照 judgement-snapshot 同样不在默认里," + "平台单独建任务传 steps=judgement-snapshot"), date: str = Query(None, description="补跑指定数据日 YYYY-MM-DD;空=最新数据日"), key: str = Query(None), x_job_key: str = Header(None),