LangGraph 多 Agent 在生产环境中崩溃的原因与代码级修复方案
26 जुलाई 2026
0
Computing/SoftwareRelated Video
10:33别再提循环工程了,图工程时代已经到来
Chase AI
Comments (0)
Log in to leave a comment
No posts yet
10:33Chase AI
Log in to leave a comment
No posts yet
从单 Prompt 链跨越到多 Agent 架构时,后端工程师经常会产生一个误区:以为只要把 Prompt 写得更好,系统就会变得稳定。然而,在生产一线爆出的问题绝大多数与 Prompt 毫无关系。状态污染、无限循环、API Rate Limit、无法排查的异步 Trace 等系统架构缺陷,才是真正的罪魁祸首。
要想把基于 LangGraph 的多 Agent 系统推上生产环境,就不能将其视为单纯的 Prompt 资源集合,而必须像对待后端系统一样,在代码层面上做好状态隔离、并发控制与链路追踪。
在基于状态的图结构中,如果多个节点直接修改同一个共享对象,就会引发竞态条件(Race Condition)。在并发执行的 Fan-out 模式下,若没有显式指定 Reducer 就直接覆盖普通字段,系统只会保留最晚结束的节点结果,其余节点的数据全都会丢失。
要解决这个问题,必须在父图状态中挂载显式的 Reducer,同时将子 Agent 完全封装为拥有独立 Schema 的子图(Subgraph)。
`python
import operator
from typing import Annotated, List, TypedDict
from langgraph.graph import END, START, StateGraph
class ParentState(TypedDict):
task_id: str
input_query: str
audit_logs: Annotated[List[str], operator.add]
final_response: str
class InternalAgentState(TypedDict):
sub_task: str
scratchpad_messages: List[str]
sub_result: str
def internal_processing_node(state: InternalAgentState) -> dict:
updated_messages = state["scratchpad_messages"] + ["내부 격리 추론 진행 중"]
return {
"scratchpad_messages": updated_messages,
"sub_result": f"하위 작업 완료: {state['sub_task']}"
}
subgraph_builder = StateGraph(InternalAgentState)
subgraph_builder.add_node("internal_processing", internal_processing_node)
subgraph_builder.add_edge(START, "internal_processing")
subgraph_builder.add_edge("internal_processing", END)
compiled_subgraph = subgraph_builder.compile()
def call_isolated_subgraph_wrapper(state: ParentState) -> dict:
subgraph_input: InternalAgentState = {
"sub_task": state["input_query"],
"scratchpad_messages": []
}
subgraph_output = compiled_subgraph.invoke(subgraph_input)
return {
"audit_logs": [f"[서브그래프 결과]: {subgraph_output['sub_result']}"]
}
parent_builder = StateGraph(ParentState)
parent_builder.add_node("isolated_agent", call_isolated_subgraph_wrapper)
parent_builder.add_edge(START, "isolated_agent")
parent_builder.add_edge("isolated_agent", END)
main_graph = parent_builder.compile()
`
在父级 ParentState 中,我们为会发生并行写入的 audit_logs 字段指定了 operator.add Reducer。而子任务则被隔离在采用独立状态 InternalAgentState 的子图中,仅通过包装函数(Wrapper)传递结果。阻断数据污染后,每周可以节省 5 小时以上的死循环排查时间。
另一个常见问题是 ReAct 反馈循环无法满足终止条件,导致程序一直转圈。此时需要在状态 Schema 中引入计数器,并在条件边中使用护栏路由(Guardrail Router)进行拦截。
`python
from typing import Literal, TypedDict
from langgraph.graph import END, START, StateGraph
class GuardedState(TypedDict):
query: str
draft: str
feedback: str
is_approved: bool
iterations: int
max_iterations: int
def drafting_node(state: GuardedState) -> dict:
return {
"draft": f"작성된 초안 (반복 회차: {state['iterations'] + 1})",
"iterations": state["iterations"] + 1
}
def review_node(state: GuardedState) -> dict:
approved = state["iterations"] >= 3
return {
"is_approved": approved,
"feedback": "승인 완료" if approved else "반려: 내용 수정 필요"
}
def loop_guardrail_router(state: GuardedState) -> Literal["drafting", "fallback_escalation", "end"]:
if state["is_approved"]:
return END
if state["iterations"] >= state["max_iterations"]:
return "fallback_escalation"
return "drafting"
def fallback_escalation_node(state: GuardedState) -> dict:
return {
"draft": "에이전트 검수 피드백 루프 최대 횟수 초과. 담당자 수동 검토 건으로 이관 처리되었습니다."
}
builder = StateGraph(GuardedState)
builder.add_node("drafting", drafting_node)
builder.add_node("review", review_node)
builder.add_node("fallback_escalation", fallback_escalation_node)
builder.add_edge(START, "drafting")
builder.add_edge("drafting", "review")
builder.add_conditional_edges(
"review",
loop_guardrail_router,
{
"drafting": "drafting",
"fallback_escalation": "fallback_escalation",
END: END
}
)
builder.add_edge("fallback_escalation", END)
guarded_graph = builder.compile()
`
当计数器(iterations)达到设定值(max_iterations)时,逻辑会自动分支到人工介入节点(fallback_escalation)。这能彻底切断因无限循环造成的 Token 浪费。
如果一次性并发触发多个子节点,极易触发 OpenAI 或 Anthropic API 的每分钟 Token 限制(TPM),引发 HTTP 429 错误。后端一旦动摇,整个事务就会陷入瘫痪。
必须使用 asyncio.Semaphore 限制并发请求数,并结合 tenacity 的指数退避(Exponential Backoff)重试机制,才能保证 API 不被冲垮。
`python
import asyncio
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
API_SEMAPHORE = asyncio.Semaphore(5)
@retry(
stop=stop_after_attempt(5),
wait=wait_exponential(multiplier=1, min=2, max=10),
retry=retry_if_exception_type(Exception),
reraise=True
)
async def safe_llm_call_with_backoff(llm: ChatOpenAI, prompt: str) -> str:
async with API_SEMAPHORE:
response = await llm.ainvoke([HumanMessage(content=prompt)])
return response.content
async def parallel_worker_node(state: dict) -> dict:
llm = ChatOpenAI(model="gpt-4o", temperature=0)
task_input = state["task_data"]
result_text = await safe_llm_call_with_backoff(llm, f"하위 작업 처리: {task_input}")
return {"results": [result_text]}
`
通过限制最大并发调用数为 5,并在失败时按 2 秒到 10 秒的指数级延长等待时间进行重试,可以完全避免因外部 API 调用失败而导致的系统中断。
在单循环模式下,只要有一个节点崩溃,就需要花 14 秒以上重新推理;但如果像这样隔离节点并加上退避机制,失败恢复时间可以降低到 10ms 左右。由于已成功节点的结果会被保留,因此不会产生多余的 Token 消耗。
对于异步交织在一起的 Agent,仅靠控制台打印日志是无法理清执行流程的。必须接入 Langfuse 等基于 OpenTelemetry 的可观测性平台,并将数据库操作或后端逻辑等非 LLM 业务也纳入 Trace Span 中,才能找出性能瓶颈。
`python
import os
from langfuse.decorators import observe, langfuse_context
from langgraph.graph import StateGraph, START, END
@observe(name="vector_store_retrieval")
def query_vector_store(query: str) -> list:
langfuse_context.update_current_observation(
input={"query": query},
metadata={"top_k": 3, "database": "pgvector"}
)
return ["문서 1: 보안 규정 예시", "문서 2: 서비스 약관"]
def retrieval_node(state: dict) -> dict:
docs = query_vector_store(state["query"])
return {"context": docs}
def execute_graph_with_tracing(app, user_query: str, session_id: str, user_id: str):
from langfuse.callback import CallbackHandler
langfuse_handler = CallbackHandler()
config = {
"configurable": {"thread_id": session_id},
"callbacks": [langfuse_handler]
}
return app.invoke({"query": user_query}, config=config)
`
通过 @observe 装饰器,Vector Search 等普通函数也可以轻松注入 Trace。在调用 Graph 时传入 CallbackHandler,即可按 Session 一目了然地查看整个 Agent 的运行轨迹与 Token 消耗量。
如果在运行到第 10 个节点时突发错误,从头重新执行会同时浪费金钱和时间。而使用 PostgresSaver 检查点管理器,节点执行的快照就会直接持久化保存到数据库中。
`python
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.graph import StateGraph
DATABASE_URL = "postgresql://postgres:postgres@localhost:5432/agent_checkpoints"
def build_app_graph():
builder = StateGraph(dict)
return builder
def resume_execution_from_failure(app, thread_id: str, fixed_payload: dict, last_valid_node: str):
config = {"configurable": {"thread_id": thread_id}}
app.update_state(
config,
values=fixed_payload,
as_node=last_valid_node
)
resumed_output = app.invoke(None, config)
return resumed_output
`
通过 get_state_history 确认最后一个正常状态后,使用 update_state 修正出问题的数据,再调用 app.invoke(None, config),就能从精准中断的位置继续向下运行。
在所有 Agent 节点上全量无脑堆 GPT-4o 是极大的预算浪费。行业标准的做法是采取 Tiering(分层)策略:主 Task Planning 使用高规格模型,而简单的分类或校验节点则挂载 Claude 3.5 Haiku 这类轻量模型。
在此基础上,配合基于 RedisVL 的 Semantic Cache(语义缓存),对于相同或相似的校验请求,甚至无需调用 LLM,在 50ms 内即可直接返回结果。
`python
from redisvl.extensions.llmcache import SemanticCache
from langchain_community.chat_models import ChatAnthropic
audit_semantic_cache = SemanticCache(
name="audit_nodes_cache",
redis_url="redis://localhost:6379",
distance_threshold=0.1,
ttl=86400
)
def audit_verification_node(state: dict) -> dict:
prompt_query = f"다음 최종 결과물의 정책 준수 여부를 검수하세요: {state['final_response']}"
cached_response = audit_semantic_cache.check(prompt=prompt_query)
if cached_response:
return {
"audit_passed": cached_response[0]["response"] == "PASSED",
"audit_logs": ["[Audit Node]: 시맨틱 캐시 데이터 활용 (LLM 호출 스킵)"]
}
audit_llm = ChatAnthropic(model="claude-3-5-haiku-20241022", temperature=0)
eval_result = audit_llm.invoke(prompt_query).content
audit_semantic_cache.store(
prompt=prompt_query,
response=eval_result,
metadata={"node": "audit_verification"}
)
return {
"audit_passed": eval_result == "PASSED",
"audit_logs": [f"[Audit Node]: 신규 모델 검수 완료 ({eval_result})"]
}
`
将 distance_threshold 严格设为 0.1 以防止误判,仅在未命中缓存时才调用轻量模型 Haiku。仅靠这一组合配置,就可以在保证整体校验质量的同时,将 API Token 费用直接砍掉最高 40%。
将 Agent 系统推向生产环境时,需要的不是炫酷的花式 Prompt 技巧,而是扎实的后端基本功——只有把状态隔离、并发控制、检查点恢复、模型分层等基础设施打牢,系统才不会在关键时刻崩盘。