•22 min read

ステートフルなAgentic RAG: グラフステートマシン、自己修正ループ、フォールバックルーティング

ステートフルなAgentic RAG: グラフステートマシン、自己修正ループ、フォールバックルーティング

このガイドでは、グラフステートマシンを活用した、回復力のあるステートフルなAgentic RAGシステムの構築について詳しく説明します。ここでは、素朴な単一ショットの検索パイプラインと、動的なマルチホップルーティング、クエリ分解、反復評価を対比させます。コアとなる実装は、LLMグレーダーでドキュメントの関連性を検証し、ハルシネーション時にウェブ検索フォールバックをトリガーし、PostgresとLangGraphを使用して会話状態のチェックポイントを維持する自己修正型検索ノードを特徴としています。

Audio Briefing
0:00 / 0:00

素朴なRAGの限界

従来のRAG実装は、しばしば単純なパターンに従います。

  1. ユーザーのクエリを受信する。
  2. クエリを埋め込む。
  3. ベクターストアから上位k個のドキュメントを検索する。
  4. ドキュメントをクエリと連結する。
  5. LLMを使用して応答を生成する。

このアプローチは脆いです。以下のことを前提としています。

  • 初期クエリは検索に完全に適している。
  • ベクターストアには必要な情報がすべて含まれている。
  • 検索されたドキュメントは常に適切で十分である。
  • 情報が不足しているか、関連性がない場合でもLLMはハルシネーションを起こさない。

現実世界のシナリオでは、これらの前提は成り立ちません。複雑なクエリには分解が必要です。情報が不足している場合は、外部ツールの使用が必要になります。関連性のないドキュメントは、質の悪い応答やハルシネーションにつながります。ステートフルなAgentic RAGは、動的な制御フロー、反復的な洗練、明示的な自己修正メカニズムを導入することで、これらの欠点を解決します。

Advertisement

アーキテクチャの概要:グラフステートマシン

私たちのアーキテクチャは、LangGraphを使用して実装されたグラフステートマシンを中心に構築されています。グラフ内の各ノードは、個別の処理ステップまたは決定ポイントを表します。状態はノード間で明示的に渡され、複雑な多ターン対話と反復的な洗練を可能にします。

主要なコンポーネントは次のとおりです。

  • 状態定義: 会話および処理状態を定義するPydanticモデル。
  • ルーターノード: 現在のクエリと状態に基づいて次のアクションを決定します(例:直接検索、クエリ分解、ウェブ検索)。
  • 検索ノード: ベクターストアのルックアップを実行します。
  • LLMグレーダーノード: 検索されたドキュメントとクエリの関連性を評価します。
  • ウェブ検索ノード: フォールバックとして外部ウェブ検索を実行します。
  • 応答生成ノード: 最終的な回答を合成します。
  • 状態永続化: グラフ状態のチェックポイントをPostgresに保存し、長時間の会話と回復を可能にします。
# rag_graph/state.py
from typing import List, Optional, Literal
from langchain_core.documents import Document
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage
from langgraph.graph import StateGraph, END
from pydantic import BaseModel, Field

class AgentState(BaseModel):
    """
    Represents the state of our RAG agent.
    This state is passed between nodes in the graph.
    """
    query: str = Field(description="The original user query.")
    chat_history: List[BaseMessage] = Field(default_factory=list, description="Full chat history.")
    documents: List[Document] = Field(default_factory=list, description="Retrieved documents.")
    generation: Optional[str] = Field(None, description="Generated LLM response.")
    retrieval_attempts: int = Field(0, description="Number of retrieval attempts.")
    web_search_performed: bool = Field(False, description="Flag indicating if web search was performed.")
    # Add a field to track the current decision path for debugging/logging
    current_path: List[str] = Field(default_factory=list, description="Path taken through the graph.")

    class Config:
        arbitrary_types_allowed = True # Allow BaseMessage

コアコンポーネントと自己修正ループ

1. クエリールーター

ルーターは動的な動作のエントリーポイントです。ユーザーのクエリを分析し、最初のアクションを決定します。これには以下が含まれます。

  • 直接検索: クエリが単純で、内部知識ベースでカバーされている可能性が高い場合。
  • クエリ分解: 複雑な複数部分の質問の場合、それらをサブクエリに分解します。(簡潔にするため、この例では完全に実装されていませんが、一般的な拡張機能です)。
  • ウェブ検索: クエリが明らかに内部知識ベースの範囲外である場合(例:「ロンドンの天気は?」)。
