•26 min read

LangGraphを本番環境で使う:Human-in-the-Loop、チェックポイント、状態永続化

LangGraphを本番環境で使う:Human-in-the-Loop、チェックポイント、状態永続化

LangGraphは、ステートフルなマルチアクターアプリケーションを構築するための堅牢なフレームワークを提供します。これらのエージェントワークフローを本番環境にデプロイするには、状態管理、フォールトトレランス、および人間による介入を慎重に検討する必要があります。このガイドでは、状態スキーマ、永続的なチェックポイント、ヒューマン・イン・ザ・ループ(HITL)メカニズム、および動的なサブグラフ管理に焦点を当て、本番環境レベルのLangGraphアプリケーションの実装について詳しく説明します。

Audio Briefing
0:00 / 0:00

コアコンセプト

LangGraphの力は、エージェントワークフローを有向非巡回グラフ(DAG)または状態を持つ巡回グラフとしてモデル化する能力に由来します。グラフ内の各ノードはステップを表し、エッジは遷移を定義します。ノード間で渡される状態オブジェクトは、コンテキストを維持し、複雑な相互作用を可能にする上で中心的な役割を果たします。

状態スキーマの定義

適切に定義された状態スキーマは、保守性、型安全性、およびデータ整合性にとって不可欠です。LangGraphは、この目的のためにTypedDictとPydanticを活用しています。Pydanticは検証、シリアライズ、デシリアライズを提供するため、本番環境に最適です。

ユーザーのクエリ、過去のやり取り、および潜在的な解決策に関する情報が必要なカスタマーサポートを支援するエージェントを考えてみましょう。

from typing import List, Literal, TypedDict, Optional
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage
from pydantic import BaseModel, Field

# Define the state using TypedDict for basic structure
class AgentState(TypedDict):
    """
    Represents the state of our agentic workflow.
    """
    chat_history: List[BaseMessage]
    user_query: str
    tool_output: Optional[str]
    next_action: Optional[Literal["call_tool", "respond_to_user", "human_review"]]
    escalation_reason: Optional[str]
    # Pydantic models can be nested within TypedDict for richer validation
    customer_profile: Optional['CustomerProfile']

# Define a Pydantic model for a nested object within the state
class CustomerProfile(BaseModel):
    customer_id: str
    tier: Literal["bronze", "silver", "gold", "platinum"] = "bronze"
    recent_tickets: List[str] = Field(default_factory=list)
    is_vip: bool = False

# Example of how to use Pydantic for validation and default values
# This Pydantic model can be used to validate and parse the AgentState
class AgentStatePydantic(BaseModel):
    chat_history: List[BaseMessage]
    user_query: str
    tool_output: Optional[str] = None
    next_action: Optional[Literal["call_tool", "respond_to_user", "human_review"]] = None
    escalation_reason: Optional[str] = None
    customer_profile: Optional[CustomerProfile] = None

    class Config:
        arbitrary_types_allowed = True # Allow BaseMessage

この例では、AgentStateが全体構造を定義し、CustomerProfileは顧客データに構造化された検証を提供するPydanticモデルです。この分離により、明確さと再利用性が向上します。

チェックポイントと状態の永続化

本番システムでは、障害からの回復、長時間実行プロセスの再開、デバッグを可能にするために状態の永続化が必要です。LangGraphは、このためのCheckpointer実装を提供します。

MemorySaver(開発/テスト)

ローカル開発やテストには、MemorySaverで十分です。これはチェックポイントをメモリに保存します。

from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import StateGraph, END

# Define a simple graph for demonstration
class SimpleAgentState(TypedDict):
    value: int

def increment_node(state: SimpleAgentState):
    return {"value": state["value"] + 1}

builder = StateGraph(SimpleAgentState)
builder.add_node("increment", increment_node)
builder.set_entry_point("increment")
builder.add_edge("increment", END)

memory_checkpoint = MemorySaver()
graph_with_memory = builder.compile(checkpointer=memory_checkpoint)

# Run the graph
config = {"configurable": {"thread_id": "thread-1"}}
result = graph_with_memory.invoke({"value": 0}, config=config)
print(f"MemorySaver Result: {result}") # {'value': 1}

