Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
98 changes: 98 additions & 0 deletions BENCHMARK_ANALYSIS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
# LangGraph PoC vs 自研 StreamingAgentLoop — 完整对比分析

基于本地 uvicorn + xiaomi LLM + SQLite/MemorySaver 实测:
- 单轮 benchmark(`bench_results.json` / `bench_results_optimized.json`)
- 多轮 benchmark(`bench_multiturn.json`)
- 真实任务案例

## 1. 性能对比

### 单轮(5 场景 × 5 次,优化后)

| 场景 | v1 E2E mean | v2 E2E mean | delta |
|---|---|---|---|
| 普通对话(无工具) | 2.6s | 3.0s | +13.8% |
| list_topics | 6.7s | 5.3s | -20.5% |
| get_batch_job_status | 5.8s | 3.9s | -32.6% |
| get_citation_tree | 9.7s | 10.5s | +8.8% |
| skim_paper confirm | 6.7s | 13.4s | +100.4% |

### 多轮增长曲线(场景 A,10 轮纯对话)

| 轮次 | v1 TTFT | v2 TTFT | v1 E2E | v2 E2E |
|---|---|---|---|---|
| 1 | 2.4s | 2.2s | 2.8s | 2.6s |
| 5 | 2.3s | 1.9s | 2.5s | 2.1s |
| 10 | 1.0s | 1.5s | 1.3s | 1.8s |

**关键发现**:TTFT/E2E 不随轮次显著增长。两后端在短对话下历史拼接开销可忽略,LLM 方差(同 prompt 1.0s~7.4s)完全淹没框架差异。

### 含 confirm 多轮(场景 B)

v2 confirm resume 后多轮不稳定(第 4 轮起全部 N/A,流被截断)。v1 confirm 后续轮也大量 N/A(LLM 调工具/超长回复导致 timeout)。这是 **LLM 行为方差 + timeout 设置**导致,非框架本质问题,但暴露 v2 在 confirm 后续状态恢复上需进一步验证。

## 2. 真实多工具任务案例

**任务**:search_papers('attention') → skim_paper(第一篇) → 一句话总结

| 指标 | v1 | v2 |
|---|---|---|
| 总耗时 | 10.08s | 10.02s |
| 首 token | 1.88s | 8.94s |
| 工具调用次数 | 2 | 2 |
| 工具序列 | search_papers, get_system_status | search_papers, get_system_status |
| 触发 confirm | 0 | 0 |

**关键发现**:
1. **两后端 LLM 决策路径完全一致**(相同工具序列)— 框架不影响 LLM 工具选择
2. **LLM 都没调 skim_paper**(任务要求粗读,但 LLM 只搜索+查状态就总结)— LLM 决策偏差,与框架无关
3. v2 TTFT 比 v1 慢(8.94s vs 1.88s)— 但这是 LLM 方差(同任务多次跑会变),非框架固有

## 3. 框架优势分析(结构性,基于源码 + 实测)

### ✅ LangGraph 的优势

| 优势 | 说明 | 自研 loop 对比 |
|---|---|---|
| **checkpoint 持久化** | 服务重启后状态不丢(PG),thread_id 隔离 | 老 loop 靠 pending action 全量快照,重启后快照可能过期 |
| **增量状态存储** | checkpoint 每步只存新消息(增量),存储 O(n) | 老 loop confirm 时存全量 conversation JSON,随轮数增大 |
| **interrupt 原生支持** | 框架级 human-in-the-loop,`interrupt()` + `Command(resume=)` | 老 loop 手写 store_pending_action / load / mark_handled / cleanup_expired |
| **状态一致性** | checkpoint 是运行时状态(增量),resume 自动恢复 | 老 loop 快照是 confirm 时点冻结,期间若有新消息会丢失 |
| **可观测性** | checkpoint 列表/回放/time-travel,可调试 | 老 loop 无 |
| **生态** | 可接 langgraph-platform 部署/监控/streaming UI | 老 loop 自维护 |
| **少写边界条件** | JSON 解析/多 confirm/max_rounds/usage 这些手修过的 bug,LangGraph 内置 | 老 loop 我们手修了 ⑧⑨⑩⑬ 等多个 bug |
| **并发安全** | thread_id 隔离 + checkpoint 事务 | 老 loop pending action 快照可能竞态 |