# rag_graph/nodes.py
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnablePassthrough
from langchain_openai import ChatOpenAI
from langchain_community.vectorstores import FAISS
from langchain_openai import OpenAIEmbeddings
from langchain_community.tools import DuckDuckGoSearchRun
from langchain_core.output_parsers import StrOutputParser
from typing import List, Dict, Any

# Assume these are initialized globally or passed in
# For production, use environment variables for API keys and proper vector store setup
llm = ChatOpenAI(model="gpt-4o", temperature=0)
embeddings = OpenAIEmbeddings()
vectorstore = FAISS.from_texts(["LangGraph is a library for building stateful, multi-actor applications with LLMs.",
                                "RAG stands for Retrieval Augmented Generation.",
                                "Self-correction in RAG improves accuracy.",
                                "Postgres can be used for state persistence."], embeddings)
retriever = vectorstore.as_retriever()
web_search_tool = DuckDuckGoSearchRun()

# --- Prompts ---
ROUTER_PROMPT = ChatPromptTemplate.from_messages([
    ("system", "You are a smart routing agent. Your goal is to determine the best next step for a user query."),
    ("human", """Given the user query and chat history, decide whether to:
    1. 'retrieve': Perform a standard RAG retrieval from our internal knowledge base.
    2. 'web_search': Perform a web search for external information.
    3. 'generate': Directly generate a response if the query is simple and doesn't require retrieval or search.

    Respond with only one of the keywords: 'retrieve', 'web_search', 'generate'.

    Chat History: {chat_history}
    User Query: {query}
    """)
])

RETRIEVAL_GRADER_PROMPT = ChatPromptTemplate.from_messages([
    ("system", "You are a document relevance grader. Your task is to assess if the retrieved documents are relevant to the user's query."),
    ("human", """User Query: {query}
    Retrieved Documents: {documents}

    Are the retrieved documents relevant to the user's query?
    Respond with 'yes' or 'no'.
    """)
])

HALLUCINATION_GRADER_PROMPT = ChatPromptTemplate.from_messages([
    ("system", "You are a hallucination grader. Your task is to assess if the generated answer is grounded in the provided documents."),
    ("human", """User Query: {query}
    Retrieved Documents: {documents}
    Generated Answer: {generation}

    Is the generated answer fully supported by the retrieved documents?
    Respond with 'yes' or 'no'.
    """)
])

ANSWER_GENERATION_PROMPT = ChatPromptTemplate.from_messages([
    ("system", "You are an AI assistant. Use the following retrieved context to answer the user's question. If the context does not contain the answer, state that you don't know."),
    ("human", """Context: {documents}
    Chat History: {chat_history}
    Question: {query}
    """)
])

# --- Nodes ---
def route_query(state: AgentState) -> str:
    """Decides the next step based on the query."""
    print("---ROUTE QUERY---")
    state.current_path.append("route_query")
    router_chain = ROUTER_PROMPT | llm | StrOutputParser()
    decision = router_chain.invoke({"query": state.query, "chat_history": state.chat_history})
    print(f"Router decision: {decision}")
    if "web_search" in decision.lower():
        return "web_search"
    elif "retrieve" in decision.lower():
        return "retrieve"
    else: # Default to generate if not explicitly retrieve or web_search
        return "generate"

def retrieve(state: AgentState) -> AgentState:
    """Retrieves documents from the vector store."""
    print("---RETRIEVE DOCUMENTS---")
    state.current_path.append("retrieve")
    state.retrieval_attempts += 1
    documents = retriever.invoke(state.query)
    state.documents = documents
    return state

def grade_documents(state: AgentState) -> str:
    """Grades the relevance of retrieved documents."""
    print("---GRADE DOCUMENTS---")
    state.current_path.append("grade_documents")
    if not state.documents:
        print("No documents retrieved, initiating web search.")
        return "no_documents"

    grader_chain = RETRIEVAL_GRADER_PROMPT | llm | StrOutputParser()
    decision = grader_chain.invoke({"query": state.query, "documents": state.documents})
    print(f"Document relevance decision: {decision}")
    if "yes" in decision.lower():
        print("Documents are relevant.")
        return "relevant"
    else:
        print("Documents are not relevant, initiating web search.")
        return "not_relevant"

def web_search(state: AgentState) -> AgentState:
    """Performs a web search and adds results to documents."""
    print("---WEB SEARCH---")
    state.current_path.append("web_search")
    state.web_search_performed = True
    web_results = web_search_tool.invoke({"query": state.query})
    # Convert web results to Document objects
    state.documents.extend([Document(page_content=web_results, metadata={"source": "web_search"})])
    return state

