270 lines
14 KiB
Python
270 lines
14 KiB
Python
"""akg-factor-bridge CLI。
|
||
|
||
python run.py views # 连通性自检:打印视图/表行数
|
||
python run.py apply-views [--dry-run] # 把插槽视图 DDL 应用到基座 PG
|
||
python run.py probe # G1 体检(只读,详见 probe.py)
|
||
python run.py freeze [--date D] # 输入冻结(G0.5,详见 freeze.py)
|
||
python run.py tracks # 赛道覆盖体检 + 成员表快照(G2)
|
||
python run.py plan [--date D] [--top N] # 每日选股计划(R4,读已落库因子表)
|
||
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
|
||
|
||
daily 模式:不给 --date 则取基座最新扫描日(盘前链口径,07-30 D 案);
|
||
start=end=date,跑完自动冻结(--no-freeze 可关)。
|
||
history 模式:需 --start/--end,默认不逐日冻结。
|
||
所有写入幂等(删涉及日期区间再插),可安全重跑。
|
||
"""
|
||
import argparse
|
||
import datetime as dt
|
||
import warnings
|
||
|
||
# pandas 用 DBAPI 连接读 SQL 会 warn(无害,功能正常)——静音保持日志干净
|
||
warnings.filterwarnings("ignore", message=".*only supports SQLAlchemy.*")
|
||
|
||
import pandas as pd
|
||
|
||
import common
|
||
import config
|
||
import db
|
||
import factors
|
||
|
||
|
||
def cmd_views():
|
||
checks = [
|
||
("PG v_factor_universe", "pg", "SELECT count(*) FROM v_factor_universe"),
|
||
("PG v_factor_consensus", "pg", "SELECT count(*) FROM v_factor_consensus"),
|
||
("PG v_factor_events", "pg", "SELECT count(*) FROM v_factor_events"),
|
||
("PG v_factor_transmission", "pg", "SELECT count(*) FROM v_factor_transmission"),
|
||
("153 stock_fund_heat_scores","heat", "SELECT count(*) FROM stock_fund_heat_scores"),
|
||
("平台 gp_day_data", "price", "SELECT count(*) FROM gp_day_data"),
|
||
("平台 factor_metadata", "factor", "SELECT count(*) FROM factor_metadata"),
|
||
]
|
||
print("连通性自检:")
|
||
for name, src, sql in checks:
|
||
try:
|
||
df = db.read_pg(sql) if src == "pg" else db.read_mysql(src, sql)
|
||
print(f" ✅ {name}: {int(df.iloc[0, 0])}")
|
||
except Exception as e: # noqa: BLE001
|
||
print(f" ❌ {name}: {e!r}")
|
||
|
||
# 视图版本自检:v2 才有的列在不在(评审第一批是否已应用)
|
||
print("\n视图版本(2026-07-26 评审 v2):")
|
||
for label, sql in (
|
||
("v_factor_transmission.n_sources",
|
||
"SELECT n_sources FROM v_factor_transmission LIMIT 1"),
|
||
("v_factor_transmission.mkt_trade_date",
|
||
"SELECT mkt_trade_date FROM v_factor_transmission LIMIT 1"),
|
||
("v_factor_events.source_type",
|
||
"SELECT source_type FROM v_factor_events LIMIT 1"),
|
||
("v_factor_segment_members(第五插槽,G2)",
|
||
"SELECT 1 FROM v_factor_segment_members LIMIT 1")):
|
||
try:
|
||
db.read_pg(sql)
|
||
print(f" ✅ {label}")
|
||
except Exception: # noqa: BLE001
|
||
print(f" ⬜ {label} —— 未就绪")
|
||
|
||
fronts = [
|
||
("热度 heat", "heat", "trade_date", "stock_fund_heat_scores"),
|
||
("一致预期 consensus", "pg", "asof_date", "v_factor_consensus"),
|
||
("事件 events", "pg", "disclosure_date", "v_factor_events"),
|
||
("传导 transmission", "pg", "scan_date", "v_factor_transmission"),
|
||
("行情 gp_day_data", "price", "`timestamp`", "gp_day_data"),
|
||
]
|
||
print("\n数据历史深度(min ~ max,distinct 天数——回填范围据此定):")
|
||
for name, src, col, tbl in fronts:
|
||
sql = f"SELECT MIN({col}), MAX({col}), COUNT(DISTINCT {col}) FROM {tbl}"
|
||
try:
|
||
df = db.read_pg(sql) if src == "pg" else db.read_mysql(src, sql)
|
||
lo, hi, n = df.iloc[0, 0], df.iloc[0, 1], df.iloc[0, 2]
|
||
print(f" {name}: {lo} ~ {hi} ({int(n)} 天)")
|
||
except Exception as e: # noqa: BLE001
|
||
print(f" {name}: ❌ {e!r}")
|
||
|
||
print(f"\n当前口径:SUBFACTOR_UNIVERSE={config.SUBFACTOR_UNIVERSE} "
|
||
f"| EVENT_SOURCE_TYPES={sorted(config.EVENT_SOURCE_TYPES)} "
|
||
f"| EVENT_MAX_PER_DOC={config.EVENT_MAX_PER_DOC}")
|
||
|
||
|
||
_META = {
|
||
"akg_upside": ("astock-kg 预期空间", "分析师一致预期目标价隐含收益率(target_mid/price-1)"),
|
||
"akg_heat": ("astock-kg 热度", "生态日频资金热度分(0~1)"),
|
||
"akg_event": ("astock-kg 事件", "利好利空事件时间衰减加权分(仅公告来源,单文档封顶)"),
|
||
"akg_transmission": ("astock-kg 传导", "板块传导未动成员传导强度(distinct源数×(1-已动比例))"),
|
||
"akg_gate": ("astock-kg 门槛档位",
|
||
"三档置信门槛(07-30拍板): 2=主榜(券商覆盖且upside>=0; 赛道C闸未启用, "
|
||
"转正后再交赛道成员), 1=观察档(无券商覆盖但在图谱传导链上, 无估值锚, "
|
||
"低置信), 0=不采纳(两锚皆无, 或upside<0)。全池出行, 让平台看得见门槛"),
|
||
"akg_score": ("astock-kg 景气度漏斗",
|
||
"先档后分: 主榜=200+传导档位x20+组内分(0.6z(-热度)+0.4z(upside), "
|
||
"两段式07-30拍板, 组内分clip±9.9故档间不重叠); "
|
||
"观察档=100+0.6z(传导)+0.4z(-热度)。仅gate>0出行, 数值直接可排序"),
|
||
}
|
||
|
||
|
||
def cmd_register():
|
||
print(f"注册因子(共 {len(_META)} 个):")
|
||
for code, (name, desc) in _META.items():
|
||
common.register(code, name, factors.FACTORS[code], ["astock-kg", code.split("_", 1)[1]], desc)
|
||
|
||
|
||
def _latest_data_day() -> str:
|
||
"""daily 不传 --date 时的默认日:基座传导台账最新 scan_date(=数据日)。
|
||
|
||
盘前链(07-30 D 案)在次日早晨补全上一交易日,"今天"多半还没有数据;
|
||
以基座刚完成的扫描日为准,两边永远对齐。查不到再退回今天。"""
|
||
try:
|
||
df = db.read_pg("SELECT MAX(scan_date) d FROM v_factor_transmission")
|
||
v = None if df.empty else df.iloc[0, 0]
|
||
if v is not None and not pd.isna(v):
|
||
return pd.Timestamp(v).date().isoformat()
|
||
except Exception as e: # noqa: BLE001 —— 视图不可达时退回今天,不阻塞构建
|
||
print(f" (取基座最新扫描日失败,默认改用今天: {e!r})")
|
||
return dt.date.today().isoformat()
|
||
|
||
|
||
def cmd_build(which, mode, start, end, date, do_freeze=True):
|
||
if mode == "daily":
|
||
d = date or _latest_data_day()
|
||
start = end = d
|
||
if not start or not end:
|
||
raise SystemExit("history 模式需要 --start 与 --end")
|
||
codes = list(factors.FACTORS) if which == "all" else [which]
|
||
frames = {}
|
||
for code in codes:
|
||
if code not in factors.BUILDERS:
|
||
raise SystemExit(f"未知因子: {code}(可选: {list(factors.FACTORS)} 或 all)")
|
||
print(f"[{code}] {mode} {start} ~ {end}")
|
||
try:
|
||
df = factors.BUILDERS[code](start, end)
|
||
frames[f"factor_{code}"] = df
|
||
common.write_factor(factors.FACTORS[code], df, mode)
|
||
except Exception as e: # noqa: BLE001 —— 一路失败不拖累其余(批量容错)
|
||
print(f" ❌ {code} 失败: {e!r}")
|
||
# 日更顺手冻结:输入与输出落在同一目录,任何一行因子值都能被逐步复算。
|
||
# history 模式默认不冻结(逐日冻结应单独跑,避免一次回填写出几百个目录)。
|
||
if do_freeze and mode == "daily":
|
||
try:
|
||
import freeze
|
||
freeze.snapshot(start, extra_frames=frames)
|
||
except Exception as e: # noqa: BLE001 —— 冻结失败不该让因子构建算失败
|
||
print(f" ❌ 冻结失败(因子已落库): {e!r}")
|
||
|
||
|
||
def main():
|
||
ap = argparse.ArgumentParser(description="akg-factor-bridge")
|
||
sub = ap.add_subparsers(dest="cmd", required=True)
|
||
sub.add_parser("views")
|
||
sub.add_parser("register")
|
||
sub.add_parser("tracks") # 赛道覆盖体检 + confirmed 成员表快照(G2)
|
||
pl = sub.add_parser("plan") # 每日选股计划(R4)
|
||
pl.add_argument("--date", help="默认取 score 表最新日")
|
||
pl.add_argument("--top", type=int, default=20, help="主榜条数")
|
||
pl.add_argument("--obs-top", type=int, default=10, help="观察档条数")
|
||
pl.add_argument("--theme-cap", type=int, default=5,
|
||
help="每个传导主题最多几条(防单板块刷屏;0=不设限)")
|
||
pp = sub.add_parser("push-pool") # 计划入池(08-03,写 Mongo 股票池分组)
|
||
pp.add_argument("--date", help="默认取 score 表最新日(与 plan 同口径)")
|
||
pp.add_argument("--top", type=int, help="计划取主榜前几只(默认读 POOL_TOP=20)")
|
||
pp.add_argument("--dry-run", action="store_true", help="只打印入池/出池明细,不写库")
|
||
pp.add_argument("--no-kick", action="store_true",
|
||
help="写完不触发决策系统增量补扫(当晚全量扫兜底)")
|
||
p = sub.add_parser("probe")
|
||
p.add_argument("--section", choices=["all", "pools", "price", "upside", "corr"],
|
||
default="all",
|
||
help="pools=池结构清单 price=行情年表 upside=分布与q档位 corr=三项相关矩阵")
|
||
av = sub.add_parser("apply-views")
|
||
av.add_argument("--file", default="sql/astock_kg_slot_views.sql")
|
||
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="计划日,默认今天;补写过去的日期会让版本史失真,会有提示")
|
||
pr = sub.add_parser("plan-review") # 复盘周报(每周五固定动作;只读,名单级,不是回测)
|
||
pr.add_argument("--date", help="以哪天为今天算滚动窗口,默认今天;调度中心统一追加的就是它")
|
||
pr.add_argument("--since", help="给了就按给定窗口跑,不再滚动")
|
||
pr.add_argument("--until")
|
||
pr.add_argument("--horizons", default="5,10,20")
|
||
pr.add_argument("--start-price", choices=["next_close", "signal_close"], default="next_close")
|
||
pr.add_argument("--out", default="data/review")
|
||
rv = sub.add_parser("review-targets") # 个股深度评析的目标名单:从已生成的计划文件补写(2026-09-09)
|
||
rv.add_argument("--date", required=True, help="计划日 YYYY-MM-DD,读 data/plan/plan_<date>.json")
|
||
f = sub.add_parser("freeze")
|
||
f.add_argument("--date", help="默认今天")
|
||
b = sub.add_parser("build")
|
||
b.add_argument("factor", help="akg_upside|akg_heat|akg_event|akg_transmission|all")
|
||
b.add_argument("--mode", choices=["daily", "history"], default="daily")
|
||
b.add_argument("--start")
|
||
b.add_argument("--end")
|
||
b.add_argument("--date")
|
||
b.add_argument("--no-freeze", action="store_true", help="daily 模式下跳过输入冻结")
|
||
a = ap.parse_args()
|
||
|
||
if a.cmd == "views":
|
||
cmd_views()
|
||
elif a.cmd == "probe":
|
||
import probe # 按需加载:一次性诊断命令,不影响常规链路
|
||
probe.run(a.section)
|
||
elif a.cmd == "apply-views":
|
||
import apply_views
|
||
raise SystemExit(1 if apply_views.apply(a.file, a.dry_run) else 0)
|
||
elif a.cmd == "freeze":
|
||
import freeze
|
||
freeze.snapshot(a.date)
|
||
elif a.cmd == "register":
|
||
cmd_register()
|
||
elif a.cmd == "plan":
|
||
import plan
|
||
plan.generate(a.date, a.top, a.obs_top, a.theme_cap)
|
||
elif a.cmd == "push-pool":
|
||
import pool
|
||
pool.push(a.date, a.top, dry_run=a.dry_run, kick=not a.no_kick)
|
||
elif a.cmd == "regime-append":
|
||
import plan
|
||
import regime
|
||
day = a.date or plan._latest_date("t_factor_akg_score") # noqa: SLF001 —— 同仓自用
|
||
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 == "review-targets":
|
||
import review_targets
|
||
print(review_targets.from_plan_file(a.date))
|
||
elif a.cmd == "plan-review":
|
||
import datetime as _dt
|
||
import plan_review
|
||
since, until = a.since, a.until
|
||
if not since and not until:
|
||
today = _dt.date.fromisoformat(a.date) if a.date else None
|
||
since, until = plan_review.weekly_window(today)
|
||
hs = tuple(int(x) for x in a.horizons.split(",") if x.strip())
|
||
print(f"[plan-review] 窗口 {since} ~ {until},期限 {hs},起算 {a.start_price}")
|
||
plan_review.run(since, until, hs, a.start_price, a.out)
|
||
elif a.cmd == "tracks":
|
||
import tracks
|
||
tracks.coverage_report()
|
||
out, df, missing = tracks.snapshot(only_confirmed=True)
|
||
print(f"\nconfirmed 成员表快照: {out}"
|
||
f"({df['ts_code'].nunique()} 只,{len(df)} 行 股票×赛道)")
|
||
if missing:
|
||
print(f"⚠️ {len(missing)} 个映射键未命中(明细见体检表各行)")
|
||
try:
|
||
tracks.gate_simulation()
|
||
except Exception as e: # noqa: BLE001 —— 演算失败不影响体检本体
|
||
print(f"(赛道闸演算失败: {e!r})")
|
||
elif a.cmd == "build":
|
||
cmd_build(a.factor, a.mode, a.start, a.end, a.date, do_freeze=not a.no_freeze)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|