起因
最近在做一个叫 vibe-writer 的 AI 写作系统。
输入一个主题,系统会自动完成:
- 搜索资料
- 规划文章结构
- 生成章节内容
- 审稿与重写
最后输出一篇完整的技术博客。
整个系统由多个 Agent 协作完成:PlannerAgent 负责规划,WriterAgent 基于 ReAct loop 写作,ReviewAgent 负责审稿,SearchAgent 提供实时搜索能力。架构大概是这样:
浏览器 (React + TypeScript)
│ POST /jobs 创建任务
│ GET /jobs/stream SSE 实时进度
│ POST /jobs/reply 大纲确认
▼
FastAPI + asyncio
└── Orchestrator(流水线调度)
├── PlannerAgent 生成大纲
├── OpinionAgent 生成章节论点 + 搜索方向
├── WriterAgent ←── SearchAgent(注入)
└── ReviewAgent 轻审(每章)+ 全文重审WriterAgent 是整个系统里最有意思的部分——它不是被动接收资料然后写作,而是基于 ReAct loop 在写作过程中自主决定什么时候需要搜索、搜什么:
# 核心 agent loop
while stop_reason == "tool_use":
response = LLM(messages, tools) # LLM 决定要不要调工具
execute tools → collect results # 执行搜索
messages.append(tool_results) # 结果反馈给 LLM
return final text # LLM 决定停止时返回前端通过 SSE 实时接收每一个 token,写作过程可以实时预览。
系统虽然能跑,但随着 workflow 越来越复杂,我开始意识到:
问题已经不是“LLM 能不能生成内容”,而是:
- workflow 该由谁控制
- 状态该由谁维护
- 多个 Agent 之间如何协作
- 系统如何 tracing 和 recover
后来去看了 deer-flow 的源码,我才意识到:
我写的并不是一个真正的 Agent workflow,更像是一个“披着 Agent 外衣的硬编码流水线”。
问题一:谁在掌控流程?
我原来的 Orchestrator 大概长这样:
async def _run_pipeline(self, job):
# Stage 1: 生成大纲
chapters = await self._planner.plan(self.topic)
# Stage 2: 并行写作
results = await asyncio.gather(*[
self._write_chapter(title, ...) for title in chapters
])
# Stage 3: 全文审稿,不通过就重写,最多两轮
full_results = await self._reviewer.review_full(...)
if rewrite_tasks:
await asyncio.gather(*rewrite_tasks)
full_results2 = await self._reviewer.review_full(...)
if rewrite_tasks2:
await asyncio.gather(*rewrite_tasks2)
# 两轮之后不管结果,直接接受流程完全由 Python 代码控制:先规划,再写作,再审稿,审稿失败重写两次,两次之后不管结果强制结束。
LLM 在这里是执行器,不是决策者。
真正的 Agent 应该是:LLM 自己判断"这章写完了需要审一下",自己看到审稿反馈后决定怎么改,自己决定什么时候满意了可以停。Python 代码只提供工具,不掌控流程。
问题二:代码耦合
把 _run_pipeline 里的每一行分个类:
- 业务逻辑(调用 PlannerAgent、WriterAgent、ReviewAgent)
- 控制流(判断要不要重写、重写几次、什么时候结束)
- 基础设施(
push_eventSSE 推流、job_store.update状态更新、log.info日志、time.monotonic()计时)
数一下,基础设施的行数最多,业务逻辑反而最少,代码能运行,但 workflow 本身无法被观察、推理、可视化。
这意味着什么?如果我想改"review 失败重写几次"这个规则,我要在一堆 push_event 和 job_store.update 中间找到那几行真正相关的代码。改一个业务规则,要读懂整个函数。
问题三:状态散落
我用了一个全局的 JobStore 来维护任务状态:
# store.py - 纯内存 dict
job = job_store.get(self.job_id)
job.stage = StageStatus.WRITE
job_store.update(job)job 对象可以在任何地方被任何人修改。job.chapters 这个字段,orchestrator、_write_chapter、review 阶段都直接改它,你知道状态被改了,但不知道是谁、什么时候、为什么改的。
用 LangGraph 重构
这三个问题,LangGraph 提供了一个统一的解法:把控制流、业务逻辑、状态管理三件事分开放。
State:单一事实来源
LangGraph 的 State 本质上是一种“显式状态机”。
每个节点只能声明“自己修改了什么”,而不是直接操作整个全局对象。
在我的项目中我首先定义了一个 WriterState,所有节点共享这一份状态:
class ChapterState(TypedDict):
title: str
content: str
review_passed: bool
review_feedback: str
rewrite_count: int
class WriterState(TypedDict):
topic: str
outline: list[str]
chapters: list[ChapterState]
rewrite_count: int
final_content: str每个节点只能返回它想修改的字段,LangGraph 负责合并:
def write_node(state: WriterState) -> dict:
# 写完章节后,只返回 chapters 字段
return {"chapters": updated_chapters}
# topic、outline、rewrite_count 等字段不动这和前端的 Redux 是同一个思路:状态集中管理,谁改了什么一目了然。
原来需要手动维护的 StageStatus(PLAN/WRITE/REVIEW/EXPORT)也消失了——LangGraph 本身就知道图跑到哪个节点了,不需要额外记录。
Node:只做业务逻辑
节点函数只关心一件事:
async def plan_node(state: WriterState) -> dict:
chapters = await planner.plan(state["topic"])
return {"outline": chapters, "chapters": [...]}
async def review_node(state: WriterState) -> dict:
results = await reviewer.review_full(state["topic"], state["chapters"])
# 更新每章的 review_passed 和 review_feedback
return {"chapters": updated, "rewrite_count": state["rewrite_count"] + 1}没有 push_event,没有 job_store.update,没有重试逻辑。节点只做业务。
Edge:控制流的归宿
原来硬编码在 Python 里的"要不要重写",现在变成一个独立的条件边函数:
def should_rewrite(state: WriterState) -> str:
failed = [ch for ch in state["chapters"] if not ch["review_passed"]]
if not failed:
return "export" # 全部通过,去导出
if state["rewrite_count"] >= 2:
return "export" # 达到上限,接受现状(显式决策,不是静默放弃)
return "write" # 还有章节没通过,回去重写原来 workflow 的控制逻辑是“隐式埋在代码里的”。
现在它第一次变成了:
- 可观察的
- 可推理的
- 可修改的
- 可视化的
graph structure。
以后想改成"最多重写三次",只改这一个函数里的数字。
组装起来
builder = StateGraph(WriterState)
builder.add_node("plan", plan_node)
builder.add_node("write", write_node)
builder.add_node("review", review_node)
builder.add_node("export", export_node)
builder.add_edge(START, "plan")
builder.add_edge("plan", "write")
builder.add_edge("write", "review")
builder.add_conditional_edges("review", should_rewrite) # 条件边
builder.add_edge("export", END)光看这段代码,整个流程一目了然。
关于 push_event(SSE 推流)
我原来以为 push_event 这类基础设施可以通过 LangGraph 的 callback 机制完全从节点里抽离出来,callback 在每个节点执行前后自动触发,节点本身不需要知道它的存在。
但实际上这行不通,不是所有基础设施都适合被 framework “自动接管”,token streaming 本身就是业务逻辑的一部分。
callback 拿到的是节点的元数据(名字、输入输出),拿不到节点内部的运行状态。对于 token 级别的流式推送——"当前写到哪个字了"——只有节点自己知道,callback 没法拦截。
去看了 deer-flow 和 vibe-blog 的源码,他们也是直接在节点里调用事件推送,没有用 callback 做这件事。
Checkpointer:两行代码换来断点续写
原来的 JobStore 是纯内存的,服务重启任务就丢了。
加 LangGraph 的 checkpointer 只需要改两处:
# graph.py:编译图时传入 checkpointer
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
checkpointer = AsyncSqliteSaver.from_conn_string("writer_checkpoints.db")
graph = builder.compile(checkpointer=checkpointer)
# jobs.py:调用图时传入 thread_id
await graph.ainvoke(
initial_state,
config={"configurable": {"thread_id": job_id}},
)每次节点执行完,LangGraph 自动把当前 WriterState 序列化写入 SQLite。服务重启后,用同一个 job_id 重新调用,从上次中断的节点继续。
一些感受
做完这次重构,最大的收获不是"学会了 LangGraph",而是理解了一个更根本的问题:
流水线和 Agent 的区别,不在于用了什么框架,而在于控制权在谁手里。
流水线是 Python 说"先做 A,再做 B,再做 C"。Agent 是 LLM 说"我现在需要做 A",Python 只是提供工具和约束。
我原来的 vibe-writer 是前者——每一步都由代码决定,LLM 只是被调用的执行器。LangGraph 解决的不是"能不能跑"的问题,而是帮你把控制流从业务逻辑里分离出来,让代码结构本身就能反映这个区别。
LangGraph 不是让系统更智能,而是让系统更可维护。
当然,vibe-writer 现在还是偏流水线的——PLAN → WRITE → REVIEW → EXPORT 这个顺序还是固定的,只是"要不要重写"这个决策交给了条件边。真正让 LLM 自主决定整个流程的走向,是下一步要探索的方向。
项目地址:vibe-writer