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
एकल प्रॉम्प्ट चेन से आगे बढ़कर मल्टी-एजेंट आर्किटेक्चर की ओर रुख करने वाले बैकएंड डेवलपर्स अक्सर एक गलतफहमी के शिकार हो जाते हैं। उन्हें लगता है कि प्रॉम्प्ट को और बेहतर लिखने से सिस्टम स्थिर हो जाएगा। लेकिन असल में प्रोडक्शन में आने वाली अधिकांश समस्याएं प्रॉम्प्ट से संबंधित ही नहीं होती हैं। स्टेट का दूषित होना (State Pollution), इनफिनिट लूप, API रेट लिमिट, और डिबग न किए जा सकने वाले एसिंक्रोनस ट्रेसेज जैसी सिस्टम आर्किटेक्चरल समस्याएं ही इसकी असली वजह हैं।
LangGraph आधारित मल्टी-एजेंट को प्रोडक्शन एनवायरनमेंट में सफलतापूर्वक चलाने के लिए, इसे सिर्फ प्रॉम्प्ट्स का एक समूह मानने के बजाय एक बैकएंड सिस्टम की तरह संभालना होगा—जहाँ स्टेट आइसोलेशन, कॉनकरेंसी कंट्रोल, और ट्रेसिंग को कोड लेवल पर हैंडल किया जाता है।
स्टेट-बेस्ड ग्राफ में जब कई नोड्स एक ही शेयर्ड ऑब्जेक्ट को सीधे मॉडिफाई करते हैं, तो रेस कंडीशन पैदा होती है। कॉनकरेंटली चलने वाले fan-out पैटर्न में, बिना किसी रिड्यूसर के सामान्य फ़ील्ड्स को ओवरराइट करने पर केवल सबसे अंत में समाप्त होने वाले नोड का परिणाम ही बचता है और बाकी सब गायब हो जाता है।
इसे रोकने के लिए, पैरेंट ग्राफ की स्टेट में एक स्पष्ट रिड्यूसर (explicit 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 फीडबैक लूप्स का निकास स्थितियों (exit conditions) को पूरा न कर पाने के कारण गोल-गोल घूमते रहना भी आम बात है। इसके लिए स्टेट स्कीमा में एक काउंटर की आवश्यकता होती है और एक गार्डरेल राउटर की जो कंडीशनल एजेस में इसे फ़िल्टर करे।
`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) पर रीडायरेक्ट करने के लिए सेट किया गया है। यह अनंत लूप के कारण होने वाले टोकन की बर्बादी को पूरी तरह से रोकता है।
यदि चाइल्ड नोड्स को एक साथ पैरेलल रूप से चालू किया जाता है, तो वे OpenAI या Anthropic API की टोकन पर मिनट (TPM) सीमाओं से टकराएंगे, जिससे HTTP 429 एरर उत्पन्न होगा। बैकएंड के डगमगाते ही पूरा ट्रांजैक्शन ठप हो जाता है।
API को क्रैश होने से बचाने के लिए, कॉनकरेंट रिक्वेस्ट्स की संख्या को सीमित करने के लिए asyncio.Semaphore का उपयोग करें और tenacity के एक्सपोनेंशियल बैकऑफ (Exponential Backoff) को लागू करें।
`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]}
`
समवर्ती कॉल्स (concurrent calls) को अधिकतम 5 तक सीमित करें और विफल होने पर वेट टाइम को 2 सेकंड से बढ़ाकर 10 सेकंड करते हुए पुनः प्रयास करें। यह बाहरी API कॉल विफलताओं के कारण होने वाले सिस्टम डाउनटाइम को पूरी तरह से रोकता है।
सिंगल-लूप दृष्टिकोण में, यदि केवल एक नोड विफल हो जाता है, तो पूरे री-रीज़निंग के लिए 14 सेकंड से अधिक का समय खर्च करना पड़ता है। हालांकि, इस तरह नोड्स को अलग करके और बैकऑफ लागू करके, आप विफलता रिकवरी समय को केवल 10ms तक कम कर सकते हैं। चूंकि सफल नोड्स के परिणाम संरक्षित रहते हैं, इसलिए टोकन का अनावश्यक पुनरुपयोग भी नहीं होता है।
एसिंक्रोनस रूप से जुड़े एजेंट्स के एग्जीक्यूशन फ्लो को केवल कंसोल आउटपुट से समझना असंभव है। Langfuse जैसे OpenTelemetry-आधारित ऑब्जर्वेबिलिटी प्लेटफॉर्म को एकीकृत करें, और अड़चनों को पहचानने के लिए नॉन-LLM लॉजिक जैसे कि DB ऑपरेशन्स या बैकएंड प्रोसेसिंग को भी ट्रेसिंग स्पैन में शामिल करें।
`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 चेकपॉइंटर का उपयोग करके, नोड निष्पादन का स्नैपशॉट सीधे DB में सहेजा जाता है।
`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 को थोपना बजट की बर्बादी है। मुख्य प्लानिंग के लिए हाई-परफॉरमेंस मॉडल का उपयोग करना और सरल वर्गीकरण या समीक्षा नोड्स के लिए Claude 3.5 Haiku जैसे हल्के मॉडल को जोड़ना मॉडल टियरिंग का मूल सिद्धांत है।
इसके अलावा, RedisVL आधारित सिमेंटिक कैश सेट करने पर, समान या मिलते-जुलते रिव्यू रिक्वेस्ट्स बिना 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})"]
}
`
गलत पॉजिटिव (false positives) को रोकने के लिए distance_threshold को कड़ाई से 0.1 पर सेट किया जाता है, और लाइटवेट मॉडल Haiku को केवल तभी कॉल किया जाता है जब डेटा कैश में उपलब्ध न हो। केवल इसी कॉन्फ़िगरेशन को लागू करने से, समग्र समीक्षा गुणवत्ता को बनाए रखते हुए API टोकन लागत में 40% तक की कमी आ सकती है।
मल्टी-एजेंट सिस्टम को प्रोडक्शन में ले जाने के लिए आकर्षक प्रॉम्प्ट तकनीकों की आवश्यकता नहीं होती है। सिस्टम को क्रैश होने से बचाने के लिए सॉलिड बैकएंड इंफ्रास्ट्रक्चर—जैसे स्टेट आइसोलेशन, कॉनकरेंसी कंट्रोल, चेकपॉइंट रीज्यूम और मॉडल टियरिंग—का होना बेहद जरूरी है।