def generate_response(state: AgentState) -> AgentState:
    """Generates a final answer based on retrieved documents and query."""
    print("---GENERATE RESPONSE---")
    state.current_path.append("generate_response")
    # Format documents for the prompt
    docs_content = "\n\n".join([doc.page_content for doc in state.documents])

    generation_chain = ANSWER_GENERATION_PROMPT | llm | StrOutputParser()
    response = generation_chain.invoke({
        "query": state.query,
        "documents": docs_content,
        "chat_history": state.chat_history
    })
    state.generation = response
    return state

def grade_generation_for_hallucination(state: AgentState) -> str:
    """Grades the generated answer for hallucination against retrieved documents."""
    print("---GRADE GENERATION FOR HALLUCINATION---")
    state.current_path.append("grade_generation_for_hallucination")
    if not state.documents: # If no documents were used, we can't grade against them
        print("No documents to grade against, assuming no hallucination for now.")
        return "no_hallucination"

    grader_chain = HALLUCINATION_GRADER_PROMPT | llm | StrOutputParser()
    decision = grader_chain.invoke({
        "query": state.query,
        "documents": state.documents,
        "generation": state.generation
    })
    print(f"Hallucination decision: {decision}")
    if "yes" in decision.lower():
        print("Generation is grounded in documents.")
        return "no_hallucination"
    else:
        print("Generation contains hallucination, attempting web search fallback.")
        return "hallucination"

def update_chat_history(state: AgentState) -> AgentState:
    """Updates the chat history with the latest query and response."""
    print("---UPDATE CHAT HISTORY---")
    state.current_path.append("update_chat_history")
    state.chat_history.append(HumanMessage(content=state.query))
    if state.generation:
        state.chat_history.append(AIMessage(content=state.generation))
    return state

2. 自己修正型検索ノード

これは重要なコンポーネントです。最初の検索後、LLMが「グレーダー」として機能し、ドキュメントの関連性を評価します。

  • 関連性がある場合: 応答生成に進みます。
  • 関連性がない場合: ウェブ検索などのフォールバックメカニズムをトリガーします。これにより、LLMが不適切または不足している内部データに基づいて回答を生成するのを防ぎます。

このループは拡張可能です。ウェブ検索も失敗した場合、システムはユーザーに明確化を求めたり、人間にエスカレートしたりすることができます。

3. ハルシネーション検出とフォールバック

生成後、別のLLMグレーダーが、生成された応答を検索されたドキュメントと照合して評価します。これは重要な自己修正ステップです。

  • 根拠がある場合: 応答は有効と見なされ、ユーザーに返されます。
  • ハルシネーションの場合: システムはウェブ検索をトリガーするか(まだ実行されていない場合)、修正されたクエリで検索を再試行できます。この反復的な洗練により、信頼性が大幅に向上します。

4. Postgresによる状態永続化

LangGraphは状態永続化の組み込みサポートを提供します。堅牢でスケーラブルなチェックポイントのためにPostgresを使用します。これにより、以下が可能になります。

  • 長時間の会話: ユーザーは数日後に会話に戻ることができます。
  • 障害からの回復: エージェントプロセスがクラッシュした場合でも、状態を再ロードできます。
  • デバッグと監査: 完全な状態履歴が利用可能です。
# rag_graph/graph.py
from langgraph.checkpoint.sqlite import SqliteSaver # For local testing
from langgraph.checkpoint.postgres import PostgresSaver # For production
from langgraph.graph import StateGraph, END
from rag_graph.state import AgentState
from rag_graph.nodes import (
    route_query, retrieve, grade_documents, web_search,
    generate_response, grade_generation_for_hallucination, update_chat_history
)
import os

# For production, configure PostgresSaver
# memory = PostgresSaver.from_conn_string(os.environ["POSTGRES_CONNECTION_STRING"])
# For local testing, use SqliteSaver
memory = SqliteSaver.from_conn_string(":memory:") # In-memory SQLite for quick testing

def build_graph():
    workflow = StateGraph(AgentState)

    # Define nodes
    workflow.add_node("retrieve", retrieve)
    workflow.add_node("grade_documents", grade_documents)
    workflow.add_node("web_search", web_search)
    workflow.add_node("generate_response", generate_response)
    workflow.add_node("grade_generation_for_hallucination", grade_generation_for_hallucination)
    workflow.add_node("update_chat_history", update_chat_history)

    # Set entry point
    workflow.set_entry_point("route_query")

    # Define edges
    workflow.add_conditional_edges(
        "route_query",
        route_query,
        {
            "retrieve": "retrieve",
            "web_search": "web_search",
            "generate": "generate_response" # Direct generation for simple queries
        }
    )

    workflow.add_edge("retrieve", "grade_documents")

    workflow.add_conditional_edges(
        "grade_documents",
        grade_documents,
        {
            "relevant": "generate_response",
            "not_relevant": "web_search",
            "no_documents": "web_search" # If retrieval yielded nothing, try web search
        }
    )

    workflow.add_edge("web_search", "generate_response") # After web search, always try to generate

    workflow.add_edge("generate_response", "grade_generation_for_hallucination")

    workflow.add_conditional_edges(
        "grade_generation_for_hallucination",
        grade_generation_for_hallucination,
        {
            "no_hallucination": "update_chat_history",
            "hallucination": "web_search" # If hallucination, try web search (if not already done)
        }
    )

    workflow.add_edge("update_chat_history", END)

    # Compile the graph
    app = workflow.compile(checkpointer=memory)
    return app

