计划 API:桥容器常驻改 uvicorn(:8300),GET /plan 取每日选股计划;

plan.py 拆装配/渲染两层供 API 与 cron 共用;requirements 加 fastapi/uvicorn
This commit is contained in:
zlt 2026-07-30 14:09:02 +08:00
parent d00c1f262a
commit 37f7ace34b
4 changed files with 210 additions and 86 deletions

64
api.py Normal file
View File

@ -0,0 +1,64 @@
"""桥侧计划 API07-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}

View File

@ -1,13 +1,16 @@
# akg-factor-bridge独立部署单元可落在任意能同时连通「基座PG/153/平台MySQL」的服务器。
# 容器常驻sleep infinity由宿主 cron 或平台 XXL-JOB 以 docker exec 触发 build
# 也可改 command 为一次性任务由外部调度拉起
# 容器常驻并提供计划 APIuvicorn :8300见 api.py07-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 即可。

221
plan.py
View File

@ -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:
"""装配一天的计划为结构化字典。数据缺失抛 RuntimeErrorapi 侧转 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)}T1 口径:"
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)

View File

@ -4,3 +4,5 @@ psycopg[binary]>=3.1
PyMySQL>=1.1
python-dotenv>=1.0
PyYAML>=6.0
fastapi>=0.110
uvicorn>=0.29