From 37f7ace34bbfa25995539456942cdbf75ffd1ab9 Mon Sep 17 00:00:00 2001 From: zlt Date: Thu, 30 Jul 2026 14:09:02 +0800 Subject: [PATCH] =?UTF-8?q?=E8=AE=A1=E5=88=92=20API=EF=BC=9A=E6=A1=A5?= =?UTF-8?q?=E5=AE=B9=E5=99=A8=E5=B8=B8=E9=A9=BB=E6=94=B9=20uvicorn(:8300)?= =?UTF-8?q?=EF=BC=8CGET=20/plan=20=E5=8F=96=E6=AF=8F=E6=97=A5=E9=80=89?= =?UTF-8?q?=E8=82=A1=E8=AE=A1=E5=88=92=EF=BC=9B=20plan.py=20=E6=8B=86?= =?UTF-8?q?=E8=A3=85=E9=85=8D/=E6=B8=B2=E6=9F=93=E4=B8=A4=E5=B1=82?= =?UTF-8?q?=E4=BE=9B=20API=20=E4=B8=8E=20cron=20=E5=85=B1=E7=94=A8?= =?UTF-8?q?=EF=BC=9Brequirements=20=E5=8A=A0=20fastapi/uvicorn?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- api.py | 64 +++++++++++++ docker-compose.yml | 9 +- plan.py | 221 ++++++++++++++++++++++++++++----------------- requirements.txt | 2 + 4 files changed, 210 insertions(+), 86 deletions(-) create mode 100644 api.py diff --git a/api.py b/api.py new file mode 100644 index 0000000..0782413 --- /dev/null +++ b/api.py @@ -0,0 +1,64 @@ +"""桥侧计划 API(07-30 用户需求):对外提供每日选股计划。 + +容器常驻命令改为 uvicorn 后随容器启动(docker-compose 已配端口,默认 8300); +cron 的 docker exec 构建/出计划照旧,互不影响。局域网内部服务,v1 无鉴权。 + + GET /health 存活 + 最新计划日 + GET /plan 最新一天的计划(JSON) + GET /plan?date=2026-07-30 指定日期 + GET /plan?format=md Markdown 原文(浏览器直接可读) + GET /plan/dates 可用日期列表 + POST /plan/refresh?date=... 重新生成该日计划文件(data/plan/*.md) + +将来接 XXL-JOB / 事件回调,触发器打这层即可,不必进容器。 +""" +import pandas as pd +from fastapi import FastAPI, HTTPException +from fastapi.responses import PlainTextResponse + +import db +import plan + +app = FastAPI(title="akg-factor-bridge · 每日选股计划", version="0.1") + + +@app.get("/health") +def health(): + try: + d = plan._latest_date("t_factor_akg_score") # noqa: SLF001 —— 桥内自用 + except Exception as e: # noqa: BLE001 —— 库连不上也要能回答"我还活着" + return {"ok": False, "error": repr(e)} + return {"ok": True, "latest_plan_date": d} + + +@app.get("/plan/dates") +def plan_dates(limit: int = 30): + df = db.read_mysql( + "factor", "SELECT DISTINCT trade_date FROM t_factor_akg_score " + "ORDER BY trade_date DESC LIMIT %s", (int(limit),)) + if df.empty: + return {"dates": []} + return {"dates": [pd.Timestamp(x).date().isoformat() + for x in df["trade_date"]]} + + +@app.get("/plan") +def get_plan(date: str | None = None, format: str = "json", + top: int = 20, obs_top: int = 10, theme_cap: int = 5): + try: + data = plan.collect(date, top, obs_top, theme_cap) + except RuntimeError as e: + raise HTTPException(status_code=404, detail=str(e)) + if format == "md": + return PlainTextResponse(plan.render_md(data), + media_type="text/markdown; charset=utf-8") + return data + + +@app.post("/plan/refresh") +def refresh(date: str | None = None): + try: + out = plan.generate(date) + except SystemExit as e: + raise HTTPException(status_code=404, detail=str(e)) + return {"ok": True, "file": out} diff --git a/docker-compose.yml b/docker-compose.yml index 12d99bd..a507347 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,13 +1,16 @@ # akg-factor-bridge:独立部署单元,可落在任意能同时连通「基座PG/153/平台MySQL」的服务器。 -# 容器常驻(sleep infinity),由宿主 cron 或平台 XXL-JOB 以 docker exec 触发 build; -# 也可改 command 为一次性任务由外部调度拉起。 +# 容器常驻并提供计划 API(uvicorn :8300,见 api.py,07-30);构建/出计划仍由宿主 +# cron 以 docker exec 触发。将来接 XXL-JOB / 事件回调,触发器打这层 API 即可。 services: akg-factor-bridge: build: . container_name: akg_factor_bridge env_file: .env restart: unless-stopped - command: sleep infinity + # uvicorn 起不来时(如镜像未重建、缺 fastapi)退回 sleep——保住 cron 的 docker exec + command: sh -c "uvicorn api:app --host 0.0.0.0 --port 8300 || sleep infinity" + ports: + - "${BRIDGE_API_PORT:-8300}:8300" volumes: - .:/app # 代码卷挂载:改码即生效,不用重构镜像(开发期)。冻结成 prod 镜像时删此行并 --build。 # 跨机部署(默认):.env 三处用 LAN IP,默认 bridge 网络出网到 LAN 即可。 diff --git a/plan.py b/plan.py index 5274e57..9d30d23 100644 --- a/plan.py +++ b/plan.py @@ -1,11 +1,15 @@ -"""每日选股计划(R4 首版):把 akg_gate / akg_score 变成一份人能读的榜单。 +"""每日选股计划(R4):数据装配 / Markdown 渲染 / 产出,三段分离。 + + collect() -> dict 结构化计划——api.py 直接当 JSON 返回 + render_md() -> str 从 dict 渲染 Markdown + generate() CLI 与 cron 的入口:collect + render + 落盘 + 打印 数据全部来自已落库的表,不重算: 平台因子表 t_factor_akg_score / _gate / _upside / _heat —— 当日截面 - 基座只读视图 v_factor_transmission —— 传导证据(主题、源数、已动比例) + 基座只读视图 v_factor_transmission —— 传导证据 基座 industry_pools —— 股票名称 -产出:终端打印 + data/plan/plan_<日期>.md。升降档一节对比前一交易日的档位表, -体现"数据到达本身是信号"(首次覆盖 / 新进传导链即升档)。 +升降档一节对比前一交易日的档位表——数据到达本身是信号(首次覆盖 / +新进传导链即升档)。 """ import json import os @@ -86,46 +90,24 @@ def _tier_label(score: float) -> str: int((score - 190.0) // 20), "?") -def _pct(v) -> str: - return "—" if v is None or pd.isna(v) else f"{v:+.0%}" +def _val(series: pd.Series, k: str): + v = series.get(k) + return None if v is None or pd.isna(v) else float(v) -def _num(v) -> str: - return "—" if v is None or pd.isna(v) else f"{v:.2f}" - - -def _pick(ranked: pd.Series, ev: dict, top: int, theme_cap: int): - """按分数从高到低取 top 条;每个传导主题最多 theme_cap 条(0=不设限)。 - - 为什么设限:传导目标是环节级,同环节全体成员共享同一条证据,不设限时 - 榜单会被三四个环节刷屏——20 个名额实际只是 4 注。限额后变成 - "多条线索 × 每条线索取最冷最便宜的几只",被挤掉的仍在完整档位表里。""" - out, cnt = [], {} - for k, s in ranked.items(): - e = ev.get(k) - theme = e[0] if e else "(无传导)" - if theme_cap and cnt.get(theme, 0) >= theme_cap: - continue - cnt[theme] = cnt.get(theme, 0) + 1 - out.append((k, s)) - if len(out) >= top: - break - return out - - -def generate(date: str | None = None, top: int = 20, obs_top: int = 10, - theme_cap: int = 5) -> str: +def collect(date: str | None = None, top: int = 20, obs_top: int = 10, + theme_cap: int = 5) -> dict: + """装配一天的计划为结构化字典。数据缺失抛 RuntimeError(api 侧转 404)。""" ds = date or _latest_date("t_factor_akg_score") if not ds: - raise SystemExit("t_factor_akg_score 还没有数据——先 build akg_score。") + raise RuntimeError("t_factor_akg_score 还没有数据——先 build akg_score。") score = _factor("t_factor_akg_score", ds) gate = _factor("t_factor_akg_gate", ds) if score.empty or gate.empty: - raise SystemExit(f"{ds} 缺 akg_score / akg_gate——先 build 该日再出计划。") + raise RuntimeError(f"{ds} 缺 akg_score / akg_gate——先 build 该日再出计划。") upside = _factor("t_factor_akg_upside", ds) if upside.empty: - # 晚间 18:40 建当日 upside 表时行情源(~19:50 发布)还没到,当日表常为空。 - # 这里现算兜底:此刻库里已有晚到的当日价,as-of 口径不变(consensus<=当日)。 + # 当日 upside 表为空时现算兜底(as-of 口径不变:consensus<=当日、当日收盘价) import factors df_up = factors.build_upside(ds, ds) if df_up is not None and not df_up.empty: @@ -140,81 +122,154 @@ def generate(date: str | None = None, top: int = 20, obs_top: int = 10, main = score[score >= _MAIN_MIN].sort_values(ascending=False) obs = score[score < _MAIN_MIN].sort_values(ascending=False) - L = [f"# 每日选股计划 · {ds}", ""] - L.append(f"主榜 {len(main)} 只 / 观察档 {len(obs)} 只 / 全池档位覆盖 {len(gate)} 只。") - stale = sorted(d for d in mkt_days if d != ds) + def _pick(ranked: pd.Series, n: int): + """分数从高到低取 n 条;每个传导主题最多 theme_cap 条(0=不设限)—— + 传导目标是环节级、同环节成员共享同一条证据,不限额会被少数环节刷屏。""" + out, cnt = [], {} + for k, s in ranked.items(): + e = ev.get(k) + theme = e[0] if e else "(无传导)" + if theme_cap and cnt.get(theme, 0) >= theme_cap: + continue + cnt[theme] = cnt.get(theme, 0) + 1 + out.append((k, s)) + if len(out) >= n: + break + return out + + def _row(rank: int, k: str, s: float, with_tier: bool) -> dict: + e = ev.get(k) + r = {"rank": rank, "code": k, "name": names.get(k), + "score": round(float(s), 2), + "evidence": ({"theme": e[0], "n_sources": e[1], + "moved_ratio": round(e[2], 4)} if e else None), + "heat": _val(heat, k), "upside": _val(upside, k)} + if with_tier: + r["tier"] = _tier_label(s) + return r + + changes = None + prev_ds = _prev_date("t_factor_akg_gate", ds) + if prev_ds: + prev = _factor("t_factor_akg_gate", prev_ds) + both = pd.concat([prev.rename("prev"), gate.rename("cur")], + axis=1).fillna(-1.0) # -1 = 当日不在面板 + lab = {-1.0: "池外", 0.0: "不采纳", 1.0: "观察档", 2.0: "主榜"} + up_df = both[both["cur"] > both["prev"]].sort_values("cur", ascending=False) + down_df = both[both["cur"] < both["prev"]].sort_values("prev", ascending=False) + + def _mv(d: pd.DataFrame): + return [{"code": k, "name": names.get(k), + "from": lab.get(r["prev"], "?"), "to": lab.get(r["cur"], "?")} + for k, r in d.iterrows()] + + changes = {"base_date": prev_ds, + "upgrades_total": int(len(up_df)), + "downgrades_total": int(len(down_df)), + "upgrades": _mv(up_df.head(15)), + "downgrades": _mv(down_df.head(15))} + + return { + "date": ds, + "counts": {"main": int(len(main)), "observe": int(len(obs)), + "gate_covered": int(len(gate))}, + "market_snapshot_days": sorted(mkt_days), + "heat_date": hd, + "theme_cap": theme_cap, + "main": [_row(i, k, s, True) + for i, (k, s) in enumerate(_pick(main, top), 1)], + "observe": [_row(i, k, s, False) + for i, (k, s) in enumerate(_pick(obs, obs_top), 1)], + "changes": changes, + "encoding": "主榜分=200+传导档位×20+组内分(还没热、还便宜);" + "观察档分=100+0.6z(传导)+0.4z(−热度)", + } + + +def _fmt_pct(v) -> str: + return "—" if v is None else f"{v:+.0%}" + + +def _fmt_num(v) -> str: + return "—" if v is None else f"{v:.2f}" + + +def _fmt_ev(e) -> str: + if not e: + return "—" + return f"{e['theme']}({e['n_sources']} 源,已动 {e['moved_ratio']:.0%})" + + +def render_md(d: dict) -> str: + L = [f"# 每日选股计划 · {d['date']}", ""] + c = d["counts"] + L.append(f"主榜 {c['main']} 只 / 观察档 {c['observe']} 只 / " + f"全池档位覆盖 {c['gate_covered']} 只。") + stale = [x for x in d["market_snapshot_days"] if x != d["date"]] if stale: - L.append(f"⚠️ 本日传导用的行情快照 = {'、'.join(stale)}(T−1 口径:" - f"\"谁已经动了\"看的是上个交易日收盘;拍点方案定版前均如此)。") + L.append(f"注:本日传导用的行情快照 = {'、'.join(stale)}" + f"(与计划日不同——历史降级日口径)。") L.append("") - cap_txt = f",每主题限额 {theme_cap}" if theme_cap else "" - main_rows = _pick(main, ev, top, theme_cap) - L.append(f"## 主榜 Top {len(main_rows)}(有券商预期、目标价不低于现价{cap_txt})") + cap_txt = f",每主题限额 {d['theme_cap']}" if d["theme_cap"] else "" + L.append(f"## 主榜 Top {len(d['main'])}(有券商预期、目标价不低于现价{cap_txt})") L.append("") L.append("| # | 代码 | 名称 | 总分 | 档位 | 传导证据 | 热度 | 预期空间 |") L.append("|---|------|------|------|------|----------|------|----------|") - for i, (k, s) in enumerate(main_rows, 1): - e = ev.get(k) - etxt = f"{e[0]}({e[1]} 源,已动 {e[2]:.0%})" if e else "—" - L.append(f"| {i} | {k} | {names.get(k, '—')} | {s:.1f} | {_tier_label(s)} " - f"| {etxt} | {_num(heat.get(k))} | {_pct(upside.get(k))} |") + for r in d["main"]: + L.append(f"| {r['rank']} | {r['code']} | {r['name'] or '—'} | {r['score']:.1f} " + f"| {r['tier']} | {_fmt_ev(r['evidence'])} " + f"| {_fmt_num(r['heat'])} | {_fmt_pct(r['upside'])} |") L.append("") - obs_rows = _pick(obs, ev, obs_top, theme_cap) - L.append(f"## 观察档 Top {len(obs_rows)}" + L.append(f"## 观察档 Top {len(d['observe'])}" f"(无券商预期、但在传导链上——没有估值锚,置信度低{cap_txt})") L.append("") L.append("| # | 代码 | 名称 | 分 | 传导证据 | 热度 |") L.append("|---|------|------|----|----------|------|") - for i, (k, s) in enumerate(obs_rows, 1): - e = ev.get(k) - etxt = f"{e[0]}({e[1]} 源,已动 {e[2]:.0%})" if e else "—" - L.append(f"| {i} | {k} | {names.get(k, '—')} | {s:.1f} " - f"| {etxt} | {_num(heat.get(k))} |") + for r in d["observe"]: + L.append(f"| {r['rank']} | {r['code']} | {r['name'] or '—'} | {r['score']:.1f} " + f"| {_fmt_ev(r['evidence'])} | {_fmt_num(r['heat'])} |") L.append("") L.append("## 今日升降档") L.append("") - prev_ds = _prev_date("t_factor_akg_gate", ds) - if not prev_ds: + ch = d["changes"] + if not ch: L.append("(没有更早的档位表可比,升降档从下一个交易日开始。)") else: - prev = _factor("t_factor_akg_gate", prev_ds) - both = pd.concat([prev.rename("prev"), gate.rename("cur")], axis=1) - both = both.fillna(-1.0) # -1 = 当日不在面板 - up = both[both["cur"] > both["prev"]] - down = both[both["cur"] < both["prev"]] - lab = {-1.0: "池外", 0.0: "不采纳", 1.0: "观察档", 2.0: "主榜"} - L.append(f"对比 {prev_ds}:升档 {len(up)} 只,降档 {len(down)} 只。" + L.append(f"对比 {ch['base_date']}:升档 {ch['upgrades_total']} 只," + f"降档 {ch['downgrades_total']} 只。" f"升档=拿到新锚(首次覆盖 / 新进传导链),本身就是值得看的信号。") - - def _rows(d: pd.DataFrame, cap: int = 15): - lines = [] - for k, r in d.iterrows(): - lines.append(f"- {k} {names.get(k, '')}:" - f"{lab.get(r['prev'], '?')} → {lab.get(r['cur'], '?')}") - if len(lines) >= cap: - lines.append(f"- ……共 {len(d)} 只,其余见档位表") - break - return lines - - if not up.empty: + if ch["upgrades"]: L.append("") L.append("**升档**:") - L += _rows(up.sort_values("cur", ascending=False)) - if not down.empty: + L += [f"- {m['code']} {m['name'] or ''}:{m['from']} → {m['to']}" + for m in ch["upgrades"]] + if ch["upgrades_total"] > len(ch["upgrades"]): + L.append(f"- ……共 {ch['upgrades_total']} 只,其余见档位表") + if ch["downgrades"]: L.append("") L.append("**降档**:") - L += _rows(down.sort_values("prev", ascending=False)) + L += [f"- {m['code']} {m['name'] or ''}:{m['from']} → {m['to']}" + for m in ch["downgrades"]] + if ch["downgrades_total"] > len(ch["downgrades"]): + L.append(f"- ……共 {ch['downgrades_total']} 只,其余见档位表") L.append("") L.append("---") - L.append("口径:主榜分 = 200 + 传导档位×20 + 组内分(还没热、还便宜);" - "观察档分 = 100 + 0.6z(传导) + 0.4z(−热度)。") + L.append(f"口径:{d['encoding']}。") + return "\n".join(L) - text = "\n".join(L) + +def generate(date: str | None = None, top: int = 20, obs_top: int = 10, + theme_cap: int = 5) -> str: + try: + data = collect(date, top, obs_top, theme_cap) + except RuntimeError as e: + raise SystemExit(str(e)) + text = render_md(data) os.makedirs("data/plan", exist_ok=True) - out = f"data/plan/plan_{ds}.md" + out = f"data/plan/plan_{data['date']}.md" with open(out, "w", encoding="utf-8") as f: f.write(text + "\n") print(text) diff --git a/requirements.txt b/requirements.txt index 8319cc0..33659cb 100644 --- a/requirements.txt +++ b/requirements.txt @@ -4,3 +4,5 @@ psycopg[binary]>=3.1 PyMySQL>=1.1 python-dotenv>=1.0 PyYAML>=6.0 +fastapi>=0.110 +uvicorn>=0.29