# Example usage (in a separate script or main block)
if __name__ == "__main__":
    app = build_graph()

    # Example 1: Simple RAG query
    print("\n--- Running Example 1: Simple RAG ---")
    config = {"configurable": {"thread_id": "1"}}
    inputs = {"query": "What is LangGraph?", "chat_history": []}
    for s in app.stream(inputs, config=config):
        print(s)
    final_state = app.get_state(config)
    print(f"\nFinal Answer (Thread 1): {final_state.values['generation']}")
    print(f"Path taken (Thread 1): {final_state.values['current_path']}")

    # Example 2: Query requiring web search (e.g., current events)
    print("\n--- Running Example 2: Web Search Fallback ---")
    config = {"configurable": {"thread_id": "2"}}
    inputs = {"query": "What is the capital of France?", "chat_history": []} # Assume internal KB doesn't have this
    for s in app.stream(inputs, config=config):
        print(s)
    final_state = app.get_state(config)
    print(f"\nFinal Answer (Thread 2): {final_state.values['generation']}")
    print(f"Path taken (Thread 2): {final_state.values['current_path']}")

    # Example 3: Query that might lead to hallucination or irrelevant docs
    print("\n--- Running Example 3: Hallucination/Irrelevant Docs ---")
    config = {"configurable": {"thread_id": "3"}}
    inputs = {"query": "Tell me about the latest advancements in quantum computing, specifically related to cold fusion.", "chat_history": []}
    for s in app.stream(inputs, config=config):
        print(s)
    final_state = app.get_state(config)
    print(f"\nFinal Answer (Thread 3): {final_state.values['generation']}")
    print(f"Path taken (Thread 3): {final_state.values['current_path']}")

アーキテクチャのトレードオフ

機能素朴なRAGステートフルなAgentic RAG (グラフステートマシン)
複雑性低高 (グラフ定義、状態管理、複数のLLM呼び出し)
堅牢性低 (ハルシネーション、不十分な検索に陥りやすい)高 (自己修正、フォールバック、反復的な洗練)
柔軟性低 (固定パイプライン)高 (動的ルーティング、ノードの追加/削除が容易、カスタムロジック)
コスト (LLM呼び出し)低 (クエリあたり1-2回の呼び出し)高 (ルーティング、グレーディング、生成、場合によってはループのための複数のLLM呼び出し)
レイテンシ低高 (複数のシーケンシャルLLM呼び出し、ツール使用)
状態管理なし (クエリごとにステートレス)明示的 (永続化された状態、多ターン会話)
デバッグシンプル複雑 (グラフ実行のトレース、状態遷移)
スケーラビリティステートレスコンポーネントのスケーリングが容易状態永続化には堅牢なデータベースが必要、グラフ実行はスレッドごとに並列化可能
Advertisement

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

  1. LLMグレーダーのドリフト: LLMグレーダー(関連性、ハルシネーション用)のパフォーマンスは、新しいLLMバージョンやプロンプトの変更によって低下する可能性があります。
    • 修正: 継続的な評価を実装します。クエリ、ドキュメント、期待されるグレーダー出力のゴールデンデータセットを維持します。毎晩テストを実行し、重大な逸脱があった場合はアラートを発します。LLMモデルのバージョンを固定します。
  2. グラフ内の無限ループ: 不適切に設計された条件付きエッジは、ノードが互いを繰り返し呼び出すことにつながる可能性があります(例:ウェブ検索が役に立たない場合、grade_documents -> web_search -> grade_documents)。
    • 修正: 状態内にループ検出と最大再試行回数を実装します。例えば、state.retrieval_attemptsやstate.web_search_performedフラグは冗長なアクションを防ぎます。LangGraphのmax_stepsも設定できます。
  3. 状態スキーマの進化: AgentState Pydanticモデルが進化するにつれて、Postgres内の既存のチェックポイントが互換性を失う可能性があります。
    • 修正: スキーマ移行を計画します。データベーススキーマの変更にはAlembicのようなツールを使用します。LangGraphのチェックポイントについては、状態のバージョン管理を行うか、古いチェックポイントの移行戦略(例:ロード、変換、保存)を検討してください。
  4. ツールのレイテンシとレート制限: ウェブ検索やその他の外部ツールは、かなりのレイテンシを導入したり、APIレート制限に達したりする可能性があります。
    • 修正: 頻繁に検索される用語のキャッシュを実装します。外部ツールには非同期呼び出しを使用します。指数関数的バックオフを伴う堅牢な再試行メカニズムを実装します。ツールの使用状況を監視し、適切なレート制限を設定します。
  5. コンテキストウィンドウのオーバーフロー: あまりにも多くのドキュメントを連結すると(特にウェブ検索後)、LLMのコンテキストウィンドウを超える可能性があります。
    • 修正: 最終的な生成ステップに渡す前に、インテリジェントなドキュメント要約または再ランキングを実装します。関連性スコアに基づいてドキュメントを優先します。ドキュメントを適切に切り詰めます。
  6. コスト超過: クエリあたりのLLM呼び出し回数が増えると、コストが急速に上昇する可能性があります。
    • 修正: トークン効率のためにプロンプトを最適化します。より単純なタスクには安価なモデルを使用します(例:ルーティング/グレーディングに十分であればgpt-3.5-turbo)。適切な場合はLLM応答のキャッシュを実装します。会話あたりのトークン使用量を監視します。
  7. 複雑なグラフパスのデバッグ: 複数ノードのグラフを介した実行フローのトレースは困難な場合があります。
    • 修正: AgentStateにcurrent_path: List[str]フィールドを追加して、訪問したノードのシーケンスをログに記録します。オブザーバビリティツール(例:LangSmith、OpenTelemetry)と統合して、グラフ実行を視覚化します。各ノード内に詳細なログを追加します。