# Resume from checkpoint
result_resume = graph_with_memory.invoke({"value": 100}, config=config) # This will overwrite if not careful
# To truly resume, you'd typically load the state and then invoke
# For MemorySaver, subsequent invokes with the same thread_id will use the last state
result_resume_correct = graph_with_memory.invoke(None, config=config) # Invoking with None uses the last checkpointed state
print(f"MemorySaver Resumed Result: {result_resume_correct}") # {'value': 2}

PostgresSaver(本番)

本番環境では、PostgreSQLのような永続ストアが不可欠です。PostgresSaverはPostgreSQLデータベースと統合してチェックポイントを保存および取得します。

まず、PostgreSQLデータベースが稼働しており、psycopg2-binaryパッケージがインストールされていることを確認してください(pip install psycopg2-binary)。

import os
from langgraph.checkpoint.postgres import PostgresSaver
from langgraph.graph import StateGraph, END
from typing import TypedDict

# Database connection string
# In production, use environment variables or a secret management system
DATABASE_URL = os.getenv("DATABASE_URL", "postgresql://user:password@localhost:5432/langgraph_db")

# Ensure your database and table exist.
# A simple table creation script:
# CREATE TABLE IF NOT EXISTS checkpoints (
#     thread_id VARCHAR(255) PRIMARY NULL,
#     checkpoint JSONB NOT NULL,
#     metadata JSONB,
#     created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
# );

# Define a simple state
class PersistentAgentState(TypedDict):
    count: int
    messages: list[str]

# Define a node function
def process_step(state: PersistentAgentState):
    new_count = state["count"] + 1
    new_messages = state["messages"] + [f"Processed step {new_count}"]
    return {"count": new_count, "messages": new_messages}

# Build the graph
builder = StateGraph(PersistentAgentState)
builder.add_node("step_one", process_step)
builder.add_node("step_two", process_step)
builder.set_entry_point("step_one")
builder.add_edge("step_one", "step_two")
builder.add_edge("step_two", END)

# Initialize PostgresSaver
postgres_checkpoint = PostgresSaver(conn_string=DATABASE_URL)
graph_with_postgres = builder.compile(checkpointer=postgres_checkpoint)

# Example usage:
thread_id = "customer_interaction_123"
config = {"configurable": {"thread_id": thread_id}}

# Initial invocation
print(f"--- Initial Invocation for thread {thread_id} ---")
initial_state = {"count": 0, "messages": ["Start"]}
result_initial = graph_with_postgres.invoke(initial_state, config=config)
print(f"Result after first run: {result_initial}")
# Expected: {'count': 2, 'messages': ['Start', 'Processed step 1', 'Processed step 2']}

# Simulate a crash and resume
print(f"\n--- Resuming Invocation for thread {thread_id} ---")
# To resume, we invoke with None, which tells LangGraph to load the last checkpoint
result_resume = graph_with_postgres.invoke(None, config=config)
print(f"Result after resuming (should be same as initial if graph ended): {result_resume}")

# Let's modify the graph to show actual resumption
# Add another step and make it loop for demonstration
builder_loop = StateGraph(PersistentAgentState)
builder_loop.add_node("step_one", process_step)
builder_loop.add_node("step_two", process_step)
builder_loop.add_node("step_three", process_step) # New step
builder_loop.set_entry_point("step_one")
builder_loop.add_edge("step_one", "step_two")
builder_loop.add_edge("step_two", "step_three")
builder_loop.add_edge("step_three", END) # Changed to END for now

graph_with_postgres_loop = builder_loop.compile(checkpointer=postgres_checkpoint)

thread_id_loop = "customer_interaction_loop_456"
config_loop = {"configurable": {"thread_id": thread_id_loop}}

print(f"\n--- Looping Invocation for thread {thread_id_loop} ---")
# First run, will go through all steps
result_loop_1 = graph_with_postgres_loop.invoke({"count": 0, "messages": ["Loop Start"]}, config=config_loop)
print(f"Result after first loop run: {result_loop_1}")
# Expected: {'count': 3, 'messages': ['Loop Start', 'Processed step 1', 'Processed step 2', 'Processed step 3']}

# Now, let's simulate a partial run and then resume
# We need to modify the graph to *not* end immediately
builder_partial = StateGraph(PersistentAgentState)
builder_partial.add_node("step_A", process_step)
builder_partial.add_node("step_B", process_step)
builder_partial.add_node("step_C", process_step)
builder_partial.set_entry_point("step_A")
builder_partial.add_edge("step_A", "step_B")
# Intentionally stop before C to demonstrate resumption
# We'll use a conditional edge to simulate a breakpoint or partial run
def should_continue(state: PersistentAgentState):
    if state["count"] < 2: # Stop after step_B
        return "continue"
    return "end"