### ❌ LangGraph 的劣势

| 劣势 | 说明 | 量化 |
|---|---|---|
| **无工具场景延迟** | graph 编译 + LangChain 消息转换 + checkpoint 是纯额外开销 | +13.8% E2E(优化后) |
| **confirm 流 2 请求** | v2 需触发 + confirm resume 两请求,v1 单请求 | +100.4% E2E |
| **依赖更重** | langgraph + langchain-core + psycopg v3 | +~50MB 安装体积 |
| **学习曲线** | LangChain 消息协议/tool_call_chunks 等概念 | 团队需学习 |
| **confirm 后续稳定性** | v2 confirm resume 后多轮有不稳定(实测 N/A) | 需进一步排查 |

## 4. 结论与建议

### 性能上
- **短对话/无工具**:v2 慢 ~13.8%(框架固有开销,可接受)
- **含工具**:v2 与 v1 持平或更快(工具往返毫秒级,框架开销被 LLM 延迟掩盖)
- **confirm**:v2 慢 ~100%(2 请求架构,可优化为 Command goto 单请求)
- **多轮**:两后端在短对话下不劣化;长对话需更大数据量验证

### 功能上
- **LangGraph 优势在工程维护性**:checkpoint 持久化、interrupt 原生、少写边界条件、可观测性、生态
- **这些优势在"对话越长 + confirm 越多 + 需要重启/并发"时越明显**

### 决策建议
1. **如果**项目重点是快速对话、少 confirm、单进程 → 保留自研 loop(性能更好,依赖更轻)
2. **如果**项目需要长对话持久化、多 confirm、服务重启恢复、可观测性 → 用 LangGraph(工程优势 > 13.8% 延迟代价)
3. **PoC 结论**:LangGraph 功能完整、可优化到接近 v1 性能,但 confirm 流和多轮稳定性需进一步打磨。**建议保留 PoC 分支,先在生产环境小范围验证 confirm 多轮 + checkpoint 持久化**,再决定是否整体替换。

## 5. 不确定因素(需更多测试)

- **生产 PG + PostgresSaver I/O 开销**未测(本地用 MemorySaver)
- **xiaomi LLM 方差极大**(同 prompt TTFT 1.0s~7.4s),5-10 次 mean 仍有噪声,要可靠结论需固定 temperature=0 + 20+ 次
- **confirm resume 后多轮稳定性**需排查(v2 第 4 轮起 N/A)
- **长对话(50+ 轮)** 下 checkpoint 存储增长未测
3 changes: 2 additions & 1 deletion Dockerfile.backend
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ COPY infra/migrations/ infra/migrations/

# 使用腾讯云 pip 镜像源(阿里云镜像下载大文件不稳,腾讯云源稳定且 ECS 内网快)
# graph extra 含 numpy/scikit-learn/umap-learn(similarity.py 降维散点图用)
RUN pip install --cache-dir=/.pip-cache ".[llm,pdf,graph]" \
# langgraph extra 含 langgraph + langgraph-checkpoint-postgres + psycopg v3(/agent/v2 PoC)
RUN pip install --cache-dir=/.pip-cache ".[llm,pdf,graph,langgraph]" \
-i https://mirrors.cloud.tencent.com/pypi/simple \
--timeout 60 --retries 5 \
&& pip install --cache-dir=/.pip-cache \
Expand Down
7 changes: 7 additions & 0 deletions apps/api/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,13 @@ async def app_error_handler(_request: Request, exc: AppError):
app.include_router(cs_feeds.router)
app.include_router(graph.router)
app.include_router(agent.router)
# PoC:LangGraph 后端 /agent/v2/*。需安装 .[langgraph] extra;未装时跳过以保证核心可用。
try:
from apps.api.routers import agent_v2