よくある質問

Q1: マルチターン会話を処理し、コンテキストを効果的に維持するにはどうすればよいですか?

A1: chat_historyフィールドはAgentStateで非常に重要です。各ターンで、ユーザーのクエリとAIの応答が追加されます。応答を生成する際、LLMはこの履歴でプロンプトされ、進行中の会話を理解できるようになります。LangGraphの状態永続化により、この履歴はセッション間で維持されます。

Q2: ルーターとジェネレーターなど、異なるノードに適切なLLMを選択するための最良の戦略は何ですか?

A2: 階層的なアプローチを使用します。ルーティングや初期グレーディングのような単純でリスクの低いタスクには、より高速で安価なモデル(例:gpt-3.5-turbo、Llama 3 8B)で十分かもしれません。複雑な推論、要約、または最終的な回答の生成には、より高性能で高価なモデル(例:gpt-4o、Claude 3 Opus)がしばしば好まれます。各ノードの特定のタスクについて、異なるモデルをベンチマークします。

Q3: ウェブ検索以外のより複雑なツール(内部APIやデータベースなど)を統合するにはどうすればよいですか?

A3: 各ツールは独自のLangGraphノード内にカプセル化できます。route_queryノードを拡張して、どのツールを呼び出すかを決定できます。たとえば、クエリが「顧客の注文状況」を尋ねる場合、ルーターは内部APIを呼び出すorder_lookupノードに指示できます。ツールが後続のノードで簡単に利用できる形式(例:Documentオブジェクト)でデータを返すようにしてください。

Q4: LangGraphは高スループット、低レイテンシの本番環境に適していますか?

A4: LangGraphは、複雑なAgenticワークフローのための堅牢なフレームワークを提供します。そのパフォーマンスは、基盤となるLLM呼び出しと外部ツールのレイテンシに大きく依存します。高スループットの場合、LLM呼び出しを最適化し(キャッシュ、より小さなモデル)、非同期操作を使用し、状態永続化層(Postgres)が非常に高性能であることを確認してください。グラフの実行自体は効率的ですが、I/O操作が通常ボトルネックになります。可能な場合はリクエストのバッチ処理を検討してください。

Q5: RAGシステムが検索されたドキュメントから機密情報を漏洩させないようにするにはどうすればよいですか?

A5: ドキュメントが生成のためにLLMに渡される前に、堅牢なPII検出と編集を実装します。これはグラフ内の専用ノードにすることができます。さらに、ベクターストアと検索メカニズムが適切なアクセス制御で構成されていることを確認してください。非常に機密性の高いデータの場合、より小さくプライベートなLLMをファインチューニングするか、差分プライバシーなどの手法を使用することを検討してください。

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