builder_partial.add_conditional_edges(
    "step_B",
    should_continue,
    {"continue": "step_C", "end": END}
)
builder_partial.add_edge("step_C", END)

graph_with_postgres_partial = builder_partial.compile(checkpointer=postgres_checkpoint)

thread_id_partial = "partial_run_789"
config_partial = {"configurable": {"thread_id": thread_id_partial}}

print(f"\n--- Partial Invocation for thread {thread_id_partial} ---")
# Run partially, it should stop after step_B
result_partial_1 = graph_with_postgres_partial.invoke({"count": 0, "messages": ["Partial Start"]}, config=config_partial)
print(f"Result after partial run 1 (should stop at count 2): {result_partial_1}")
# Expected: {'count': 2, 'messages': ['Partial Start', 'Processed step 1', 'Processed step 2']}

# Resume the partial run. It should pick up from where it left off (after step_B)
print(f"\n--- Resuming Partial Invocation for thread {thread_id_partial} ---")
result_partial_2 = graph_with_postgres_partial.invoke(None, config=config_partial)
print(f"Result after resuming partial run (should complete step_C): {result_partial_2}")
# Expected: {'count': 3, 'messages': ['Partial Start', 'Processed step 1', 'Processed step 2', 'Processed step 3']}

ヒューマン・イン・ザ・ループ(HITL)

HITLは、複雑で、リスクが高く、または曖昧なワークフローにとって非常に重要です。LangGraphは、interruptメカニズムを介してこれを容易にします。エージェントは実行を一時停止し、人間の入力を待ってから再開できます。

from langgraph.graph import StateGraph, END
from typing import TypedDict, Literal, Optional
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage

class HumanReviewState(TypedDict):
    chat_history: List[BaseMessage]
    user_query: str
    agent_response: Optional[str]
    human_feedback: Optional[str]
    review_needed: Literal["yes", "no"]

def agent_decides_review(state: HumanReviewState):
    # Simulate agent logic to decide if human review is needed
    if "sensitive" in state["user_query"].lower():
        return {"review_needed": "yes", "agent_response": "I've drafted a response, but it might be sensitive. Awaiting human review."}
    return {"review_needed": "no", "agent_response": "Here's my direct response."}

def generate_response(state: HumanReviewState):
    # Simulate generating a response
    if state["review_needed"] == "no":
        return {"agent_response": f"Agent's direct response to: {state['user_query']}"}
    # If review is needed, the agent_response would have been set by agent_decides_review
    return state

def human_review_node(state: HumanReviewState):
    # This node is where the human would provide feedback.
    # In a real system, this would be an API endpoint or UI interaction.
    # For demonstration, we'll assume feedback is provided externally.
    print(f"\n--- HUMAN REVIEW REQUIRED ---")
    print(f"User Query: {state['user_query']}")
    print(f"Agent's Draft/Decision: {state['agent_response']}")
    print(f"Please provide feedback for thread_id: {config['configurable']['thread_id']}")
    # The graph will be interrupted here.
    # The human would then call the graph with updated state via `update_state` or `invoke` with feedback.
    return state # State remains unchanged until human input

def process_human_feedback(state: HumanReviewState):
    if state["human_feedback"]:
        final_response = f"Agent incorporated human feedback: '{state['human_feedback']}'. Final response: {state['agent_response']}"
        return {"agent_response": final_response, "review_needed": "no"}
    return state

# Build the graph
builder = StateGraph(HumanReviewState)
builder.add_node("decide_review", agent_decides_review)
builder.add_node("generate_response", generate_response)
builder.add_node("human_review", human_review_node)
builder.add_node("process_feedback", process_human_feedback)

builder.set_entry_point("decide_review")

# Conditional edge for review
builder.add_conditional_edges(
    "decide_review",
    lambda state: state["review_needed"],
    {
        "yes": "human_review",
        "no": "generate_response",
    }
)

builder.add_edge("generate_response", END)
builder.add_edge("human_review", "process_feedback")
builder.add_edge("process_feedback", END)

