存钱小帮手 Agent 的设计与实现
在线体验:https://agent.xiaofenglei.cn/(红黑榜页面:/markets)


这篇文章复盘我做的一个「存钱小帮手」Agent 产品(内部代号 SoundFinancialPlanning,对应小程序「稳健理财」Tab)。它不是一个简单的 ChatBot:用户输入一只股票或一个理财问题,后端会先把基本面数据采集齐全,再让段永平、巴菲特、芒格等十位「投资大师」各自从自己的框架出发做一轮圆桌发言,最后汇总成可执行的建议。另外还有一个「红黑榜」,每小时自动跑一遍全市场,把强势股和弱势股的详情预热好。
下面按「整体架构 → 前后端技术栈 → 后端 LLM 底层调用 → Skill 加载 → Workflow 编排 → 定时器」的顺序展开。
一、整体架构
产品分三个 sibling 工程,同处一个目录下,各自独立 git 仓:
| 工程 | 角色 | 技术 |
|---|---|---|
SoundFinancialPlanning/ |
微信小程序客户端 | TypeScript + LESS,自定义 tabBar |
SoundFinancialAgentFE/ |
Web 前端(agent.xiaofenglei.cn) | pnpm monorepo:Next.js 16 + React 19,packages/{domain,agent,data,db} |
SoundFinancialPlanningBE/ |
后端 API + 多 Agent 投研 Worker | Python 3.11 / FastAPI + RQ + Redis + MySQL |
小程序和 Web 共用同一套后端。后端是整个 Agent 的核心,所以下文重点讲后端,前端只点技术栈。
两个进程,一个镜像
docker-compose.yml 把同一个 Dockerfile 构建出的镜像起成两个服务:
- api 进程:
gunicorn app.main:app -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000 --timeout 0,处理 HTTP + WebSocket。--timeout 0是故意的——单次/analysis可能跑 10 分钟以上,但 WS 是长连接,不能被 gunicorn 提前 kill。 - worker 进程:
python -m app.workers.analysis_worker,跑 RQ consumer 监听任务队列。
两者共享 Redis(任务队列 + pubsub + 行情缓存)和 MySQL。业务数据落 MySQL,./data 承载 TradingAgents 的 logs/cache。
二、前后端技术栈
Web 前端(SoundFinancialAgentFE)
pnpm monorepo,包边界很严格:
apps/web Next.js 16 (App Router) + React 19.2 + Tailwind 4
packages/domain 领域模型(零 React/Next 依赖)
packages/agent Agent 调用封装
packages/data 数据层
packages/db Supabase 访问层
关键依赖:next@16.2.6、react@19.2.4、@supabase/ssr、openai@^6、zod@^4。测试用 Vitest + Playwright。packages/* 必须零 React/Next 依赖,靠 pnpm typecheck 在 CI 里卡住边界。
小程序端(SoundFinancialPlanning)
TypeScript + LESS,自定义 tabBar,页面有:holdings(我的持仓)、index(稳健理财)、chat(大师圆桌对话)、master-portfolio(大师组合)、analysis、redblack-detail(红黑榜详情)、mine。
小程序的 wxml 模板按字段名锁死,所以后端的 schema 字段必须严格对齐——后端适配前端,而不是反过来。
后端(SoundFinancialPlanningBE)
- Python 3.11+,FastAPI 0.115,gunicorn + uvicorn workers
- TradingAgents-astock 框架 vendored 进
app/agents/tradingagents/(不走pip install -e) - MySQL 8 + aiomysql(生产),SQLite + aiosqlite 仅本地兜底;无 migrations,启动
create_all - Redis 7:RQ 任务队列 + pubsub + 行情缓存
- LLM:MiniMax(OpenAI-compatible,
https://api.minimaxi.com/v1)+ Anthropic Claude,双 Provider 互为兜底
三、数据源:先说这个,是因为它比 LLM 还容易炸
LLM 挂了最多是"答得不准",但数据源挂了,整个圆桌的「基本面」就直接空掉,详情页会出现 当前价格: N/A; 市值: N/A; PE/PB/PS: None,这种页面对用户等于不可用。所以在做 LLM 调用之前,先把数据层稳下来。
app/services/data_collectors/ 目录是这一块的真实落点(提交历史里从 fundamentals_collector.py 一个 600+ 行的大文件,拆成了四个文件):
app/services/data_collectors/
__init__.py 模块导出
base.py StockQuote / OHLCVBar / Fundamentals / NewsItem dataclass
astock.py 377 行 A 股:mootdx + akshare + eastmoney
hk_us.py 675 行 港美股:yfinance + akshare fallback + sina quote fallback
坑一:HK/US ticker 通了,数据采集没跟上。 我加了 US ticker alias 识别(feat(ticker): recognize standalone US ticker symbols),AAPL/NVDA/NBIS 能解析了。但下游的 _l1_roundtable_fundamentals 不分市场,一股脑走 a_stock.get_fundamentals / get_news / get_concept_blocks / get_fund_flow,mootdx/akshare/eastmoney 只认 6 位 A 股代码,HK/US ticker 全部返回空。表现就是用户输入 AAPL,圆桌里出来 当前价格: N/A; 市值: N/A; 相关概念: Baidu PAE error。修法:在 roundtable_pipeline.py 里加了 _collect_fundamentals_hk_us(),专门走 data_collectors.hk_us 的 yfinance 通道,把返回 shape 重新映射成 astock 路径用的同一个 schema。A 股路径不动。
坑二:yfinance 限流。 美股大票还行,小票和港股基本动不动就 429。两层兜底:行情走 yfinance.fast_info → history → sina 三级链;财务数据再退到 akshare 的美股接口(feat: add akshare for US stock data fallback)。美股 ticker 池也只采前 200(市值排序),而不是全 universe,避免触发限流。
坑三:港股 ticker 池根本采不到。 早期用 eastmoney 拉港股,ECS 海外出口被挡,返回全空。修法是直接硬编码一份港股蓝筹列表作为 primary,再退到市值前 200 兜底。
坑四:info_rating 永远是 C。 verify_financial_data 原本只看 news_count 这一项给信息完整度打分。HK/US 在 yfinance 限流下新闻经常为 0,于是 info_rating 永远 C,详情页看起来"数据不全"。改成复合判断:
# app/services/financial_verification.py(提交 5226a21)
def _composite_info_rating(news_count: int, fills: dict[str, Any]) -> str:
fields = [fills.get("price"), fills.get("market_cap"),
fills.get("industry"), fills.get("pe"),
fills.get("pb"), fills.get("roe")]
core_ok = sum(1 for f in fields if f is not None) >= 4
if news_count >= 3 and core_ok: return "A"
if core_ok or news_count >= 1: return "B"
return "C"
US/HK 在有 yfinance 财务 + 0 新闻的情况下,现在能拿 B 而不是 C。info_rating 是详情页"数据完整度"标签的源头,这一改之后 AAPL、NVDA 这些标的详情页就能正常显示完整度,而不是永远 C。
整个数据层有一个贯穿的设计原则:单一 ticker 维度,所有源降级到统一 schema,由 BaseCollector 接口约束。这样上层(圆桌、红黑榜)不用知道当前数据来自哪一源,限流、源宕机、字段缺失都被压在了 collector 内部。
四、后端 LLM 底层调用
这是整个 Agent 最底层的一块。我没有把 LLM 调用写死在某一个 SDK 上,而是做了一套「多 Provider + 重试 + 超时 + 降级」的调用链。核心代码在 app/services/master_agent.py。
1. 凭证统一注入
每次调用前先把 settings 里的凭证刷进环境变量,避免不同入口(API 进程 / Worker 进程 / 脚本)凭证不一致:
def _ensure_llm_credentials() -> None:
settings = get_settings()
if settings.anthropic_api_key and not os.getenv("ANTHROPIC_API_KEY"):
os.environ["ANTHROPIC_API_KEY"] = settings.anthropic_api_key
if settings.minimax_api_key and not os.getenv("MINIMAX_API_KEY"):
os.environ["MINIMAX_API_KEY"] = settings.minimax_api_key
# ... base_url / model 同理
2. 双 Provider 优先级链
根据 llm_provider 配置决定优先用谁,另一个做兜底。比如默认 MiniMax 优先、Anthropic 兜底:
provider = str(settings.llm_provider or "minimax").strip().lower()
provider_order = ["anthropic", "minimax"] if provider == "anthropic" else ["minimax", "anthropic"]
for current in provider_order:
if current == "minimax" and os.getenv("MINIMAX_API_KEY") and OpenAI is not None:
client = OpenAI(api_key=os.getenv("MINIMAX_API_KEY"),
base_url=os.getenv("MINIMAX_BASE_URL", "https://api.minimaxi.com/v1"))
completion = client.chat.completions.create(
model=os.getenv("LLM_QUICK_THINK_MODEL", "MiniMax-Text-01"),
temperature=0.2,
messages=[{"role": "system", "content": system},
{"role": "user", "content": json.dumps(prompt, ensure_ascii=False)}],
)
content = (completion.choices[0].message.content or "").strip()
if content:
break
if current == "anthropic" and os.getenv("ANTHROPIC_API_KEY") and Anthropic is not None:
client = Anthropic(api_key=os.getenv("ANTHROPIC_API_KEY"))
message = client.messages.create(
model=os.getenv("ANTHROPIC_MODEL", "claude-3-5-sonnet-latest"),
max_tokens=1400, system=system,
messages=[{"role": "user", "content": json.dumps(prompt, ensure_ascii=False)}],
)
content = "\n".join(getattr(b, "text", "") for b in message.content).strip()
if content:
break
注意 MiniMax 走的是 OpenAI SDK(OpenAI-compatible),Anthropic 走原生 SDK。两者 message 结构不同,所以分别构造。try: from openai import AsyncOpenAI 这种「软导入」是为了在某一个 SDK 没装时不让整个模块崩掉。
3. 重试 + 超时(首 token 前才重试)
圆桌场景下每个大师的回复是要流式推给前端的,所以重试策略很讲究。_attempt_llm_stream_with_retry 的核心规则是:
- 首 token 之前抛异常或超时 → 重试(最多
ROUNDTABLE_MASTER_RETRIES次,默认 3),退避min(2**attempt, 4)秒。 - 首 token 之后再超时或异常 → 不重试,用已经流出去的 partial 文本收尾。因为部分文字已经显示在用户屏幕上了,重试会重复输出。
- 超时本身不重试(产品决策),直接 fall through 到下一个 Provider。
async def _consume() -> None:
async for delta in gen_fn(master, model_prompt):
started = True # 第一个 token 一来就置位
accumulated.append(delta)
await on_delta(delta)
try:
await asyncio.wait_for(_consume(), timeout=timeout)
except asyncio.TimeoutError:
if started:
return "".join(accumulated).strip() or None # partial 收尾
if attempt < retries:
await asyncio.sleep(min(2 ** attempt, 4)); continue
return None
4. 三级降级
如果 MiniMax 和 Anthropic 都失败(或没配凭证),最后还有一层本地启发式兜底 analyze_stock_with_masters:用 hash + 关键词正则(ai|云|半导体 加成长分,股息|银行|煤炭 加价值分)生成一个结构化分数。这样即使 LLM 全挂,接口依然返回结构化数据,只是 source 字段标记成 heuristic。
return {
"stockName": q[:16] or "未命名标的",
"bullScore": bull, "bearScore": bear,
"masters": masters,
"source": "heuristic", # 或 "model"
}
5. 调用结果熔断监控
record_master_call_outcome 维护一个 30 分钟滑动窗口,按大师维度统计失败率。样本数 ≥6 且失败率 ≥30% 时打 ERROR 日志告警。Fail-open,只记日志不阻断业务:
_OUTCOME_WINDOW_SECONDS = 30 * 60
_OUTCOME_MIN_SAMPLES = 6
_OUTCOME_FAILURE_RATE_THRESHOLD = 0.30
五、Skill 如何加载
「大师」的人格不是写死在 prompt 里的,而是用 Claude Code 的 Skill 机制承载的。每个大师对应一个 SKILL.md,社区仓库里有人维护好了(Panmax/buffett-skill、Panmax/duanyongping-skill 等),我直接 sync 进来用。
1. 目录结构
app/master_skills/
catalog.py # 大师元数据登记表
__init__.py
library/ # commit 进仓的 skill 原文兜底
buffett/SKILL.md
duan/SKILL.md
dalio/SKILL.md
...(共 11 位)
.claude/skills/ # agent-sdk 实际读取的位置(gitignored,sync 后才有)
2. Catalog 登记表
catalog.py 用 dataclass 描述每个 skill 的 key、中文名、头像、repo 来源和别名。它实际只有这五个字段,没有 version / author / tags,版本和署名由 SKILL.md 的 frontmatter 承担(实际上 frontmatter 也只有 name + description,下面会看到):
@dataclass(frozen=True)
class MasterSkillSpec:
key: str
name: str
avatar: str
repo: str | None = None
aliases: tuple[str, ...] = ()
11 位大师的元数据都是同一种结构,但来源不全是社区仓。李录(li)的 repo=None —— 是我手工写的,没用现成 skill 包;段永平、巴菲特、达利欧、索罗斯、彼得·林奇、马克斯、博格、蒂尔、格雷厄姆这 9 位都标了 Panmax/*-skill,芒格标的是 alchaincyf/munger-skill:
MASTER_SKILL_SPECS = {
"duan": MasterSkillSpec(key="duan", name="段永平",
avatar="/assets/avatars/duan.png", repo="Panmax/duanyongping-skill",
aliases=("duanyongping",)),
"buffett": MasterSkillSpec(key="buffett", name="巴菲特",
avatar="/assets/avatars/buffett.png", repo="Panmax/buffett-skill"),
"munger": MasterSkillSpec(key="munger", name="芒格",
avatar="/assets/avatars/munger.png", repo="alchaincyf/munger-skill"),
"li": MasterSkillSpec(key="li", name="李录",
avatar="/assets/avatars/li.png", repo=None, aliases=("lilu", "li_lu")),
# ... dalio / soros / lynch / marks / bogle / thiel / graham
}
MASTER_ALIAS_TO_KEY 把所有别名(含 key 本身)展平成一张反查表,用户输入 peter-lynch / duanyongping / grahamben 都能归一到同一个 key。这个反查表只在程序世界里有效 —— 别名不写进 SKILL.md frontmatter,LLM 不会"看到"别名。
3. SKILL.md 真实长什么样
去掉我美化过的版本,一个 SKILL.md 真实的 frontmatter 长这样(巴菲特为例):
---
name: buffett-perspective
description: |
巴菲特的投资思维框架与商业智慧。基于伯克希尔年度致股东信(1965-至今)...
用途:作为思维顾问,用巴菲特的视角分析投资决策、商业判断、人生选择、风险评估等问题。
当用户提到「用巴菲特的视角」「巴菲特会怎么看」时使用。
不要在用户只是问「帮我看看代码」「这个bug怎么修」等纯技术问题时触发。
---
description 是一段很长的 YAML block scalar,里面写了触发条件(“用户提到 X 时使用”)和反触发条件(“不要在 Y 时触发”)。这部分是真正决定 LLM 在 agent-sdk 路径下何时激活这个 skill 的关键。
4. 文本加载(双位置,没有 lru_cache)
.claude/skills/ 是 gitignored 的,只有跑过 sync 脚本才存在;为了让没 sync 的环境也能跑,load_master_skill_text 会先找 sync 位置,再 fallback 到 commit 进仓的 library/ 副本:
def load_master_skill_text(master: str | None) -> str:
spec = get_master_skill_spec(master)
if spec is None:
return ""
candidates = [
spec.local_skill_path, # .claude/skills/<key>/SKILL.md
_SKILL_LIBRARY_DIR / spec.key / "SKILL.md", # app/master_skills/library/<key>/SKILL.md
]
for path in candidates:
if path.exists():
return path.read_text(encoding="utf-8").strip()
return ""
诚实的几点:
- 这里没有
@lru_cache。每次调用都从盘读,代价不高(SKILL.md 都 < 10KB),但也没有失效机制 —— 改完 SKILL.md 要重启 worker 进程才生效。 - skill 不带工具(
tools:字段为空)。它只是人格 prompt + 触发条件;真正的工具调用走 round-table pipeline,skill 不挂工具。 - 11 个 skill 是「同一套人格描述 + 不同 frontmatter 触发条件」,不是「11 个独立工作流」。catalog 里的
key / name / avatar / aliases是给前端路由和程序用的,SKILL.md 里的description触发条件是给 LLM 用的,两套数据没有自动一致性校验 —— 改了 catalog 没改 SKILL.md,路由没问题但人格解释漂移;反过来同理。
5. System Prompt 组装
build_master_system_prompt 把 skill 原文拼成最终 system prompt,并加上投资研究场景的额外约束(中文、不编造数据、信息不足要说不确定):
def build_master_system_prompt(master: str | None) -> str:
key = canonical_master_key(master)
if key == ROUNDTABLE_KEY: # "all" → 圆桌助理
return build_roundtable_prompt()
spec = MASTER_SKILL_SPECS[key]
skill_text = load_master_skill_text(key)
if not skill_text:
return f"你现在扮演{spec.name}风格的投资研究助手。..." # skill 丢失时的兜底
return (
f"你现在扮演{spec.name}风格的投资研究助手。\n"
"下面是必须遵循的 skill 原文,请把它作为你的主要行为和风格约束。\n\n"
f"{skill_text}\n\n"
"额外约束:\n1. 使用中文回答。\n2. 优先讨论商业模式、估值、风险和仓位。\n"
"3. 不要编造数据,不要声称拥有内幕信息。\n4. 信息不足要明确说不确定。"
)
6. 接入 Claude Agent SDK
更深一层的玩法在 master_chat_service.py:单大师对话会优先走 claude_agent_sdk 的 query(),把 skill 名直接作为 skills 参数传进去,让 Agent SDK 自己加载 skill 并跑工具循环(max_turns=6)。SDK 跑挂了再 fallback 到裸 Anthropic messages.create:
async def _claude_agent_reply(master, prompt, resume_session_id):
skills = None if master == "all" else [master]
options_kwargs = {
"cwd": str(PROJECT_ROOT),
"model": settings.anthropic_model,
"setting_sources": ["project"],
"max_turns": 6,
}
if skills:
options_kwargs["skills"] = skills # 让 SDK 加载这个大师的 skill
if resume_session_id:
options_kwargs["resume"] = resume_session_id
async for message in query(prompt=prompt, options=ClaudeAgentOptions(**options_kwargs)):
...
也就是说,skill 有两条加载路径:一条是自己读 SKILL.md 文本拼进 system prompt(轻量、可控、双 Provider 通用),另一条是交给 Claude Agent SDK 原生加载(重,但能用上 SDK 的工具循环和会话恢复)。
六、Workflow 如何实现
圆桌是整个产品的核心 Workflow:N 个大师并行发言,每个大师的发言又要流式推给前端,还要保证一个大师崩了不影响其他人。实现在 app/services/roundtable_pipeline.py::run_roundtable_sync。
1. 编排总览
async def run_roundtable_sync(task_id, query, masters, ...):
selected_masters = _normalize_masters(masters) # 10 位大师
shared_analysis = await _load_roundtable_fundamentals(...) # 先采集共享基本面
concurrency = int(os.getenv("ROUNDTABLE_MASTER_CONCURRENCY", str(MASTER_CONCURRENCY))) # 默认 3
semaphore = asyncio.Semaphore(concurrency)
master_tasks = []
async def _run_one(master, index):
async with semaphore: # 信号量限并发
await _publish(task_id, {"type": "progress",
"stage": f"master:{master}",
"message": f"{_display_name(master)} 正在发言"}) # 诚实播报
base_analysis = analyze_stock_with_masters(prompt_for_masters, [master])
reply, analysis, llm_ok = await roundtable_master_reply_stream(
master, prompt_for_masters, base_analysis, on_delta=_on_delta)
return master, reply, analysis, llm_ok
for index, master in enumerate(selected_masters):
task = asyncio.create_task(_run_one(master, index)) # 全部 create_task
master_tasks.append(task)
for finished_task in asyncio.as_completed(master_tasks): # 谁先好谁先推
try:
master, raw_reply, analysis, llm_ok = await finished_task
except Exception:
logger.exception(...); continue # 单大师崩溃不影响其他人
...
await _publish(task_id, {"type": "report", "section": f"master_{master}", ...})
几个关键设计:
asyncio.Semaphore(3)限并发:10 个大师不会同时打满 LLM,最多 3 个在飞。asyncio.as_completed:谁先回来谁先推report帧,前端按完成顺序渲染,不用等最慢的那个。- 诚实播报:「正在发言」是在拿到 semaphore 之后才发的,而不是
create_task时就发。否则 10 个大师一瞬间全部「正在发言」,但其实 7 个在排队——体验上是骗人的。 - 单大师隔离:
try/except包住每个finished_task,一个大师的 LLM 调用挂了,循环继续,只是这一位用兜底模板。
2. 流式 delta 推送
每个大师发言时,_on_delta 回调把 LLM 产出的每个 token chunk 包成 report_delta 帧推给前端。第一个 chunk 先发 report_start 开占位消息,后续 chunk 不断 append:
async def _on_delta(delta):
if not stream_started[0]:
stream_started[0] = True
await _publish(task_id, {"type": "report_start",
"section": f"master_{master}", "masterKey": master, ...})
await _publish_live(task_id, {"type": "report_delta",
"section": f"master_{master}", "masterKey": master, "delta": delta})
帧通过 Redis pubsub 推送(pubsub.publish),晚加入的 WS 客户端会先 replay history list 再 live subscribe,所以刷新页面也能看到完整发言。
3. 共享基本面,避免重复采集
圆桌开始前先跑一次 _load_roundtable_fundamentals,把这只股票的行情/估值/财务/资金流/概念板块采集齐,作为 shared_analysis 注入给每一位大师。否则 10 个大师各自采集一遍,既慢又容易被行情源限流。采集完还会做一次 verify_financial_data 交叉验证(比如市值对不对得上),算一个 info_rating。
4. 数据优先模式
master_chat_service.chat_with_master 里还有个细节:用户提问里如果带 A 股 6 位代码(正则 \b([036]\d{5})\b 命中),会先并行触发 fundamentals_collector.collect_all 采集基本面,再用 four_masters_framework_analysis(段永平看生意本质、巴菲特看护城河、芒格看逆向风险、李录看文明趋势)做结构化分析,把结果作为 context 注入 system prompt。持仓上下文每次提问都注入但落库,避免多轮会话累积重复的持仓 dump。
5. 价格的真相:LLM 只给两个粗标量,剩下全是规则算的
红黑榜详情页有一个 UI 槽位:「加仓区间 / 目标价 / 止损位」。这个槽位必须填,不允许留空。要回答"价格是怎么出来的",必须先讲清楚这条链上每一步是谁做的——否则很容易把"LLM 推理出的价"和"程序算出来的价"混在一起。
为什么要避免 LLM 直接给价
我一开始打算让 LLM 自己在分析里出价,实测下来有两个问题:
- 大模型对一只股票直接给"目标价 156.30 元"这种数字非常不靠谱 —— 它没有实时报价,能给出的只有"区间",但用户拿到的是一个虚假精度的数字;
- 不同 prompt 下,模型要么爱给(无理由地编数字)、要么偏保守(一个价都不出),分布不可控。
价格的 7 步因果链
每条 ticker 的最终 trade_setup,实际走的是这 7 步(都在 app/services/research_leaderboard.py):
[1] LLM 在 candidate screening 阶段对每只票给出两个 0-100 整数
bull_score = "多头逻辑强度"
bear_score = "空头风险强度"
↓
[2] 程序算分差
spread = bull_score - bear_score
↓
[3] 阈值映射成 verdict(程序,不是 LLM)
spread ≥ 25 → Buy
spread ≥ 10 → Accumulate
spread ≤ -20 → Reduce
其余 → Hold
↓
[4] 程序读"当前价"
price = _coerce_float(metrics.get("currentPrice") or item.get("price"))
—— 来自 §三 那一节的数据采集层(astock.py / hk_us.py),不是 LLM
↓
[5] 程序按 verdict 查乘数表算 4 个价位(程序,不是 LLM)
rating == "Buy" → addLow=price*0.97, addHigh=price*1.01, target=price*1.12, stop=price*0.92
rating == "Accumulate" → 0.95 / 0.99 / 1.08 / 0.91
rating == "Reduce" → 0.90 / 0.94 / 0.95 / 0.88
rating == "Hold" → 0.94 / 0.98 / 1.05 / 0.90
↓
[6] 程序组装 trade_setup dict
{"action": rating, "addLow": ..., "addHigh": ..., "targetPrice": ..., "stopLoss": ...}
↓
[7] 程序写入 ResearchLeaderboardDetailTask 的 payload
详情页直接读 targetPrice / stopLoss / addLow / addHigh
关键代码(_derive_trade_setup):
rating = "Buy" if spread >= 25 else ("Accumulate" if spread >= 10
else ("Reduce" if spread <= -20 else "Hold"))
price = _coerce_float(metrics.get("currentPrice") or item.get("price"))
add_low = add_high = target_price = stop_loss = None
if price and price > 0:
if rating == "Buy":
add_low, add_high = price*0.97, price*1.01
target_price, stop_loss = price*1.12, price*0.92
elif rating == "Accumulate":
add_low, add_high = price*0.95, price*0.99
target_price, stop_loss = price*1.08, price*0.91
elif rating == "Reduce":
add_low, add_high = price*0.90, price*0.94
target_price, stop_loss = price*0.95, price*0.88
else: # Hold
add_low, add_high = price*0.94, price*0.98
target_price, stop_loss = price*1.05, price*0.90
trade_setup = {"action": rating,
"addLow": add_low, "addHigh": add_high,
"targetPrice": target_price, "stopLoss": stop_loss}
7 步里 LLM 只参与了第 [1] 步,并且只是给两个粗标量。[3] [5] 是规则做的,[2] [4] [6] [7] 是程序做的。LLM 从来不出数字,详情页却永远有数字,这就是「强制填充」的真正机制。
LLM 整体挂了怎么办:第 [1] 步如果 LLM 没出 bull/bear(或 model path 整体 rate-limited),代码 fallback 到 _heuristic_core(_heuristic_core 是 fallback 路径,trigger 条件:leaderboard got zero model hits / insufficient refined finalists / leaderboard model path rate-limited)—— 用 bull = 48 + cap_bonus + style_bias + max(momentum,0) 这种纯规则算分,然后 [2]–[7] 完全一样。fallback 不只是"返回同样的 payload",是连 verdict 阈值映射和价位乘数表都共用——这样 LLM 在线时和 LLM 离线时的产品形状是一致的,不会因为模型抖动导致 UI 行为变化。
为什么这条链上还有"两套价":Vendored 在 app/agents/tradingagents/agents/schemas.py 里的 TraderProposal 也定义了 entry_price / stop_loss,但那个 Pydantic 模型是给 TradingAgents 长文分析报告用的(提交报告里的"Entry Price: 156.30"那段文字),没有落到详情页的 trade_setup。也就是说,同样叫"价格",两套机制并存:长文里是 LLM 自由给的数字(不展示给用户,是分析师报告原文),详情页里是规则按当前价 × 系数算的(展示给用户)。把这两个分开,用户看到的永远是确定性的数字,长文里的"自由数字"只是模型推理痕迹。
master_agent.py 里另一段文案 "如果没有明确买入价和止损纪律,情绪波动会放大回撤" 看起来像在"算价",其实它是反向 prompt —— 当用户没给价格时,提示模型主动反问用户拿输入。和"算价"是两件事,别混。
乘数表本身是有产品的取舍在里面
Buy 1.12 / Hold 1.05 / Reduce 0.95 —— Hold 也给 +5% 的乐观目标价,是为了让「等待」也有可视化锚点;Buy 的 12% / Reduce 的 -5% 是不对称的,因为底部敢重锤比顶部敢跳水的概率更低。这个表改过一次(从 1.15 收到 1.12),改了之后 A 股用户反馈"目标价比现价高一点更能让人拿住"。这种数值不是 LLM 决定的,是产品决定。乘数表、阈值(25 / 10 / -20)、"Hold 也给乐观价"这三条是产品的边界条件,和模型完全解耦。
6. verdict 为什么稳定:准过的不是 LLM,是架构
把上面 7 步再串起来看:"verdict 值得参考"这件事在代码里不是由 LLM 的可信度承担的,是由架构承担的。
- 唯一的 LLM 输入是 0-100 的两个整数(bull / bear)——LLM 不是在做"该买该卖"的方向判断,是在给"多头逻辑强弱 / 空头风险强弱"两个标量打分。让模型给一个粗标量比让它直接给方向或价格稳定得多:标量是连续的,方向和价格是离散的,连续空间对 LLM 的抖动更鲁棒。
- 方向判断是阈值映射(
spread ≥ 25Buy /≥ 10Accumulate /≤ -20Reduce / 否则 Hold)——小幅分差波动不会跳档。LLM 偶尔给spread=23不会变成 Buy,给spread=12不会跳出 Accumulate。这是抗 LLM 抖动的第一道闸。 - 价位是乘数表(Buy 1.12 / Hold 1.05 / Reduce 0.95)——Buy 和 Hold 的差距不在于"LLM 觉得能涨多少",在于"产品认为 Buy 比 Hold 更值得下手"。价位的差异是相对的、固定的,不是 LLM 推理出来的具体数字。这是抗 LLM 抖动的第二道闸。
- fallback 链——
_heuristic_core在 LLM 整体挂了时按规则给 bull/bear(_heuristic_core不是主路径,trigger 条件见 §六.5),下游的 [3] [5] 完全共用,所以 LLM 在线时和离线时的产品形状一致。这是抗 LLM 完全不可用的第三道闸。 - 大师自己的 verdict 也是同一套加权(
score = bull*0.65 + (100-bear)*0.35 + style_offsets[key][style])——大师 SKILL 的人格只在reason字段的模板里出现,没有真正影响 verdict 数字。人格是叙事层,verdict 是数据层,这是抗"人格发挥过度"的第四道闸。
把这四道闸加起来:LLM 只需要给"多头 70 / 空头 30"这种粗标量就够了,剩下的方向、价位、fallback、风格偏移全部由程序规则决定。这意味着 verdict 的稳定性不依赖"选哪个 LLM"——provider 可以是 MiniMax,可以是 Anthropic,也可以是 GPT-4o,这三者之间的差别不影响 verdict 的最终形状,只影响 bull/bear 这两个标量本身的质量。
这件事的产品意义:用户拿到的红黑榜分数、verdict、价位三件套,是经过四道规则过滤的 LLM 推理,不是裸 LLM 输出。这才是"详情页可以稳定显示买价"的真相 —— 不是因为 LLM 准,是因为架构把 LLM 的不确定性挡在了外面。
七、定时器如何实现
红黑榜需要一个定时器:每小时把全市场强势/弱势股的详情(大师发言 + 分析师报告)预热好,这样用户点开详情页是秒开的。难点是「定时器跑在哪、多容器不重复触发、失败可重试」。
1. 跑在 api 进程的 lifespan 里
没有引入额外的 cron 或 scheduler 服务,而是把定时循环挂在 FastAPI 的 lifespan 里,作为后台 asyncio.Task:
@asynccontextmanager
async def lifespan(app: FastAPI):
await init_db()
warm_graph() # 预热 TradingAgents graph
# 定时预热循环
interval = max(60, int(settings.redblack_hourly_interval_seconds or 3600))
if interval > 0:
import asyncio as _asyncio
from scripts.redblack_hourly_prewarm import hourly_loop
stop_event = _asyncio.Event()
hourly_task = _asyncio.create_task(hourly_loop(stop_event))
logger.info("init: hourly redblack prewarm loop registered (interval=%ds)", interval)
try:
yield
finally:
if stop_event: stop_event.set() # 优雅停止
if hourly_task: hourly_task.cancel()
hourly_loop 是一个经典的「事件 + 超时」循环,开机先跑一次(today 的 incomplete 详情立刻补),之后每隔 interval 秒跑一次:
async def hourly_loop(stop_event: asyncio.Event) -> None:
interval = max(60, int(settings.redblack_hourly_interval_seconds or 3600))
await enqueue_hourwarm_with_lock(initial=True) # 开机立即跑一次
while not stop_event.is_set():
try:
await asyncio.wait_for(stop_event.wait(), timeout=interval)
return # stop_event 被 set → 退出
except asyncio.TimeoutError:
pass
await enqueue_hourly_prewarm()
用 asyncio.wait_for(stop_event.wait(), timeout=interval) 而不是 asyncio.sleep(interval),是为了收到停止信号能立刻退出,不用等满一个 interval。
2. Redis 分布式锁,多容器不重复
生产上 api 可能起多副本,每个副本都跑 hourly_loop 就会重复触发。用 Redis SET NX EX 按 market 加锁,TTL 55 分钟(REDBLACK_HOURLY_LOCK_TTL,略小于 interval):
_HOURLY_LOCK_KEY = "redblack:hourly_prewarm:tick:{market}"
for market in SUPPORTED_BOARD_MARKETS:
lock_key = _HOURLY_LOCK_KEY.format(market=market)
if not force:
acquired = r.set(lock_key, "1", ex=_HOURLY_LOCK_TTL, nx=True)
if not acquired:
summary[market] = {"status": "locked"} # 别的副本在跑,跳过
continue
...
3. api 只 enqueue,worker 干活
定时器本身很轻:它只扫一遍 board,把 incomplete 的 ticker 丢进 per-market 的 RQ 队列(redblack_detail_{market}),真正跑 LLM 的是 worker 进程。这样 api 进程不会被长耗时任务拖住:
queue_name = f"redblack_detail_{market}"
queue = Queue(queue_name, connection=r)
for ticker in queued:
queue.enqueue(build_redblack_detail, ticker, today, market,
job_timeout=min(180, max(120, int(settings.redblack_detail_timeout_seconds or 180))),
result_ttl=86400, failure_ttl=86400)
4. 任务表 + 完整性闸门
每个 ticker 的详情构建对应 ResearchLeaderboardDetailTask 表里一行,状态机是 queued → running → done/error。build_redblack_detail 里有一道严格的完整性闸门:大师发言和分析师报告都齐了才算 done,否则标 error 等下一轮 hourly 重试。还有一个「building lock」(redblack_cache.mark_building),防止两个 worker 同时建同一个 ticker——抢不到锁的会 requeue 1 分钟后再试。
if not await redblack_cache.mark_building(ticker, board_date, market=normalized, ttl=ttl):
# 另一个 worker 正在建,requeue 1 分钟后
row.status = "queued"; row.next_retry_at = _utcnow() + timedelta(minutes=1)
return {"taskId": task.id, "ticker": ticker, "status": "locked"}
5. 兜底:外部 cron 也能触发
如果不想让 api 进程承担定时职责,也可以外部 cron 直接跑脚本,入口是幂等的(锁还在就不重复):
python -m scripts.redblack_hourly_prewarm # 普通触发
FORCE=1 python -m scripts.redblack_hourly_prewarm # 强制全量重建(部署后用)
八、红黑榜链路串起来
把上面几块串起来看红黑榜的完整链路:
- 定时器(lifespan
hourly_loop,Redis 锁防重)每小时触发enqueue_hourly_prewarm。 - api 进程扫 board,把 incomplete ticker enqueue 到
redblack_detail_{market}RQ 队列。 - worker 进程拉队列,跑
build_redblack_detail:加 building 锁 → 调get_red_black_item_detail(内部走上面的多 Provider LLM 链路,让大师们发言)→ 完整性闸门校验 → 写 Redis 详情缓存 + 持久化到 job + 更新 task 状态。 - 用户打开红黑榜详情页时,接口直接读 Redis 缓存,秒开。
定时器负责「让它一直是热的」,Workflow 负责「一次怎么把十位大师的发言跑齐」,LLM 底层调用负责「单次调用怎么稳」,Skill 负责「大师的人格从哪来」——四层各司其职。
九、几点反思
- Skill 机制值得借。 把人格/领域知识从 prompt 里抽出来变成可版本化、可替换的
SKILL.md,比在代码里堆 system prompt 字符串干净太多。社区已经有人维护好各大师/各行业的 skill,直接 sync 就能用。 - 多 Provider 不是「配置一个 base_url」就完事。 OpenAI-compatible 和 Anthropic 原生 SDK 的 message 结构、流式协议、重试语义都不一样,需要一层适配。尤其流式场景的「首 token 后不重试」是个反直觉但必须做的决策。
- 定时器不一定非要 cron/celery-beat。 量不大的时候,lifespan 里一个
asyncio.create_task+ RedisSET NX分布式锁就够用,少引入一个组件。但前提是任务本身要幂等、要能被外部 cron 兜底。 - 诚实的进度播报。 「正在发言」要在真正拿到并发槽位之后再发,而不是任务创建时就发。这种细节决定了 Agent 产品的「可信感」。
- 降级要贯穿到底。 LLM 挂 → 启发式兜底;skill 文件丢 → 默认 prompt 兜底;一个大师崩 → 其他大师继续;详情没建完 → 下一轮重试。每一层都有 fallback,产品才不会因为某个外部依赖抖动而整体不可用。
十、为什么不是「裸用 ChatGPT 类聊天」
这里说的"裸用",指的是直接把用户问题丢给一个通用大模型对话界面、让它自由回答。这套产品从来没考虑过这种形态,原因不是模型能力不够,而是产品形态对不上:
-
多 persona 并发。 圆桌一次要让 10 位大师同时发言。ChatGPT 类聊天界面一次只能跑一段对话;要做多 persona,必须在客户端或者后端自己造一层编排,而这就回到 §五 的 Workflow 那一套了 ——"裸用"在这里不是省事,反而是放弃了已有的能力。
-
Skill 注入没有标准协议。 我想要的不是"问巴菲特怎么看 600519",而是"以巴菲特视角,并且以巴菲特视角触发条件为前提,并且只在用户问投资判断时激活,并且可以挂载 skill 自带的资料"——这一切在 ChatGPT 类聊天界面里都需要靠"提示词模板 + 用户自律",要么每次重写要么塞到 system prompt 里硬编码;skill 文件 + frontmatter 触发条件这种声明式规范在这里没有落地位置。Claude Code 的 skill 机制(
description里写触发 / 反触发条件,agent-sdk 自己加载)是我目前找到的唯一一套"模型外部能精确控哪些 skill 被加载"的协议。 -
数据接入与价格后处理。 §3.5 的数据采集和 §5.5 的"按当前价 × 系数算价位"都是模型外部的工作。裸聊天界面没法先并行采集 7 个数据源(行情/估值/财务/资金流/概念板块/新闻/资金流),也没法在 LLM 给出"BUY"之后做确定性的价位计算并落到结构化字段。LLM 给的是方向,规则给的是数字,二者必须分开。
-
分布式预热 + 详情缓存。 红黑榜详情页是提前一小时算好的(§六 的 hourly prewarm),用户点开是秒开的。这个生命周期是模型之外的(RQ 队列、Redis 详情缓存、worker 并发),也是裸聊天界面完全管不到的。
-
多轮状态隔离。 用户问"以巴菲特视角持有比亚迪和茅台哪个更好",需要先调出持仓、调出财务、调出两段 LLM 分析、调出 trade_setup,再以巴菲特视角综合。ChatGPT 类聊天的多轮是单线上下文,做不了这种"先做 A → 再做 B → 再用 C 视角看 A+B"。§五 的并发圆桌、§三 的数据优先模式、§五.5 的价格后处理,本质上都是在多 persona / 多数据源 / 多阶段之间保持状态隔离,这是工作流编排的活,不是聊天的活。
总结一句:模型是原材料,skill + workflow 才是产品。Provider 可以是 MiniMax,可以是 Anthropic,也可以是 GPT-4o,这三者之间的差别是能力 / 成本 / 中文 / 区域可达,但都不能替代 skill 注入、workflow 编排、数据采集、价格后处理、分布式预热这五件事。把这五件事做掉之后,"用什么模型"反而是次要问题。
在线体验:https://agent.xiaofenglei.cn/,红黑榜在 /markets 页面。文中代码均来自实际仓库 SoundFinancialPlanning,有删减以突出主干。