RAGとエージェントの中核機能の実装
このトラック全体で学んだパターンに従い、ドキュメント取り込みパイプライン、ベクトルストアのインデックス作成、リランキング付き検索、エージェントのツール統合を構築します
「RAGとエージェントの中核機能の実装」はCoddyKit上の無料AI Engineering Academyレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはAI Engineering Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 AI Engineering Academyコースには全4レッスンが含まれています。
実装の順序と依存関係
システムはボトムアップで構築します。外部依存関係のないコンポーネントから始め、それらに依存するコンポーネントを上に積み重ねます。RAG + agentシステムの場合、順序は次のとおりです。(1) vector store schema、(2) ingestion pipeline、(3) retriever、(4) basic Q&A chain、(5) streaming endpoint、(6) agent with tools、(7) caching layer、(8) tracing instrumentation。各コンポーネントをパイプラインに統合する前に、単体でテストしてください。
# Build order:
IMPL_ORDER = [
'pgvector_schema', # prerequisite for everything
'document_ingestion', # populate the vector store
'hybrid_retriever', # test retrieval in isolation
'qa_chain_basic', # integrate LLM with retrieval
'streaming_endpoint', # expose via API
'function_calling', # add agent tool calls
'semantic_cache', # reduce repeat API calls
'langsmith_tracing', # add after core works
'injection_filter', # harden before load testing
]ベクトルストアのセットアップ
ドキュメントを取り込む前に、適切なスキーマでpgvectorテーブルを作成します。ベクトル、チャンクテキスト、すべてのメタデータフィールド、選択的な再インデックスに使うupdated_atタイムスタンプの列を含めてください。embedding列にはすぐにベクトルインデックス(HNSWまたはIVFFlat)を作成します。数百万行を格納した後にインデックスを追加するより、空のテーブルにあらかじめ追加するほうがはるかに高速です。
-- PostgreSQL schema with pgvector
CREATE EXTENSION IF NOT EXISTS vector;
CREATE TABLE document_chunks (
id BIGSERIAL PRIMARY KEY,
doc_id TEXT NOT NULL,
chunk_text TEXT NOT NULL,
embedding VECTOR(1536) NOT NULL,
source_file TEXT,
page_number INT,
section TEXT,
doc_type TEXT,
tenant_id TEXT NOT NULL, -- for data isolation
created_at TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX idx_chunks_hnsw ON document_chunks
USING hnsw (embedding vector_cosine_ops)
WITH (m = 16, ef_construction = 64);
CREATE INDEX idx_chunks_tenant ON document_chunks(tenant_id);取り込みパイプラインの構築
取り込み処理は、ファイルパスまたはURLを受け取り、インデックス化したチャンク数を返す単一のasync関数として実装します。ファイル種別ごとにLangChainのdocument loadersを使用し、チャンク分割にはrecursive character text splitterを使用します。embeddingの呼び出しはバッチ化して、入力の2048トークン制限内に収めるとともに、API呼び出しを数千回から数十回に減らします。
from langchain_community.document_loaders import PyPDFLoader
from langchain.text_splitter import RecursiveCharacterTextSplitter
from openai import AsyncOpenAI
async def ingest_document(file_path: str, doc_id: str, tenant_id: str) -> int:
loader = PyPDFLoader(file_path)
pages = loader.load()
splitter = RecursiveCharacterTextSplitter(chunk_size=800, chunk_overlap=100)
chunks = splitter.split_documents(pages)
# Batch embed
texts = [c.page_content for c in chunks]
client = AsyncOpenAI()
embeddings_response = await client.embeddings.create(
model='text-embedding-3-small',
input=texts
)
embeddings = [e.embedding for e in embeddings_response.data]
await insert_chunks_to_pgvector(chunks, embeddings, doc_id, tenant_id)
return len(chunks)ハイブリッドretrieverの構築
密ベクトル検索とBM25キーワード検索を組み合わせ、reciprocal rank fusionを使って結果を統合します。retrieverは、単一のretrieve(query, tenant_id, top_k)メソッドを持つクラスとして実装します。内部ではasyncio.gatherを使って両方の検索を同時実行し、RRFでランキング済みリストを統合し、チャンクIDで重複を除去して、ソースメタデータ付きの上位top_k件を返します。
import asyncio
from rank_bm25 import BM25Okapi
class HybridRetriever:
def __init__(self, pool, k_rrf: int = 60):
self.pool = pool
self.k_rrf = k_rrf
async def retrieve(self, query: str, tenant_id: str, top_k: int = 10) -> list:
dense_results, sparse_results = await asyncio.gather(
self._dense_search(query, tenant_id, top_k * 3),
self._bm25_search(query, tenant_id, top_k * 3)
)
merged = self._rrf_merge(dense_results, sparse_results)
return merged[:top_k]
def _rrf_score(self, rank: int) -> float:
return 1.0 / (self.k_rrf + rank + 1)コアQAチェーンの実装
LangChain LCELを使ってコアQ&Aチェーンを構築します。チェーンは質問と取得したチャンクを受け取り、提供されたコンテキストだけを使うよう指示した拡張プロンプトを作成し、応答をストリーミングします。ドキュメント名とページ番号で出典を示すこと、取得したコンテキストに答えがない場合は「I don't know」と答えることを、モデルに明示的に指示してください。
from langchain.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain.schema.output_parser import StrOutputParser
RAG_PROMPT = ChatPromptTemplate.from_messages([
('system', 'You are a precise assistant. Answer ONLY using the provided context. '
'Cite sources as [DocName, p.N]. If the answer is not in the context, say "I don\'t have information about that."'),
('user', 'Context:\n{context}\n\nQuestion: {question}')
])
llm = ChatOpenAI(model='gpt-4o', temperature=0, streaming=True)
qa_chain = RAG_PROMPT | llm | StrOutputParser()
async def answer_question(question: str, chunks: list) -> str:
context = '\n\n'.join(f'[{c["source"]}]\n{c["text"]}' for c in chunks)
return await qa_chain.ainvoke({'context': context, 'question': question})エージェントへのFunction Callingの追加
リアルタイムデータや計算を必要とするクエリに対応するため、function callingを使って基本的なQ&Aシステムを拡張します。現在の情報を検索するWeb検索、安全な読み取り専用データベースに対してSQLクエリを実行する処理、IDで特定のレコードを検索する処理のツールを定義してください。エージェントは質問に基づいて呼び出すツールを決定し、実行して、その結果を最終回答に組み込みます。
from openai import AsyncOpenAI
import json
TOOLS = [
{
'type': 'function',
'function': {
'name': 'search_knowledge_base',
'description': 'Search the internal document knowledge base for relevant information',
'parameters': {
'type': 'object',
'properties': {
'query': {'type': 'string', 'description': 'Search query'},
'top_k': {'type': 'integer', 'default': 5}
},
'required': ['query']
}
}
}
]
async def agent_with_tools(question: str, tenant_id: str) -> str:
client = AsyncOpenAI()
messages = [{'role': 'user', 'content': question}]
while True:
resp = await client.chat.completions.create(
model='gpt-4o', messages=messages, tools=TOOLS)
if resp.choices[0].finish_reason != 'tool_calls':
return resp.choices[0].message.content
tool_call = resp.choices[0].message.tool_calls[0]
args = json.loads(tool_call.function.arguments)
result = await dispatch_tool(tool_call.function.name, args, tenant_id)
messages.append({'role': 'tool', 'tool_call_id': tool_call.id, 'content': result})エージェント応答のストリーミング
エージェントをサーバー送信イベント対応の FastAPI StreamingResponse でラップし、フロントエンドにトークンが届くたびに表示されるようにします。ツール呼び出しを含むエージェントの応答では、ツールの実行中に進捗インジケーター(「ナレッジベースを検索中…」)をストリームし、その後、最終回答をトークン単位でストリームします。これにより、ツール実行のレイテンシー中にページがフリーズしたように見えるのを防げます。
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import asyncio
app = FastAPI()
async def event_stream(question: str, tenant_id: str):
yield 'data: {"type": "start"}\n\n'
chunks = await retriever.retrieve(question, tenant_id)
yield 'data: {"type": "retrieving", "count": ' + str(len(chunks)) + '}\n\n'
async for token in qa_chain.astream({'context': format_context(chunks), 'question': question}):
yield f'data: {{"type": "token", "content": {repr(token)}}}\n\n'
yield 'data: {"type": "done"}\n\n'
@app.post('/query/stream')
async def stream_query(question: str, tenant_id: str):
return StreamingResponse(event_stream(question, tenant_id), media_type='text/event-stream')Semantic Cache の組み込み
クエリパイプラインの検索前ステップとして semantic cache を追加します。retriever と LLM にアクセスする前に、ユーザーの質問を埋め込み、キャッシュを確認します。類似度がしきい値を超えるキャッシュヒットがあった場合は、cached: true フラグを付けてキャッシュ済みの回答を直ちに返します。キャッシュミスの場合は完全なパイプラインを実行し、生成された回答を今後の類似クエリに備えてキャッシュに保存します。
async def query_pipeline(question: str, tenant_id: str) -> dict:
# 1. Check semantic cache
cache_hit = await semantic_cache.lookup(question, tenant_id, threshold=0.92)
if cache_hit:
return {'answer': cache_hit.answer, 'cached': True, 'sources': cache_hit.sources}
# 2. Retrieve
chunks = await retriever.retrieve(question, tenant_id, top_k=5)
chunks = await reranker.rerank(question, chunks, top_n=3)
# 3. Generate
answer = await answer_question(question, chunks)
sources = [c['source'] for c in chunks]
# 4. Cache result
await semantic_cache.store(question, tenant_id, answer, sources)
return {'answer': answer, 'cached': False, 'sources': sources}LangSmith トレーシングの追加
2つの環境変数を設定して、パイプラインを LangSmith で計測します。すべての LangChain chain の呼び出しが、トークン数、レイテンシー、入力、出力、エラーの有無とともに自動的にトレースされます。カスタムの非 LangChain コード(検索や再ランキングなど)については、呼び出しを @traceable デコレーターでラップしてトレースに含めます。これにより、パイプラインの各ステップをエンドツーエンドで完全に把握できます。
import os
from langsmith import traceable
os.environ['LANGCHAIN_TRACING_V2'] = 'true'
os.environ['LANGCHAIN_API_KEY'] = os.environ['LANGSMITH_API_KEY']
os.environ['LANGCHAIN_PROJECT'] = 'document-qa-production'
# Wrap non-LangChain steps with @traceable
@traceable(name='hybrid_retrieval')
async def traced_retrieval(question: str, tenant_id: str, top_k: int) -> list:
return await retriever.retrieve(question, tenant_id, top_k)
@traceable(name='cohere_reranking')
async def traced_reranking(question: str, chunks: list) -> list:
return await reranker.rerank(question, chunks)
# LangChain LCEL chains are automatically traced — no extra code needed統合テストの実行
質問の入力から最終回答まで、クエリパイプライン全体を実行する統合テストを作成します。既知のドキュメントを格納した小規模なテスト用ベクトルストアを使用すると、取得されるチャンクや回答に含まれるべき内容について決定論的なアサーションを記述できます。本番環境のインフラストラクチャを再現しつつ、別のベクトルストアと LLM キーを使用するステージング環境に対して統合テストを実行します。
import pytest
@pytest.mark.asyncio
async def test_full_pipeline_returns_grounded_answer():
# Setup: ingest known document
await ingest_document('tests/fixtures/policy.pdf', 'policy_v1', 'test_tenant')
# Query with a question that has a known answer in the document
result = await query_pipeline(
question='What is the cancellation policy?',
tenant_id='test_tenant'
)
assert result['answer'] is not None
assert len(result['answer']) > 50
assert '24 hours' in result['answer'] # known fact in document
assert 'policy.pdf' in str(result['sources'])
assert result['cached'] is False # fresh query検索品質のベースライン測定
何かを最適化する前に、検索のベースラインを確立します。評価用テストセットを使って、ヒット率(正しいドキュメントが上位5件の結果に含まれるか)と MRR(何位にランクインするか)を測定します。このベースラインは、初期ベクトルストアを設定した後、再ランキングやハイブリッド検索を追加する前に実行します。ベースラインによって、各最適化が実際にどの程度の改善をもたらしたかが分かるため、採用すべき手法を根拠に基づいて判断できます。
async def measure_retrieval_baseline(test_cases: list) -> dict:
hits = 0
reciprocal_ranks = []
for case in test_cases:
results = await retriever.retrieve(case['question'], case['tenant_id'], top_k=5)
result_docs = [r['doc_id'] for r in results]
if case['relevant_doc'] in result_docs:
hits += 1
rank = result_docs.index(case['relevant_doc']) + 1
reciprocal_ranks.append(1.0 / rank)
else:
reciprocal_ranks.append(0.0)
return {
'hit_rate_at_5': hits / len(test_cases),
'mrr': sum(reciprocal_ranks) / len(reciprocal_ranks)
}クイックチェック
RAG とエージェントの中核機能の構築について、理解度を確認します。
レッスンのまとめ
このレッスンでは、ボトムアップの実装順序により、統合前に各コンポーネントを分離してテストできること、密な検索と BM25 検索を並行して実行するハイブリッド検索を RRF で組み合わせると、最も高い検索品質が得られること、そして@traceable デコレーターを使った LangSmith トレーシングによってパイプライン全体を完全に可視化できることを学びました。次は、セキュリティ、キャッシュ、信頼性のパターンを導入してシステムを堅牢化します。
よくある質問
「RAGとエージェントの中核機能の実装」レッスンは無料ですか?
はい。「RAGとエージェントの中核機能の実装」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、AI Engineering Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 AI Engineering Academyコースには全4レッスンが含まれています。
「RAGとエージェントの中核機能の実装」で何を学びますか?
このトラック全体で学んだパターンに従い、ドキュメント取り込みパイプライン、ベクトルストアのインデックス作成、リランキング付き検索、エージェントのツール統合を構築します ブラウザで直接実行するハンズオンコードでAI Engineering Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
AI Engineering Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのAI Engineering Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。
「RAGとエージェントの中核機能の実装」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このAI Engineering Academyレッスンでコードを書いて実行できますか?
はい。すべてのAI Engineering Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- 本番アーキテクチャの設計
- RAGとエージェントの中核機能の実装
- 堅牢化:セキュリティ、キャッシュ、信頼性
- 評価、デプロイ、振り返り