# Compile with a checkpointer and interrupt_before
# Interrupt before 'human_review' node
memory_checkpoint = MemorySaver()
graph_with_hitl = builder.compile(
    checkpointer=memory_checkpoint,
    interrupt_before=["human_review"] # Interrupt before this node
)

# --- Scenario 1: No human review needed ---
thread_id_no_review = "hitl_thread_1"
config = {"configurable": {"thread_id": thread_id_no_review}}
print(f"\n--- Running Scenario 1: No Human Review Needed ({thread_id_no_review}) ---")
result_no_review = graph_with_hitl.invoke(
    {"chat_history": [], "user_query": "What is your return policy?", "review_needed": "no"},
    config=config
)
print(f"Final state (no review): {result_no_review}")
# Expected: {'chat_history': [], 'user_query': 'What is your return policy?', 'agent_response': "Agent's direct response to: What is your return policy?", 'human_feedback': None, 'review_needed': 'no'}

# --- Scenario 2: Human review needed ---
thread_id_review = "hitl_thread_2"
config = {"configurable": {"thread_id": thread_id_review}}
print(f"\n--- Running Scenario 2: Human Review Needed ({thread_id_review}) ---")
initial_state_review = {"chat_history": [], "user_query": "I have a sensitive issue with my account.", "review_needed": "yes"}

# First invocation: will run up to 'human_review' and interrupt
print("First invocation (interrupting at human_review)...")
try:
    # LangGraph raises a StopIteration when interrupted
    graph_with_hitl.invoke(initial_state_review, config=config)
except StopIteration as e:
    print(f"Graph interrupted at: {e.args[0]['configurable']['checkpoint_id']}")
    # The state at interruption can be retrieved from the checkpointer
    interrupted_state = graph_with_hitl.get_state(config)
    print(f"State at interruption: {interrupted_state.current}")

    # Simulate human providing feedback
    human_provided_feedback = "Ensure empathy and offer a direct contact number."
    print(f"\n--- Human provides feedback: '{human_provided_feedback}' ---")

    # Resume the graph with human feedback
    print("Resuming graph with human feedback...")
    final_result_review = graph_with_hitl.invoke(
        {"human_feedback": human_provided_feedback}, # Only provide the delta
        config=config
    )
    print(f"Final state after human review: {final_result_review}")
    # Expected: {'chat_history': [], 'user_query': 'I have a sensitive issue with my account.', 'agent_response': "Agent incorporated human feedback: 'Ensure empathy and offer a direct contact number.'. Final response: I've drafted a response, but it might be sensitive. Awaiting human review.", 'human_feedback': 'Ensure empathy and offer a direct contact number.', 'review_needed': 'no'}

compileのinterrupt_beforeパラメータは、指定されたノードに入る前に実行を一時停止するようにLangGraphに指示します。StopIteration例外はこの一時停止を通知します。その後、人間が入力を行い、更新された状態で再度呼び出すことでグラフを再開できます。

状態のロールバック

チェックポイントにより状態のロールバックが可能になります。エージェントが望ましくない決定を下したり、誤った状態に陥ったりした場合、人間のオペレーターまたは自動化されたシステムが以前のチェックポイントに戻すことができます。

# Assuming graph_with_postgres_partial from above is compiled with PostgresSaver
# thread_id_partial = "partial_run_789"

# Get the history of checkpoints for a thread
history = graph_with_postgres_partial.get_state_history(config_partial)
print(f"\n--- Checkpoint History for {thread_id_partial} ---")
for i, checkpoint in enumerate(history):
    print(f"Checkpoint {i}: {checkpoint.config['configurable']['checkpoint_id']} - State: {checkpoint.state.current}")

# Let's say we want to rollback to the state before the last step
# The history is ordered from oldest to newest.
# To rollback, we need the ID of the checkpoint we want to restore.
# For demonstration, let's assume we want to go back to the state after 'step_B'
# which was the state before the final 'step_C' in the 'partial_run_789' example.
# This would be the second to last checkpoint in the history.

