| name | ai-data-engineering-rag-pipeline |
| description | Build production-grade local RAG pipelines with BM25 search, hierarchical chunking, and retrieval evaluation from Nahid's hands-on roadmap |
| triggers | ["how do I build a RAG pipeline from scratch","implement BM25 baseline search engine","create hierarchical chunking for document retrieval","evaluate retrieval performance with Recall@10","design retrieval contracts and golden datasets","set up local RAG with chunk granularity analysis","build inverted index for document search","implement parent-child metadata linkage for chunks"] |
AI Data Engineering RAG Pipeline
Skill by ara.so — Data Skills collection.
This skill helps you build production-grade local RAG (Retrieval-Augmented Generation) pipelines using the architectural patterns and implementations from Nahid Mahmud's AI & Data Engineering Roadmap. The project focuses on baseline search engines, hierarchical chunking strategies, and rigorous evaluation methodologies.
What This Project Provides
- Day 01: Retrieval contract design, governed corpus creation, golden dataset (
questions.jsonl), and evaluation frameworks
- Day 02: BM25 baseline search engine with inverted index, Okapi BM25 ranking, and Recall@10 evaluation
- Day 03: RAG pipeline with hierarchical chunking (document, section, paragraph levels), deterministic chunk IDs, parent-child metadata linkage, and failure mode analysis
Installation & Setup
Clone the repository:
git clone https://github.com/Nahid-mahmud555/ai-data-engineering-roadmap.git
cd ai-data-engineering-roadmap
Install dependencies (each day has its own requirements):
cd Day_02
pip install -r requirements.txt
cd Day_03
pip install -r requirements.txt
Common dependencies across modules:
numpy - numerical operations
rank-bm25 - BM25 implementation
sentence-transformers - embeddings (Day 03+)
faiss-cpu - vector search (Day 03+)
Project Structure
ai-data-engineering-roadmap/
├── Day_01/ # Retrieval contracts & golden datasets
├── Day_02/ # BM25 baseline search engine
│ ├── baseline_bm25.py
│ ├── corpus/ # Document collection
│ └── questions.jsonl
├── Day_03/ # RAG pipeline with hierarchical chunking
│ ├── pipeline.py
│ ├── corpus/ # Multi-domain documents
│ └── questions.jsonl
Day 02: BM25 Baseline Search Engine
Core Implementation
The BM25 baseline provides a deterministic search engine for retrieval evaluation:
from rank_bm25 import BM25Okapi
import json
corpus_docs = []
doc_ids = []
import os
corpus_dir = "corpus"
for filename in os.listdir(corpus_dir):
if filename.endswith(".txt"):
with open(os.path.join(corpus_dir, filename), 'r') as f:
content = f.read()
corpus_docs.append(content)
doc_ids.append(filename)
tokenized_corpus = [doc.lower().split() for doc in corpus_docs]
bm25 = BM25Okapi(tokenized_corpus)
query = "what is machine learning"
tokenized_query = query.lower().split()
doc_scores = bm25.get_scores(tokenized_query)
import numpy as np
top_n = 10
top_indices = np.argsort(doc_scores)[::-1][:top_n]
results = [(doc_ids[i], corpus_docs[i], doc_scores[i]) for i in top_indices]
for rank, (doc_id, content, score) in enumerate(results, 1):
print(f"Rank {rank}: {doc_id} (Score: {score:.4f})")
print(f"Preview: {content[:200]}...\n")
Recall@10 Evaluation
Evaluate retrieval quality against golden dataset:
import json
with open('questions.jsonl', 'r') as f:
questions = [json.loads(line) for line in f]
def calculate_recall_at_k(bm25_index, tokenized_corpus, doc_ids, questions, k=10):
"""
Calculate Recall@K for retrieval evaluation
questions format: [{"query": "...", "relevant_docs": ["doc1.txt", "doc2.txt"]}]
"""
total_recall = 0
num_queries = len(questions)
for q in questions:
query = q['query']
relevant_docs = set(q['relevant_docs'])
tokenized_query = query.lower().split()
scores = bm25_index.get_scores(tokenized_query)
top_k_indices = np.argsort(scores)[::-1][:k]
retrieved_docs = set([doc_ids[i] for i in top_k_indices])
relevant_retrieved = len(relevant_docs.intersection(retrieved_docs))
recall = relevant_retrieved / len(relevant_docs) if relevant_docs else 0
total_recall += recall
print(f"Query: {query}")
print(f"Recall@{k}: {recall:.2%}")
print(f"Retrieved: ")
()
avg_recall = total_recall / num_queries
()
avg_recall
recall = calculate_recall_at_k(bm25, tokenized_corpus, doc_ids, questions, k=)
Day 03: Hierarchical Chunking RAG Pipeline
Chunk Granularity Strategy
Implement multi-level chunking with parent-child metadata:
import hashlib
import uuid
def generate_deterministic_chunk_id(content, parent_id=None, level="document"):
"""
Generate deterministic chunk ID based on content hash
"""
content_hash = hashlib.sha256(content.encode('utf-8')).hexdigest()[:16]
return f"{level}_{content_hash}"
def hierarchical_chunking(document, doc_id):
"""
Create document, section, and paragraph level chunks
Returns: List of chunk dictionaries with metadata
"""
chunks = []
doc_chunk_id = generate_deterministic_chunk_id(document, level="doc")
chunks.append({
"chunk_id": doc_chunk_id,
"content": document,
"level": "document",
"parent_id": None,
"doc_id": doc_id,
"metadata": {"granularity": "document"}
})
sections = document.split("\n\n")
for sec_idx, section in enumerate(sections):
if len(section.strip()) < 50:
continue
sec_chunk_id = generate_deterministic_chunk_id(section, level="sec")
chunks.append({
"chunk_id": sec_chunk_id,
: section,
: ,
: doc_chunk_id,
: doc_id,
: {
: ,
: sec_idx
}
})
paragraphs = section.split()
para_idx, paragraph (paragraphs):
(paragraph.strip()) < :
para_chunk_id = generate_deterministic_chunk_id(paragraph, level=)
chunks.append({
: para_chunk_id,
: paragraph,
: ,
: sec_chunk_id,
: doc_id,
: {
: ,
: sec_idx,
: para_idx
}
})
chunks
sample_doc =
chunks = hierarchical_chunking(sample_doc, )
chunk chunks:
()
()
()
RAG Pipeline with Vector Search
Combine BM25 with semantic search:
from sentence_transformers import SentenceTransformer
import faiss
import numpy as np
class HybridRAGPipeline:
def __init__(self, embedding_model="all-MiniLM-L6-v2"):
self.embedding_model = SentenceTransformer(embedding_model)
self.chunks = []
self.bm25_index = None
self.faiss_index = None
def index_documents(self, documents):
"""
Index documents with both BM25 and vector embeddings
"""
all_chunks = []
for doc_id, doc_content in documents.items():
chunks = hierarchical_chunking(doc_content, doc_id)
all_chunks.extend(chunks)
self.chunks = all_chunks
tokenized_chunks = [chunk['content'].lower().split() for chunk in all_chunks]
self.bm25_index = BM25Okapi(tokenized_chunks)
chunk_texts = [chunk['content'] for chunk in all_chunks]
embeddings = self.embedding_model.encode(chunk_texts, show_progress_bar=True)
dimension = embeddings.shape[1]
self.faiss_index = faiss.IndexFlatIP(dimension)
faiss.normalize_L2(embeddings)
.faiss_index.add(embeddings)
():
filtered_indices = [i i, chunk (.chunks)
chunk[] == granularity]
filtered_indices:
filtered_indices = (((.chunks)))
tokenized_query = query.lower().split()
bm25_scores = .bm25_index.get_scores(tokenized_query)
bm25_scores_filtered = bm25_scores[filtered_indices]
bm25_scores_filtered.() > :
bm25_scores_filtered = bm25_scores_filtered / bm25_scores_filtered.()
query_embedding = .embedding_model.encode([query])
faiss.normalize_L2(query_embedding)
distances, indices = .faiss_index.search(query_embedding, (.chunks))
vector_scores = distances[]
vector_scores_filtered = vector_scores[filtered_indices]
hybrid_scores = alpha * bm25_scores_filtered + ( - alpha) * vector_scores_filtered
top_k_local = np.argsort(hybrid_scores)[::-][:top_k]
top_k_global = [filtered_indices[i] i top_k_local]
results = []
idx top_k_global:
results.append({
: .chunks[idx],
: (bm25_scores[idx]),
: (vector_scores[idx]),
: (alpha * bm25_scores[idx] + ( - alpha) * vector_scores[idx])
})
results
pipeline = HybridRAGPipeline()
documents = {
: sample_doc,
:
}
pipeline.index_documents(documents)
query =
()
results = pipeline.retrieve(query, top_k=, alpha=, granularity=)
i, result (results, ):
()
()
()
()
Retrieval Evaluation Framework
Evaluate chunking strategies against golden dataset:
def evaluate_chunking_strategy(pipeline, questions, granularity="paragraph", k=10):
"""
Evaluate retrieval performance for specific chunk granularity
"""
metrics = {
"recall": [],
"precision": [],
"mrr": []
}
for q in questions:
query = q['query']
relevant_docs = set(q['relevant_docs'])
results = pipeline.retrieve(query, top_k=k, granularity=granularity)
retrieved_docs = set([r['chunk']['doc_id'] for r in results])
relevant_retrieved = len(relevant_docs.intersection(retrieved_docs))
recall = relevant_retrieved / len(relevant_docs) if relevant_docs else 0
metrics['recall'].append(recall)
precision = relevant_retrieved / len(retrieved_docs) if retrieved_docs else 0
metrics['precision'].append(precision)
for rank, result in enumerate(results, 1):
if result['chunk']['doc_id'] in relevant_docs:
metrics[].append( / rank)
:
metrics[].append()
{
: granularity,
: np.mean(metrics[]),
: np.mean(metrics[]),
: np.mean(metrics[])
}
granularity [, , ]:
eval_results = evaluate_chunking_strategy(
pipeline, questions, granularity=granularity, k=
)
()
()
()
()
Configuration Patterns
Golden Dataset Format (questions.jsonl)
{"query": "what is machine learning", "relevant_docs": ["ml_basics.txt", "ai_intro.txt"]}
{"query": "explain supervised learning algorithms", "relevant_docs": ["ml_basics.txt"]}
{"query": "difference between classification and regression", "relevant_docs": ["ml_basics.txt", "supervised_learning.txt"]}
Chunk Metadata Schema
chunk_schema = {
"chunk_id": "str (deterministic hash)",
"content": "str (actual text content)",
"level": "str (document|section|paragraph)",
"parent_id": "str|None (parent chunk_id)",
"doc_id": "str (source document identifier)",
"metadata": {
"granularity": "str",
"section_index": "int (optional)",
"paragraph_index": "int (optional)",
"custom_fields": "dict (extensible)"
}
}
Common Patterns
Parent-Child Retrieval
Retrieve at fine granularity but return parent context:
def retrieve_with_parent_context(pipeline, query, child_granularity="paragraph", top_k=5):
"""
Retrieve chunks and include parent section context
"""
results = pipeline.retrieve(query, top_k=top_k, granularity=child_granularity)
enriched_results = []
for result in results:
chunk = result['chunk']
parent_id = chunk['parent_id']
parent = next((c for c in pipeline.chunks if c['chunk_id'] == parent_id), None)
enriched_results.append({
"match": chunk['content'],
"context": parent['content'] if parent else chunk['content'],
"score": result['hybrid_score'],
"level": chunk['level']
})
return enriched_results
Failure Mode Analysis
Identify queries with poor retrieval:
def analyze_failure_modes(pipeline, questions, threshold=0.3):
"""
Identify queries where retrieval fails (Recall@10 < threshold)
"""
failures = []
for q in questions:
query = q['query']
relevant_docs = set(q['relevant_docs'])
results = pipeline.retrieve(query, top_k=10)
retrieved_docs = set([r['chunk']['doc_id'] for r in results])
recall = len(relevant_docs.intersection(retrieved_docs)) / len(relevant_docs)
if recall < threshold:
failures.append({
"query": query,
"recall": recall,
"expected": relevant_docs,
"retrieved": retrieved_docs,
"top_result": results[0]['chunk']['content'][:200] if results else None
})
return failures
failures = analyze_failure_modes(pipeline, questions, threshold=0.3)
print(f"\nFound {len(failures)} failure cases:")
for f in failures[:5]:
print(f"\nQuery: {f[]}")
()
()
()
Troubleshooting
BM25 Returns Low Scores
Issue: All BM25 scores are near zero or negative.
Solution: Check tokenization and ensure corpus is properly preprocessed:
sample_query = "machine learning"
tokenized = sample_query.lower().split()
print(f"Tokenized query: {tokenized}")
print(f"First doc tokens: {tokenized_corpus[0][:10]}")
from rank_bm25 import BM25Okapi
bm25_custom = BM25Okapi(tokenized_corpus, k1=1.5, b=0.75)
Empty Chunks After Hierarchical Split
Issue: Some documents produce no paragraph-level chunks.
Solution: Adjust minimum length thresholds and splitting logic:
def hierarchical_chunking(document, doc_id, min_section_len=30, min_para_len=10):
for sec_idx, section in enumerate(sections):
if len(section.strip()) < min_section_len:
continue
FAISS Index Dimension Mismatch
Issue: Error: dimension mismatch when searching FAISS index.
Solution: Ensure query embeddings match index dimension:
print(f"Index dimension: {pipeline.faiss_index.d}")
query_emb = pipeline.embedding_model.encode([query])
print(f"Query embedding shape: {query_emb.shape}")
assert query_emb.shape[1] == pipeline.faiss_index.d, "Dimension mismatch!"
Poor Recall on Golden Dataset
Issue: Recall@10 is unexpectedly low.
Solution: Validate golden dataset format and tune hybrid weights:
with open('questions.jsonl', 'r') as f:
for i, line in enumerate(f, 1):
try:
q = json.loads(line)
assert 'query' in q and 'relevant_docs' in q
except Exception as e:
print(f"Line {i} error: {e}")
for alpha in [0.0, 0.3, 0.5, 0.7, 1.0]:
results = pipeline.retrieve(query, alpha=alpha)
print(f"Alpha={alpha}: Top result from {results[0]['chunk']['doc_id']}")
Environment Variables
The project uses local models by default but can be configured via environment:
export EMBEDDING_MODEL="all-MiniLM-L6-v2"
export CORPUS_DIR="./corpus"
export GOLDEN_DATASET="./questions.jsonl"
export MIN_SECTION_LENGTH="50"
export MIN_PARAGRAPH_LENGTH="20"
Best Practices
- Always validate golden datasets before evaluation runs
- Use deterministic chunk IDs for reproducibility and debugging
- Experiment with granularity levels based on query complexity
- Tune hybrid weights (alpha) for your specific domain
- Maintain parent-child linkage for context expansion
- Log failure modes to identify corpus or chunking issues
- Normalize embeddings before FAISS indexing for cosine similarity
References