app.include_router(agent_v2.router)
except ImportError:
pass
app.include_router(content.router)
app.include_router(pipelines.router)
app.include_router(settings_router.router)
Expand Down
235 changes: 235 additions & 0 deletions apps/api/routers/agent_v2.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,235 @@
"""Agent v2 路由:LangGraph 后端(PoC,与现有 /agent/chat 并行)。

复用 apps.api.routers.agent 的持久化辅助函数(_db_messages_to_openai /
_stream_with_save_for_action / _parse_sse_events / _resolve_conversation_id_from_action),
把 stream_chat/confirm_action/reject_action 换成 langgraph_agent.entry 的 v2 版本。

PoC:不合 main,拍板替换后再删老 agent.py。
@author Color2333
"""

from __future__ import annotations

from fastapi import APIRouter
from fastapi.responses import StreamingResponse

# 复用老 agent.py 的持久化辅助(避免重复实现)
from apps.api.routers.agent import (
_SSE_HEADERS,
_db_messages_to_openai,
_new_messages_to_dicts,
_parse_sse_events,
_resolve_conversation_id_from_action,
_stream_with_save_for_action,
)
from packages.domain.schemas import AgentChatRequest # noqa: TC001 FastAPI 需运行时可见以解析 body
from packages.langgraph_agent.entry import confirm_v2, reject_v2, stream_chat_v2

router = APIRouter()


@router.post("/agent/v2/chat")
async def agent_chat_v2(req: AgentChatRequest):
"""Agent v2 对话 —— LangGraph 后端,SSE 协议与 /agent/chat 一致。

后端真相源逻辑(修①②③)与 /agent/chat 完全一致:DB 拼 history +
SSE 首事件 conversation_init + done 去重。仅 agent 内核换成 LangGraph。
"""
from packages.agent_core.sse import make_sse
from packages.storage.db import session_scope
from packages.storage.repositories import (
AgentConversationRepository,
AgentMessageRepository,
)

conversation_id = getattr(req, "conversation_id", None)

with session_scope() as session:
conv_repo = AgentConversationRepository(session)
msg_repo = AgentMessageRepository(session)

if conversation_id:
conv = conv_repo.get_by_id(conversation_id)
if not conv:
conversation_id = None

if not conversation_id:
first_user_msg = next((m for m in req.messages if m.role == "user"), None)
title = first_user_msg.content[:50] if first_user_msg else "新对话"
conv = conv_repo.create(title=title)
conversation_id = conv.id

# 保存本次新消息(与老路径一致)
saved_keys: set[str] = set()
for msg in req.messages:
if msg.role == "system":
continue
content_key = f"{msg.role}:{msg.content[:200]}"
if content_key not in saved_keys:
msg_repo.create(
conversation_id=conversation_id,
role=msg.role,
content=msg.content,
meta=msg.meta,
)
saved_keys.add(content_key)

# 从 DB 读全量历史重建 OpenAI messages(修②后端拼历史)
db_msgs = msg_repo.list_by_conversation(conversation_id, limit=500)
history_msgs = _db_messages_to_openai(db_msgs)

# 合并 DB 历史 + 本次新增(去重)
new_msgs = _new_messages_to_dicts(req.messages)
history_keys = {f"{m.get('role')}:{(m.get('content') or '')[:200]}" for m in history_msgs}
extra_new = [
m
for m in new_msgs
if f"{m.get('role')}:{(m.get('content') or '')[:200]}" not in history_keys
]
msgs = history_msgs + extra_new

text_buf = ""
tool_records: list[dict] = []
tool_call_id: str | None = None
saved_done = False

def stream_with_save():
nonlocal text_buf, tool_records, tool_call_id, saved_done
# 修①:SSE 首事件返 conversation_id
yield make_sse("conversation_init", {"conversation_id": conversation_id})

