存钱小帮手 Agent 的设计与实现

存钱小帮手 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.6react@19.2.4@supabase/ssropenai@^6zod@^4。测试用 Vitest + Playwright。packages/* 必须零 React/Next 依赖,靠 pnpm typecheck 在 CI 里卡住边界。

小程序端(SoundFinancialPlanning)

TypeScript + LESS,自定义 tabBar,页面有:holdings(我的持仓)、index(稳健理财)、chat(大师圆桌对话)、master-portfolio(大师组合)、analysisredblack-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-skillPanmax/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_sdkquery(),把 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 ≥ 25 Buy / ≥ 10 Accumulate / ≤ -20 Reduce / 否则 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/errorbuild_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  # 强制全量重建(部署后用)

八、红黑榜链路串起来

把上面几块串起来看红黑榜的完整链路:

  1. 定时器(lifespan hourly_loop,Redis 锁防重)每小时触发 enqueue_hourly_prewarm
  2. api 进程扫 board,把 incomplete ticker enqueueredblack_detail_{market} RQ 队列。
  3. worker 进程拉队列,跑 build_redblack_detail:加 building 锁 → 调 get_red_black_item_detail(内部走上面的多 Provider LLM 链路,让大师们发言)→ 完整性闸门校验 → 写 Redis 详情缓存 + 持久化到 job + 更新 task 状态。
  4. 用户打开红黑榜详情页时,接口直接读 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 + Redis SET NX 分布式锁就够用,少引入一个组件。但前提是任务本身要幂等、要能被外部 cron 兜底。
  • 诚实的进度播报。 「正在发言」要在真正拿到并发槽位之后再发,而不是任务创建时就发。这种细节决定了 Agent 产品的「可信感」。
  • 降级要贯穿到底。 LLM 挂 → 启发式兜底;skill 文件丢 → 默认 prompt 兜底;一个大师崩 → 其他大师继续;详情没建完 → 下一轮重试。每一层都有 fallback,产品才不会因为某个外部依赖抖动而整体不可用。

十、为什么不是「裸用 ChatGPT 类聊天」

这里说的"裸用",指的是直接把用户问题丢给一个通用大模型对话界面、让它自由回答。这套产品从来没考虑过这种形态,原因不是模型能力不够,而是产品形态对不上:

  1. 多 persona 并发。 圆桌一次要让 10 位大师同时发言。ChatGPT 类聊天界面一次只能跑一段对话;要做多 persona,必须在客户端或者后端自己造一层编排,而这就回到 §五 的 Workflow 那一套了 ——"裸用"在这里不是省事,反而是放弃了已有的能力。

  2. Skill 注入没有标准协议。 我想要的不是"问巴菲特怎么看 600519",而是"以巴菲特视角,并且以巴菲特视角触发条件为前提,并且只在用户问投资判断时激活,并且可以挂载 skill 自带的资料"——这一切在 ChatGPT 类聊天界面里都需要靠"提示词模板 + 用户自律",要么每次重写要么塞到 system prompt 里硬编码;skill 文件 + frontmatter 触发条件这种声明式规范在这里没有落地位置。Claude Code 的 skill 机制(description 里写触发 / 反触发条件,agent-sdk 自己加载)是我目前找到的唯一一套"模型外部能精确控哪些 skill 被加载"的协议。

  3. 数据接入与价格后处理。 §3.5 的数据采集和 §5.5 的"按当前价 × 系数算价位"都是模型外部的工作。裸聊天界面没法先并行采集 7 个数据源(行情/估值/财务/资金流/概念板块/新闻/资金流),也没法在 LLM 给出"BUY"之后做确定性的价位计算并落到结构化字段。LLM 给的是方向,规则给的是数字,二者必须分开。

  4. 分布式预热 + 详情缓存。 红黑榜详情页是提前一小时算好的(§六 的 hourly prewarm),用户点开是秒开的。这个生命周期是模型之外的(RQ 队列、Redis 详情缓存、worker 并发),也是裸聊天界面完全管不到的。

  5. 多轮状态隔离。 用户问"以巴菲特视角持有比亚迪和茅台哪个更好",需要先调出持仓、调出财务、调出两段 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,有删减以突出主干。


文章作者: 小风雷
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 小风雷 !
评论
 本篇
存钱小帮手 Agent 的设计与实现 存钱小帮手 Agent 的设计与实现
从一个「稳健理财」小程序出发,拆解大师圆桌 Agent 的前后端技术栈、LLM 多 Provider 底层调用、Skill 加载机制、圆桌 Workflow 编排,以及基于 lifespan + Redis 分布式锁的定时预热器。
2026-07-21
下一篇 
ClaudeCode 的执行机制:ReAct 与 Plan-and-Execute 在 TypeScript 源码中的实现 ClaudeCode 的执行机制:ReAct 与 Plan-and-Execute 在 TypeScript 源码中的实现
基于 ClaudeCode TypeScript 源码,拆解它如何用 query 主循环实现 ReAct,如何通过 plan mode、plan file、approval flow 叠加出 Plan-and-Execute,以及这两套机制分别在什么条件下触发、判断依据是什么、如何持续传递给大模型。
2026-04-06
  目录