Building a RAG Pipeline
The Complete RAG Pipeline
You understand what RAG is and how embeddings work. Now you will build a complete, production-ready RAG pipeline end to end. This covers document ingestion, chunking, embedding, storage, retrieval, and response generation.
We will build a customer support knowledge base that answers questions based on company documentation.
Phase 1: Document Ingestion
The first phase processes your source documents and loads them into the vector database.
Step 1: Load and parse documents
// Run in Node.js: npm install openai chromadb
import fs from 'fs';
import path from 'path';
function loadDocuments(directory) {
const docs = [];
const files = fs.readdirSync(directory);
for (const file of files) {
if (!file.endsWith('.txt') && !file.endsWith('.md')) continue;
const content = fs.readFileSync(path.join(directory, file), 'utf-8');
docs.push({
id: file.replace(/[^a-z0-9]/gi, '_'),
filename: file,
content,
metadata: {
source: file,
category: inferCategory(file),
created_at: fs.statSync(path.join(directory, file)).mtime.toISOString(),
}
});
}
return docs;
}
function inferCategory(filename) {
const lower = filename.toLowerCase();
if (lower.includes('refund') || lower.includes('return')) return 'returns';
if (lower.includes('ship') || lower.includes('delivery')) return 'shipping';
if (lower.includes('account') || lower.includes('login')) return 'account';
if (lower.includes('product')) return 'products';
return 'general';
}
Step 2: Chunk documents
function chunkDocument(doc, chunkSize = 400, overlap = 50) {
const words = doc.content.split(/\s+/).filter(w => w.length > 0);
const chunks = [];
let chunkIndex = 0;
for (let i = 0; i < words.length; i += (chunkSize - overlap)) {
const chunkWords = words.slice(i, i + chunkSize);
if (chunkWords.length < 30) continue; // Skip very small final chunks
chunks.push({
id: doc.id + '_chunk_' + chunkIndex,
text: chunkWords.join(' '),
metadata: {
...doc.metadata,
chunk_index: chunkIndex,
parent_doc: doc.filename,
word_count: chunkWords.length,
}
});
chunkIndex++;
}
return chunks;
}
Step 3: Generate embeddings and store
import OpenAI from 'openai';
import { ChromaClient } from 'chromadb';
const openai = new OpenAI({ apiKey: process.env.OPENAI_API_KEY });
const chroma = new ChromaClient({ path: 'http://localhost:8000' });
async function ingestDocuments(documents) {
// Get or create collection
let collection;
try {
collection = await chroma.getCollection({ name: 'support_kb' });
console.log('Using existing collection');
} catch {
collection = await chroma.createCollection({ name: 'support_kb' });
console.log('Created new collection');
}
// Process in batches to avoid rate limits
const BATCH_SIZE = 50;
let totalChunks = 0;
for (const doc of documents) {
const chunks = chunkDocument(doc);
for (let i = 0; i < chunks.length; i += BATCH_SIZE) {
const batch = chunks.slice(i, i + BATCH_SIZE);
// Generate embeddings for the batch
const embeddingResponse = await openai.embeddings.create({
model: 'text-embedding-3-small',
input: batch.map(c => c.text),
});
// Store in ChromaDB
await collection.add({
ids: batch.map(c => c.id),
embeddings: embeddingResponse.data.map(e => e.embedding),
documents: batch.map(c => c.text),
metadatas: batch.map(c => c.metadata),
});
totalChunks += batch.length;
console.log('Ingested', totalChunks, 'chunks...');
// Rate limit pause
if (i + BATCH_SIZE < chunks.length) {
await new Promise(r => setTimeout(r, 200));
}
}
}
console.log('Ingestion complete:', totalChunks, 'total chunks stored');
return collection;
}
Phase 2: Query and Retrieval
async function retrieveContext(collection, query, options = {}) {
const {
topK = 5,
minimumScore = 0.5,
categoryFilter = null,
} = options;
// Embed the query
const queryEmbedding = await openai.embeddings.create({
model: 'text-embedding-3-small',
input: query,
});
// Build metadata filter
const whereClause = categoryFilter ? { category: categoryFilter } : undefined;
// Search the vector database
const results = await collection.query({
queryEmbeddings: [queryEmbedding.data[0].embedding],
nResults: topK,
where: whereClause,
include: ['documents', 'metadatas', 'distances'],
});
// Convert distances to similarity scores and filter
const chunks = results.documents[0]
.map((doc, i) => ({
text: doc,
metadata: results.metadatas[0][i],
score: 1 - results.distances[0][i],
}))
.filter(chunk => chunk.score >= minimumScore);
return chunks;
}
Phase 3: Response Generation
async function generateAnswer(query, retrievedChunks) {
if (retrievedChunks.length === 0) {
return {
answer: "I could not find relevant information in our knowledge base for that question. Please contact our support team for assistance.",
sources: [],
confidence: 'low',
};
}
// Format context for the prompt
const context = retrievedChunks
.map((chunk, i) => `[Source ${i + 1}: ${chunk.metadata.source}]\n${chunk.text}`)
.join('\n\n---\n\n');
const response = await openai.chat.completions.create({
model: 'gpt-4o-mini',
temperature: 0,
messages: [
{
role: 'system',
content: `You are a helpful customer support assistant. Answer questions using ONLY the provided documentation.
Rules:
- Only use information from the provided sources
- If the answer is not in the sources, say "I don't have that information in our documentation"
- Cite which source(s) your answer comes from
- Be concise and direct
- Never speculate or add information not in the sources`
},
{
role: 'user',
content: `Documentation:\n${context}\n\nQuestion: ${query}`
}
],
max_tokens: 500,
});
const sources = [...new Set(retrievedChunks.map(c => c.metadata.source))];
const confidence = retrievedChunks[0].score > 0.8 ? 'high' : retrievedChunks[0].score > 0.6 ? 'medium' : 'low';
return {
answer: response.choices[0].message.content,
sources,
confidence,
tokensUsed: response.usage.total_tokens,
};
}
Phase 4: Complete Pipeline
async function answerQuestion(query, collection) {
console.log('\nQuery:', query);
// Step 1: Retrieve relevant context
const chunks = await retrieveContext(collection, query, {
topK: 5,
minimumScore: 0.4,
});
console.log('Retrieved', chunks.length, 'relevant chunks');
// Step 2: Generate answer
const result = await generateAnswer(query, chunks);
console.log('Answer:', result.answer);
console.log('Sources:', result.sources);
console.log('Confidence:', result.confidence);
return result;
}
// Usage
const collection = await chroma.getCollection({ name: 'support_kb' });
await answerQuestion('What is your refund policy for digital products?', collection);
await answerQuestion('How long does international shipping take?', collection);
Testing Your RAG Pipeline
Metric 1: Retrieval Quality
For each test question, manually verify that the retrieved chunks actually contain the answer:
- Precision: Of the chunks retrieved, what percentage are relevant?
- Recall: Does the retrieved set contain the answer?
Metric 2: Answer Quality
Compare the generated answer to the ground truth answer:
- Is the answer factually correct?
- Does it stay within the boundaries of the retrieved context?
- Is it cited properly?
Metric 3: Answer When No Answer Exists
Ask questions outside the knowledge base scope. The system should respond "I don't have that information" rather than hallucinating.
Building a Golden Test Set
Create 20-50 question-answer pairs from your knowledge base. Run every change to your pipeline against this set. Track metrics over time to catch regressions.
Common RAG Failures and Fixes
Low retrieval precision (wrong documents retrieved):
- Improve chunking (try different sizes)
- Use metadata filters to narrow search scope
- Add hybrid search (keyword + semantic)
Answer hallucination despite retrieved context:
- Strengthen the system prompt constraint ("use ONLY the provided sources")
- Lower temperature to 0
- Add explicit "if not in sources, say so" instruction
Missing answers (answer exists but not retrieved):
- Lower the minimum score threshold
- Increase topK retrieval count
- Improve query preprocessing (expand query, remove stop words)
Key Takeaways
- A complete RAG pipeline has four phases: document ingestion (load, chunk, embed, store), retrieval (embed query, similarity search), response generation (augmented prompt), and evaluation.
- Batch your embedding API calls (50 documents at a time) to avoid rate limits and reduce cost.
- Always filter by minimum score to avoid including low-relevance chunks that confuse the model.
- Build a golden test set of 20-50 question-answer pairs to systematically evaluate pipeline quality after every change.
- The system prompt must strongly constrain the model to use only retrieved sources to prevent hallucination.
Try it yourself
Key Takeaways
- The complete RAG pipeline has four phases: ingestion (load, chunk, embed, store), retrieval (embed query, search), generation (augmented prompt), and evaluation.
- Batch embedding API calls to manage rate limits and cost -- process 50 chunks at a time with brief pauses between batches.
- Filter retrieved chunks by minimum similarity score to prevent low-relevance content from confusing the model.
- Build a golden test set of manually verified Q&A pairs to systematically evaluate pipeline quality and catch regressions.
- The system prompt must strongly constrain the model to use only retrieved sources -- this is the primary defence against hallucination in RAG systems.
Quick Quiz
1.What is the correct order of steps in a RAG ingestion pipeline?
2.Why should you filter retrieved chunks by a minimum similarity score?
3.What should a RAG system respond when no relevant chunks are found for a question?
4.What is a 'golden test set' in the context of RAG evaluation?
Ready to go further?
CareerEx gives you structured 12-week training, live classes every Saturday and Sunday, real tutor feedback, and a certificate. Join the next cohort.
Join CareerEx