115 lines
4.9 KiB
Python
115 lines
4.9 KiB
Python
"""akg-factor-bridge CLI。
|
||
|
||
python run.py views # 连通性自检:打印视图/表行数
|
||
python run.py register # 注册四子因子到 factor_metadata
|
||
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 则取今天;start=end=date。history 模式:需 --start/--end。
|
||
所有写入幂等(删涉及日期区间再插),可安全重跑。
|
||
"""
|
||
import argparse
|
||
import datetime as dt
|
||
import warnings
|
||
|
||
# pandas 用 DBAPI 连接读 SQL 会 warn(无害,功能正常)——静音保持日志干净
|
||
warnings.filterwarnings("ignore", message=".*only supports SQLAlchemy.*")
|
||
|
||
import common
|
||
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}")
|
||
|
||
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}")
|
||
|
||
|
||
_META = {
|
||
"akg_upside": ("astock-kg 预期空间", "分析师一致预期目标价隐含收益率(target_mid/price-1)"),
|
||
"akg_heat": ("astock-kg 热度", "生态日频资金热度分(0~1)"),
|
||
"akg_event": ("astock-kg 事件", "利好利空事件时间衰减加权分"),
|
||
"akg_transmission": ("astock-kg 传导", "板块传导未动成员传导强度(路径数×(1-已动比例))"),
|
||
}
|
||
|
||
|
||
def cmd_register():
|
||
print("注册四子因子:")
|
||
for code, (name, desc) in _META.items():
|
||
common.register(code, name, factors.FACTORS[code], ["astock-kg", code.split("_", 1)[1]], desc)
|
||
|
||
|
||
def cmd_build(which, mode, start, end, date):
|
||
if mode == "daily":
|
||
d = date or dt.date.today().isoformat()
|
||
start = end = d
|
||
if not start or not end:
|
||
raise SystemExit("history 模式需要 --start 与 --end")
|
||
codes = list(factors.FACTORS) if which == "all" else [which]
|
||
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)
|
||
common.write_factor(factors.FACTORS[code], df, mode)
|
||
except Exception as e: # noqa: BLE001 —— 一路失败不拖累其余(批量容错)
|
||
print(f" ❌ {code} 失败: {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")
|
||
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")
|
||
a = ap.parse_args()
|
||
|
||
if a.cmd == "views":
|
||
cmd_views()
|
||
elif a.cmd == "register":
|
||
cmd_register()
|
||
elif a.cmd == "build":
|
||
cmd_build(a.factor, a.mode, a.start, a.end, a.date)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|