# Get the checkpoint ID of the state we want to restore
# In a real system, you'd have a UI or API to select this.
# For this example, let's manually pick the checkpoint ID from the history.
# Assuming the history has at least 2 checkpoints for 'partial_run_789'
if len(history) >= 2:
    checkpoint_to_restore_id = history[-2].config['configurable']['checkpoint_id']
    print(f"\n--- Rolling back to checkpoint ID: {checkpoint_to_restore_id} ---")

    # To rollback, we invoke with a specific checkpoint_id in the config
    rollback_config = {
        "configurable": {
            "thread_id": thread_id_partial,
            "checkpoint_id": checkpoint_to_restore_id
        }
    }
    # Invoking with None will load the specified checkpoint
    rolled_back_state = graph_with_postgres_partial.invoke(None, config=rollback_config)
    print(f"State after rollback: {rolled_back_state}")
    # Expected: {'count': 2, 'messages': ['Partial Start', 'Processed step 1', 'Processed step 2']}

    # Now, if we run the graph again without a specific checkpoint_id,
    # it will continue from the rolled-back state.
    print(f"\n--- Resuming after rollback ---")
    result_after_rollback = graph_with_postgres_partial.invoke(None, config=config_partial)
    print(f"Result after resuming from rolled-back state: {result_after_rollback}")
    # Expected: {'count': 3, 'messages': ['Partial Start', 'Processed step 1', 'Processed step 2', 'Processed step 3']}
else:
    print("Not enough checkpoints in history to demonstrate rollback for 'partial_run_789'.")

動的サブグラフ

LangGraphは、動的なグラフ構築または変更を可能にし、適応型ワークフローを実現します。これは、コンテキストに基づいてツールやサブエージェントを動的に選択する必要があるエージェントにとって特に役立ちます。実行時に完全な動的グラフ構造の変更は複雑ですが、サブグラフの動的な呼び出しや条件付きルーティングは一般的です。

from langgraph.graph import StateGraph, END, START
from typing import TypedDict, Literal, Optional
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage

class DynamicAgentState(TypedDict):
    chat_history: List[BaseMessage]
    current_task: Optional[str]
    tool_output: Optional[str]
    next_step: Literal["plan", "execute_search", "execute_calculator", "respond", "end"]

# Define some dummy tool nodes
def search_tool_node(state: DynamicAgentState):
    query = state["current_task"]
    print(f"Executing search for: {query}")
    # Simulate API call
    return {"tool_output": f"Search results for '{query}': Found 10 relevant documents."}

def calculator_tool_node(state: DynamicAgentState):
    expression = state["current_task"]
    print(f"Executing calculator for: {expression}")
    # Simulate API call
    try:
        result = eval(expression) # DANGER: Never use eval with untrusted input in production!
        return {"tool_output": f"Calculator result for '{expression}': {result}"}
    except Exception as e:
        return {"tool_output": f"Calculator error: {e}"}

def planner_node(state: DynamicAgentState):
    user_query = state["chat_history"][-1].content if state["chat_history"] else ""
    if "calculate" in user_query.lower() or "math" in user_query.lower():
        return {"current_task": "2 + 2 * 3", "next_step": "execute_calculator"}
    elif "search" in user_query.lower() or "find" in user_query.lower():
        return {"current_task": "latest AI news", "next_step": "execute_search"}
    else:
        return {"next_step": "respond"}

def responder_node(state: DynamicAgentState):
    response = f"Understood. {state.get('tool_output', 'No specific tool output.')} How else can I help?"
    return {"chat_history": state["chat_history"] + [AIMessage(content=response)], "next_step": "end"}

# Build the graph
builder = StateGraph(DynamicAgentState)
builder.add_node("planner", planner_node)
builder.add_node("execute_search", search_tool_node)
builder.add_node("execute_calculator", calculator_tool_node)
builder.add_node("responder", responder_node)

builder.set_entry_point("planner")

# Conditional routing based on planner's decision
builder.add_conditional_edges(
    "planner",
    lambda state: state["next_step"],
    {
        "execute_search": "execute_search",
        "execute_calculator": "execute_calculator",
        "respond": "responder",
    }
)

# After tool execution, go to responder
builder.add_edge("execute_search", "responder")
builder.add_edge("execute_calculator", "responder")
builder.add_edge("responder", END)

graph_dynamic = builder.compile(checkpointer=MemorySaver())

# --- Scenario 1: Search query ---
thread_id_search = "dynamic_thread_search"
config_search = {"configurable": {"thread_id": thread_id_search}}
print(f"\n--- Running Dynamic Scenario 1: Search ({thread_id_search}) ---")
result_search = graph_dynamic.invoke(
    {"chat_history": [HumanMessage(content="Please search for the latest AI news.")], "next_step": "plan"},
    config=config_search
)
print(f"Final state (search): {result_search}")
# Expected: tool_output with search results, chat_history updated.