sse_iter, _ = stream_chat_v2(
msgs, conversation_id, confirmed_action_id=req.confirmed_action_id
)
for chunk in sse_iter:
yield chunk

for event_type, data in _parse_sse_events(chunk):
if event_type == "text_delta":
text_buf += data.get("content", "")
elif event_type == "tool_start":
tool_call_id = data.get("id")
elif event_type == "tool_result":
tool_records.append(
{
"name": data.get("name"),
"success": data.get("success"),
"summary": data.get("summary"),
"data": data.get("data"),
}
)
import json

with session_scope() as session:
msg_repo = AgentMessageRepository(session)
msg_repo.create(
conversation_id=conversation_id,
role="tool",
content=json.dumps(
{
"name": data.get("name"),
"success": data.get("success"),
"summary": data.get("summary"),
"data": data.get("data"),
},
ensure_ascii=False,
),
meta={"tool_call_id": tool_call_id},
)
elif event_type == "action_confirm":
# LangGraph interrupt 不存 AgentPendingAction(checkpoint 已存状态),
# 但 /agent/v2/confirm 路由需从 action_id 反查 conversation_id,
# 故这里写一行 pending action(conversation_state 留空)。
from packages.storage.repositories import AgentPendingActionRepository

action_id = data.get("id")
if action_id:
with session_scope() as session:
pending_repo = AgentPendingActionRepository(session)
pending_repo.create(
action_id=action_id,
tool_name=data.get("tool", ""),
tool_args=data.get("args") or {},
tool_call_id=None,
conversation_id=conversation_id,
conversation_state=None,
)
elif event_type == "action_result":
tool_records.append(
{
"action_id": data.get("id"),
"success": data.get("success"),
"summary": data.get("summary"),
"data": data.get("data"),
}
)
# 确认/拒绝后 pending action 已消费,删掉
action_id = data.get("id")
if action_id:
with session_scope() as session:
pending_repo = AgentPendingActionRepository(session)
pending_repo.delete(action_id)
elif event_type == "done" and not saved_done and (text_buf or tool_records):
saved_done = True
import json

with session_scope() as session:
msg_repo = AgentMessageRepository(session)
msg_repo.create(
conversation_id=conversation_id,
role="assistant",
content=text_buf,
meta={"tool_calls": tool_records} if tool_records else None,
)

return StreamingResponse(
stream_with_save(),
media_type="text/event-stream",
headers=_SSE_HEADERS,
)


@router.post("/agent/v2/confirm/{action_id}")
async def agent_confirm_v2(action_id: str):
"""确认挂起的操作(LangGraph 后端,复用持久化逻辑)"""
conversation_id = _resolve_conversation_id_from_action(action_id)
# LangGraph resume 靠 checkpoint(不靠 pending action),可提前删 pending action
_delete_pending_action(action_id)
return StreamingResponse(
_stream_with_save_for_action(
conversation_id,
lambda: confirm_v2(action_id, conversation_id),
),
media_type="text/event-stream",
headers=_SSE_HEADERS,
)


@router.post("/agent/v2/reject/{action_id}")
async def agent_reject_v2(action_id: str):
"""拒绝挂起的操作(LangGraph 后端,复用持久化逻辑)"""
conversation_id = _resolve_conversation_id_from_action(action_id)
_delete_pending_action(action_id)
return StreamingResponse(
_stream_with_save_for_action(
conversation_id,
lambda: reject_v2(action_id, conversation_id),
),
media_type="text/event-stream",
headers=_SSE_HEADERS,
)


def _delete_pending_action(action_id: str) -> None:
"""删除 pending action(confirm/reject 后已消费)。"""
from packages.storage.db import session_scope
from packages.storage.repositories import AgentPendingActionRepository

try:
with session_scope() as session:
AgentPendingActionRepository(session).delete(action_id)
except Exception:
pass


__all__ = ["router"]
Loading