Почему мультиагенты LangGraph падают в продакшене и как восстановить их на уровне кода
26 июля 2026 г.
0
Computing/SoftwareComments (0)
Log in to leave a comment
No posts yet
Log in to leave a comment
No posts yet
У бэкенд-разработчиков, переходящих от простых цепочек промптов к мультиагентным системам, часто возникает одно заблуждение. Они думают, что если написать промпт получше, система станет стабильнее. Но большинство проблем на практике вообще не связаны с промптами. Настоящие причины — это системные архитекурные изъяны: загрязнение состояния, бесконечные циклы, API Rate Limit и невозможность отладки асинхронных трассировок.
Чтобы вывести мультиагентную систему на базе LangGraph в продакшен, нужно относиться к ней не как к набору ресурсов с промптами, а как к бэкенд-системе, обеспечивая изоляцию состояний, управление конкурентностью и трассировку прямо на уровне кода.
В графе, основанном на состояниях, когда несколько узлов напрямую модифицируют один и тот же общий объект, возникает состояние гонки (race condition). В паттерне fan-out при параллельном выполнении, если перезаписывать обычные поля без специального редюсера (reducer), сохранятся только результаты узла, завершившегося последним, а остальные данные просто потеряются.
Чтобы предотвратить это, нужно прикрепить явный редюсер к состоянию верхнего уровня, а под-агентов полностью инкапсулировать в подграфы с независимой схемой.
`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. Подзадачи изолированы в подграф с собственным состоянием InternalAgentState, а обмен данными происходит только через функцию-обертку. Благодаря предотвращению повреждения данных время, затрачиваемое на отладку бесконечных циклов, сокращается более чем на 5 часов в неделю.
Также часто встречается проблема, когда цикл обратной связи ReAct зацикливается, не имея возможности выполнить условие выхода. В этом случае в схему состояния помещают счетчик и используют роутер-гардрейл, фильтрующий выполнение на условных ребрах.
`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). Это надежно предотвращает бесполезную трату токенов из-за бесконечного цикла.
Если запускать дочерние узлы одновременно в параллельном режиме, система упрется в лимит токенов в минуту (TPM) API OpenAI или Anthropic, выбивая ошибку HTTP 429. В момент, когда бэкенд начинает сбоить, вся транзакция парализуется.
Чтобы предотвратить падение API, необходимо ограничить количество одновременных запросов с помощью asyncio.Semaphore и применить экспоненциальную задержку (Exponential Backoff) из библиотеки tenacity.
`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 секунд на полный повторный вывод. Но если изолировать узлы и настроить экспоненциальную задержку, время восстановления после сбоя снижается до уровня ~10 мс. Результаты уже успешных узлов сохраняются, поэтому повторного расхода токенов не происходит.
Понимание потока выполнения асинхронно взаимосвязанных агентов невозможно получить лишь по консольному выводу. Узкие места становятся видны только тогда, когда вы подключаете платформу наблюдаемости на базе OpenTelemetry (например, Langfuse) и включаете в спаны трассировки не только LLM, но и бэкенд-логику вроде операций с БД.
`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 даже обычные функции, такие как векторный поиск, попадают в трассировку. Передавая CallbackHandler при вызове графа, можно наглядно отслеживать поведение всей агентной системы и расход токенов по сессиям.
Если посреди выполнения на 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), выполнение возобновляется ровно с того места, где произошел сбой.
Использование GPT-4o во всех узлах агентов — это пустая трата бюджета. Базовым подходом является уровнирование (tiering): использование высокопроизводительной модели для основного планирования и легких моделей, таких как Claude 3.5 Haiku, для простых узлов классификации или проверки.
Если к этому добавить семантический кэш на базе RedisVL, то идентичные или похожие запросы на проверку будут возвращаться всего за 50 мс без единого вызова LLM.
`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-токены до 40%, сохраняя при этом общее качество проверки.
Чтобы вывести агентную систему в продакшен, нужны не хитроумные трюки с промптами. Система не рухнет только тогда, когда заложен прочный бэкенд-фундамент: изоляция состояний, контроль конкурентности, возобновление по чекпоинтам и уровнирование моделей.