# --- Scenario 2: Calculator query ---
thread_id_calc = "dynamic_thread_calc"
config_calc = {"configurable": {"thread_id": thread_id_calc}}
print(f"\n--- Running Dynamic Scenario 2: Calculator ({thread_id_calc}) ---")
result_calc = graph_dynamic.invoke(
    {"chat_history": [HumanMessage(content="Can you calculate 5 * 8 + 10?")], "next_step": "plan"},
    config=config_calc
)
print(f"Final state (calculator): {result_calc}")
# Expected: tool_output with calculation result, chat_history updated.

# --- Scenario 3: Direct response ---
thread_id_direct = "dynamic_thread_direct"
config_direct = {"configurable": {"thread_id": thread_direct}}
print(f"\n--- Running Dynamic Scenario 3: Direct Response ({thread_id_direct}) ---")
result_direct = graph_dynamic.invoke(
    {"chat_history": [HumanMessage(content="Hello there!")], "next_step": "plan"},
    config=config_direct
)
print(f"Final state (direct): {result_direct}")
# Expected: tool_output None, chat_history updated with direct response.

この例は、planner_nodeの出力に基づいた動的ルーティングを示しており、動的なサブグラフ実行パスを効果的に作成しています。

Advertisement

アーキテクチャ比較:チェックポインター

機能MemorySaverPostgresSaver
永続性なし(インメモリ)永続的(PostgreSQL)
スケーラビリティ低(単一プロセス)高(データベースバックエンド)
並行性制限あり高(データベースが処理)
ユースケース開発、テスト、一時的なタスク本番、長時間実行ワークフロー、フォールトトレランス
セットアップ簡単PostgreSQLインスタンスとスキーマが必要
データ整合性低(揮発性)高(ACID特性)
ロールバック可能(履歴が保持されている場合)堅牢(チェックポイントID経由)

本番環境での落とし穴とトラブルシューティング

  1. スキーマの進化: 本番環境でTypedDictまたはPydanticの状態スキーマを変更するには、既存のチェックポイントに対する移行戦略が必要です。
    • 落とし穴: デフォルト値のない非オプションフィールドを追加すると、古いチェックポイントの読み込みが失敗します。
    • 修正: 新しいフィールドは常にOptionalとして、またはデフォルト値とともに​​追加してください。必須の新しいフィールドの場合、データベース内の古いチェックポイントを更新するためのデータ移行スクリプトを作成します。
  2. データベース接続管理: PostgresSaverは各操作で新しい接続を作成します。高スループットのアプリケーションでは、これにより接続が枯渇する可能性があります。
    • 落とし穴: PostgreSQLログに「Too many connections」エラーが表示される。
    • 修正: 接続プール(例:SQLAlchemyのcreate_engineをpool_sizeとmax_overflowと共に使用)を実装し、接続オブジェクトをPostgresSaverに渡すか、外部で管理します。接続が適切に閉じられていることを確認してください。
  3. 大規模な状態オブジェクト: 非常に大きなオブジェクト(例:広範なチャット履歴、埋め込みドキュメント)を状態に直接保存すると、パフォーマンスが低下し、データベースの制限を超える可能性があります。
    • 落とし穴: チェックポイントの遅延、大規模なデータベースストレージ、潜在的なJSONBサイズ制限。
    • 修正: 状態には大規模オブジェクトへの参照(例:ID)を保存し、ノードが必要なときに専用のドキュメントストアまたはオブジェクトストレージ(S3、GCS)から実際のデータを取得します。
  4. thread_idでの並行性問題: 複数のインスタンスまたはリクエストが、適切なロックなしに同じthread_idを同時に更新しようとすると、競合状態が発生する可能性があります。
    • 落とし穴: 状態の不整合、更新の喪失。
    • 修正: LangGraphのcheckpointerは、単一のinvoke呼び出しの基本的なアトミック性を処理します。同じスレッドに対する複雑な並行更新の場合、アプリケーション層が適切なロックまたはキューイングメカニズムを実装していることを確認してください。特定のthread_idの更新をシリアライズするためにメッセージキュー(Kafka、RabbitMQ)の使用を検討してください。
  5. 本番環境でのinterrupt_before: 強力ではありますが、ウェブサービスコンテキストでHITLのためにStopIterationのみに依存するのは難しい場合があります。
    • 落とし穴: StopIterationは通常の戻り値ではなく例外です。APIエンドポイントでこれを適切に処理するには、特定の例外処理が必要です。
    • 修正: StopIterationをキャッチし、checkpoint_idと現在の状態を抽出し、クライアントに返すようにAPIを設計します。クライアントはこれを人間に提示し、更新された状態を返送して再開します。再開にはcheckpoint_idが使用されることを確認してください。
  6. 動的サブグラフにおけるeval()のセキュリティ: 動的サブグラフの例で述べたように、動的なコード実行にeval()を使用することは、重大なセキュリティ脆弱性です。
    • 落とし穴: リモートコード実行。
    • 修正: 信頼できない入力でeval()を絶対に使用しないでください。動的なツール選択には、事前に定義されたツールのセットを使用し、文字列名を関数呼び出しにマッピングするか、任意のコード実行が本当に必要な場合(これはまれです)は安全なサンドボックス環境を使用してください。

