はじめに
RAG エンジンは、ドキュメント プロセッサ、埋め込みプロバイダー、およびベクター ストアを完全なパイプラインに接続するオーケストレーターです。この記事では、取り込み、取得、再ランキングを備えた完全な RAG エンジンを実装します。
1. RAG エンジン クラス
// packages/core/src/rag/rag-engine.ts
export class RagEngine {
private docProcessor: DocumentProcessor;
private embedding: EmbeddingProvider;
private vectorStore: VectorStore;
private webCrawler: WebCrawler;
constructor(
docProcessor: DocumentProcessor,
embedding: EmbeddingProvider,
vectorStore: VectorStore,
) {
this.docProcessor = docProcessor;
this.embedding = embedding;
this.vectorStore = vectorStore;
this.webCrawler = new WebCrawler();
}
// Ingest plain text
async ingestText(
text: string,
collectionId: string,
metadata?: Record<string, unknown>,
): Promise<{ chunksCreated: number }> {
const chunks = await this.docProcessor.process(text, 'text/plain', {
source: 'direct-text',
sourceType: 'text',
chunkSize: 1000,
chunkOverlap: 200,
});
const vectors = await this.embedding.embed(chunks.map(c => c.content));
const documents: VectorDocument[] = chunks.map((chunk, i) => ({
id: chunk.id,
vector: vectors[i],
content: chunk.content,
metadata: { ...chunk.metadata, ...metadata, collectionId },
}));
await this.vectorStore.upsert(documents);
return { chunksCreated: chunks.length };
}
// Ingest from URL
async ingestUrl(
url: string,
collectionId: string,
): Promise<{ chunksCreated: number }> {
const text = await this.webCrawler.crawl(url);
const chunks = await this.docProcessor.process(text, 'text/html', {
source: url,
sourceType: 'url',
chunkSize: 1000,
chunkOverlap: 200,
});
const vectors = await this.embedding.embed(chunks.map(c => c.content));
const documents: VectorDocument[] = chunks.map((chunk, i) => ({
id: chunk.id,
vector: vectors[i],
content: chunk.content,
metadata: { ...chunk.metadata, collectionId },
}));
await this.vectorStore.upsert(documents);
return { chunksCreated: chunks.length };
}
}
2. 検索と再ランキング
// Retrieve relevant documents
async retrieve(
query: string,
collectionId: string,
options: { topK?: number; minScore?: number } = {},
): Promise<SearchResult[]> {
const [queryVector] = await this.embedding.embed([query]);
const results = await this.vectorStore.search(queryVector, {
collection: collectionId,
topK: options.topK || 5,
minScore: options.minScore || 0.7,
filter: { 'metadata.collectionId': collectionId },
});
return results;
}
// Search with reranking — better accuracy
async searchWithReranking(
query: string,
collectionId: string,
options: { topK?: number; rerankTopK?: number } = {},
): Promise<SearchResult[]> {
// Step 1: Retrieve more candidates
const candidates = await this.retrieve(query, collectionId, {
topK: options.rerankTopK || 20,
minScore: 0.5, // Lower threshold for candidates
});
if (candidates.length === 0) return [];
// Step 2: Rerank using cross-encoder or LLM
const reranked = await this.rerankResults(query, candidates);
// Step 3: Return top K
return reranked.slice(0, options.topK || 5);
}
private async rerankResults(
query: string,
candidates: SearchResult[],
): Promise<SearchResult[]> {
// Simple reranking: score by keyword overlap + vector score
return candidates
.map(result => {
const queryTerms = query.toLowerCase().split(/\s+/);
const contentLower = result.content.toLowerCase();
const keywordScore = queryTerms.filter(t => contentLower.includes(t)).length / queryTerms.length;
const combinedScore = result.score * 0.7 + keywordScore * 0.3;
return { ...result, score: combinedScore };
})
.sort((a, b) => b.score - a.score);
}
3. コレクションの管理
// Collection CRUD
async createCollection(name: string, tenantId: string): Promise<string> {
const collectionId = crypto.randomUUID();
await this.vectorStore.createCollection(collectionId);
// Save metadata to MongoDB
await this.db.collection('rag_collections').insertOne({
_id: collectionId,
name,
tenantId,
documentCount: 0,
chunkCount: 0,
createdAt: new Date(),
});
return collectionId;
}
async listCollections(tenantId: string) {
return this.db.collection('rag_collections')
.find({ tenantId })
.sort({ createdAt: -1 })
.toArray();
}
async deleteCollection(collectionId: string) {
await this.vectorStore.delete(
await this.getChunkIds(collectionId),
);
await this.db.collection('rag_collections').deleteOne({ _id: collectionId });
}
4. エージェント フローの RAG
// Khi Agent nhận được message, RAG context được inject:
const ragContext = await ragEngine.searchWithReranking(
userMessage,
activeCollectionId,
{ topK: 5 },
);
if (ragContext.length > 0) {
messages.push({
role: 'system',
content: [
'### Relevant Knowledge Base Context',
'',
...ragContext.map((r, i) =>
`**[${i + 1}]** (score: ${r.score.toFixed(2)})\n${r.content}`
),
'',
'Use the above context to answer the user\'s question. Cite sources using [1], [2], etc.',
].join('\n'),
});
}
5. まとめ
| コンポーネント | 役割 |
|---|---|
| ドキュメントプロセッサ | ドキュメントの解析とチャンク |
| 埋め込みプロバイダー | テキスト → ベクトル |
| VectorStore | ベクトルの保存と検索 |
| RAG エンジン | パイプラインを調整する |
| 再ランキング | 検索精度の向上 |
次の記事: ワークフロー エンジン — AI タスクのビジュアル自動化。