لماذا تنهار وكلاء LangGraph المتعددون في بيئة الإنتاج وكيفية إصلاحها على مستوى الكود
July 26, 2026
0
Computing/SoftwareComments (0)
Log in to leave a comment
No posts yet
Log in to leave a comment
No posts yet
هناك وهم شائع يقع فيه مطورو الجزء الخلفي (Backend Developers) الذين ينتقلون من سلسلة التوجيهات الأحادية (Single Prompt Chains) إلى الوكلاء المتعددين (Multi-Agents). يكمن هذا الوهم في الاعتقاد بأن كتابة مطالبات (Prompts) أفضل ستجعل النظام مستقرًا. لكن معظم المشاكل التي تظهر في بيئة العمل الفعلية لا علاقة لها بالمطالبات على الإطلاق. السبب الحقيقي يرجع إلى مشاكل في بنية النظام، مثل تلوث الحالة (State Contamination)، والحلقات اللانهائية (Infinite Loops)، وحدود معدل API (Rate Limits)، والتتبعات غير المتزامنة (Async Traces) المستحيل تصحيح أخطائها.
لتشغيل وكلاء متعددين قائمين على LangGraph في بيئة الإنتاج، يجب عليك التعامل مع النظام كمنظومة خلفية لضبط إزاحة الحالات، والتحكم في التزامن، والتتبع على مستوى الكود، بدلاً من التعامل معه كمجرد مجموعة من موارد المطالبات.
في الرسوم البيانية القائمة على الحالة، يؤدي تعديل كائن مشترك واحد مباشرة بواسطة عدة عقد إلى حدوث ظروف السباق (Race Conditions). في نمط التفرع للخارج (Fan-out) المنفذ بالتزامن، إذا تم الكتابة فوق حقل عادي دون استخدام مخفض (Reducer) منفصل، فستبقى نتيجة العقدة التي تنتهي أخيرًا فقط، وتضيع بقية النتائج.
لمنع ذلك، يجب إرفاق مخفض صريح بحالة الرسم البياني الأعلى، وتغليف الوكلاء الفرعيين بالكامل كرسوم بيانية فرعية ذات مخططات (Schemas) مستقلة.
`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 العلوية، قمنا بتحديد المخفض operator.add للحقل audit_logs الذي تحدث فيه كتابة متوازية. وتم تعزيل المهمة الفرعية كرسهم بياني فرعي يستخدم حالته الخاصة InternalAgentState، وتبادل النتائج يتم فقط عبر دالة غلاف (Wrapper Function). من خلال منع تلوث البيانات، يمكن تقليل وقت تصحيح أخطاء الحلقات اللانهائية بأكثر من 5 ساعات أسبوعيًا.
تُعتبر مشكلة الدوران المفرط لحلقة التغذية الراجعة ReAct بسبب عدم استيفاء شروط الإنهاء من المشاكل الشائعة أيضًا. هناك حاجة إلى موجه حماية (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). هذا يقطع بوضوح هدر الرموز (Tokens) الناتجة عن التكرار اللانهائي.
إطلاق العقد الفرعية بالتوازي دفعة واحدة يؤدي إلى تجاوز حد الرموز في الدقيقة (TPM) لـ OpenAI أو Anthropic API، مما يسبب ظهور خطأ HTTP 429. في اللحظة التي يتأرجح فيها الجزء الخلفي، تتوقف المعاملة بأكملها.
يجب تقييد عدد الطلبات المتزامنة باستخدام asyncio.Semaphore وتطبيق التراجع الأسي (Exponential Backoff) من مكتبة tenacity لمنع انهيار 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، وفي حالة الفشل يتم مضاعفة وقت الانتظار من ثانيتين إلى 10 ثوانٍ مع إعادة المحاولة. هذا يمنع تمامًا توقف النظام الناتج عن فشل استدعاء API الخارجي.
في طريقة الحلقة الأحادية، يؤدي انهيار عقدة واحدة إلى إعادة الاستدلال بالكامل واستغراق أكثر من 14 ثانية، ولكن عزل العقد وتطبيق التراجع الأسي بهذا الشكل يقلل وقت التعافي من الفشل إلى مستوى 10 ميلي ثانية. نظرًا للحفاظ على نتائج العقد الناجحة كما هي، لا يوجد استهلاك غير ضروري للرموز.
لا يمكن تتبع تدفق الوكلاء المتشابكين بشكل غير متزامن عن طريق طباعة وحدة التحكم (Console Output) فقط. يجب ربط منصة مراقبة قائمة على OpenTelemetry مثل Langfuse، وإدراج المنطق غير المتعلق بنماذج اللغة (Non-LLM) مثل عمليات قاعدة البيانات ومعالجة الجزء الخلفي في نطاق التتبع (Tracing 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، يتم إدراج الدوار المعتادة مثل البحث المتجهي في التتبع. عند تمرير CallbackHandler أثناء استدعاء الرسم البياني، يمكن التحقق من تحركات جميع الوكلاء واستهلاك الرموز لكل جلسة بلمحة واحدة.
إذا حدث خطأ في العقدة العاشرة أثناء التنفيذ وتم إعادة التشغيل من البداية، فسيضيع المال والوقت معًا. باستخدام PostgresSaver لحفظ نقاط التحقق، يتم حفظ لقطة (Snapshot) تنفيذ العقدة كما هي في قاعدة البيانات.
`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 بعقد التصنيف البسيط أو المراجعة.
إذا تم تطبيق الذاكرة المؤقتة الدلالية (Semantic Cache) القائمة على 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% مع الحفاظ على جودة المراجعة الإجمالية.
ما تحتاجه عند نقل نظام الوكلاء إلى بيئة الإنتاج ليس حيل المطالبات العجيبة. بل البنية التحتية الصلبة للجزء الخلفي مثل عزل الحالة، والتحكم في التزامن، واستئناف نقاط التحقق، وتقسيم مستويات النماذج لضمان عدم انهيار النظام.