よくある質問

  1. ヒューマン・イン・ザ・ループのステップで認証と認可をどのように処理しますか? LangGraphの呼び出しをラップするAPIレイヤーを実装します。interrupt_beforeノードに到達すると、APIは現在の状態とcheckpoint_idを返します。UI/クライアントはこれを認証されたユーザーに提示します。ユーザーがフィードバックを提供すると、APIはそれを受け取り、ユーザーを認証し、その特定のthread_idに対してアクションを実行することを承認し、その後、更新された状態とcheckpoint_idでLangGraphを呼び出して再開します。
  2. チェックポイントにMongoDBやRedisのような別のデータベースを使用できますか? LangGraphは、MemorySaverとPostgresSaverをすぐに利用できるように提供しています。他のデータベースの場合、カスタムのBaseCheckpointSaverクラスを実装する必要があります。これには、get、put、およびlistなどのメソッドの実装が含まれます。MongoDBの場合、pymongoを使用してチェックポイントをドキュメントとして保存できます。Redisの場合、状態をJSONにシリアライズしてハッシュまたは文字列として保存できます。
  3. 本番環境でLangGraphワークフローを監視する最良の方法は何ですか? オブザーバビリティツールと統合します。
    • ロギング: ノード内で構造化ロギングを使用して、主要なイベント、状態変更、ツール呼び出しをキャプチャします。
    • トレーシング: LangChainはLangSmithと統合されており、エージェントのステップ、LLM呼び出し、ツール呼び出しの詳細なトレーシングを提供します。これは、複雑なワークフローのデバッグと理解に非常に役立ちます。
    • メトリクス: カスタムメトリクス(例:ノード実行時間、人間による介入の数、エラー率)をPrometheus/GrafanaまたはDatadogのセットアップに送信します。
  4. LangGraphワークフローをバージョン管理するにはどうすればよいですか? グラフ定義をコードとして扱い、バージョン管理(Git)で管理します。デプロイ時には、デプロイされたグラフバージョンが既存のチェックポイントスキーマと互換性があることを確認してください。主要なスキーマ変更の場合、グラフの新しいバージョンをデプロイし、古いスレッドを新しいスキーマに移行するか、古いグラフ定義で古いスレッドを実行する可能性があります。thread_idのバージョン管理スキームを使用するか、チェックポイントメタデータ内にグラフバージョンを保存することを検討してください。
  5. LangGraphアプリケーションが遅いです。パフォーマンスをデバッグするにはどうすればよいですか?
    • LLM呼び出し: 主なボトルネックは、LLMの推論時間であることがよくあります。プロンプトを最適化したり、より高速なモデルを使用したり、一般的なLLM呼び出しのキャッシュを実装したりします。
    • ツール呼び出し: ツール関数をプロファイリングします。外部APIは遅いですか?データベースクエリは最適化されていますか?
    • 状態サイズ: 前述のように、大規模な状態オブジェクトはチェックポイントを遅くする可能性があります。
    • グラフの複雑さ: 多数のノードと条件付きエッジを持つ非常に複雑なグラフはオーバーヘッドが発生する可能性があります。可能な場合は簡素化してください。
    • チェックポインター: データベース(PostgresSaverの場合)が高性能で適切にインデックス付けされていることを確認してください。
    • トレーシング: LangSmithまたは同様のツールを使用して、グラフ実行の最も遅い部分を特定します。
Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement