
Ai Data Engineering
- 50 installs
- 426 repo stars
- Updated December 11, 2025
- ancoleman/ai-design-components
AI Data Engineering is a Claude Code skill that provides architecture patterns for RAG pipelines, feature stores, and embedding generation in AI/ML systems.
About
AI Data Engineering is a Claude Code skill for building data infrastructure for AI/ML systems, including RAG pipelines, feature stores, and embedding generation. It provides a 5-stage RAG architecture, chunking strategy defaults, embedding-provider comparisons, RAGAS evaluation metrics, and feature-store patterns using Feast. A developer uses it when building retrieval-augmented generation, semantic search, or ML feature serving for production applications.
- 5-stage RAG pipeline architecture: ingestion, indexing, retrieval, generation, evaluation
- Chunking (512 tokens/50-100 overlap) and embedding provider guidance (Voyage AI, OpenAI)
- RAGAS evaluation metrics and feature-store patterns (Feast) with orchestration tools
Ai Data Engineering by the numbers
- 50 all-time installs (skills.sh)
- Ranked #929 of 2,064 Data Science & ML skills by installs in the Skillselion catalog
- Data as of Aug 5, 2026 (Skillselion catalog sync)
ai-data-engineering capabilities & compatibility
- Capabilities
- architecting data · administering linux · assembling components
- Works with
- openai · postgres · redis · snowflake · databricks
- Use cases
- data analysis · research · database · orchestration
- Pricing
- Bring your own API key
What ai-data-engineering says it does
Data pipelines, feature stores, and embedding generation for AI/ML systems.
RAG pipelines have 5 distinct stages. Understanding this architecture is critical for production implementations.
Chunking is the most critical decision for RAG quality. Poor chunking breaks retrieval.
npx skills add https://github.com/ancoleman/ai-design-components --skill ai-data-engineeringAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 50 |
|---|---|
| repo stars | ★ 426 |
| Last updated | December 11, 2025 |
| Repository | ancoleman/ai-design-components ↗ |
What it does
Build RAG pipelines, embedding generation, and feature stores for production AI/ML applications.
Who is it for?
Developers building RAG, semantic search, or ML feature serving for production.
Skip if: Traditional CRUD apps, simple key-value storage, or applications with no AI/ML component.
When should I use this skill?
You are building a RAG pipeline, an embedding pipeline, or a feature store.
What you get
A staged RAG pipeline with tuned chunking, chosen embeddings, RAGAS evaluation, and a feature store.
By the numbers
- 5-stage RAG pipeline
- 4 core RAGAS metrics
- default 512-token chunks with 50-100 token overlap
Files
AI Data Engineering
Purpose
Build data infrastructure for AI/ML systems including RAG pipelines, feature stores, and embedding generation. Provides architecture patterns, orchestration workflows, and evaluation metrics for production AI applications.
When to Use
Use this skill when:
- Building RAG (Retrieval-Augmented Generation) pipelines
- Implementing semantic search or vector databases
- Setting up ML feature stores for real-time serving
- Creating embedding generation pipelines
- Evaluating RAG quality with RAGAS metrics
- Orchestrating data workflows for AI systems
- Integrating with frontend skills (ai-chat, search-filter)
Skip this skill if:
- Building traditional CRUD applications (use databases-relational)
- Simple key-value storage (use databases-nosql)
- No AI/ML components in the application
RAG Pipeline Architecture
RAG pipelines have 5 distinct stages. Understanding this architecture is critical for production implementations.
┌─────────────────────────────────────────────────────────────┐
│ RAG Pipeline (5 Stages) │
├─────────────────────────────────────────────────────────────┤
│ │
│ 1. INGESTION → Load documents (PDF, DOCX, Markdown) │
│ 2. INDEXING → Chunk (512 tokens) + Embed + Store │
│ 3. RETRIEVAL → Query embedding + Vector search + Filters │
│ 4. GENERATION → Context injection + LLM streaming │
│ 5. EVALUATION → RAGAS metrics (faithfulness, relevancy) │
│ │
└─────────────────────────────────────────────────────────────┘For complete RAG architecture with implementation patterns, see:
references/rag-architecture.md- Detailed 5-stage breakdownexamples/langchain-rag/basic_rag.py- Working implementation
Chunking Strategies
Chunking is the most critical decision for RAG quality. Poor chunking breaks retrieval.
Default Recommendation:
- Size: 512 tokens
- Overlap: 50-100 tokens
- Method: Fixed token-based
Why these values:
- Too small (<256 tokens): Loses context, requires many retrievals
- Too large (>1024 tokens): Includes irrelevant content, hits token limits
- Overlap prevents information loss at chunk boundaries
Alternative strategies for special cases:
# Code-aware chunking (preserves functions/classes)
from langchain.text_splitter import RecursiveCharacterTextSplitter
code_splitter = RecursiveCharacterTextSplitter.from_language(
language="python",
chunk_size=512,
chunk_overlap=50
)
# Semantic chunking (splits on meaning, not tokens)
from langchain.text_splitter import SemanticChunker
semantic_splitter = SemanticChunker(
embeddings=embeddings,
breakpoint_threshold_type="percentile" # Split at semantic boundaries
)See: references/chunking-strategies.md for complete decision framework
Embedding Generation
Embedding quality directly impacts retrieval accuracy. Voyage AI is currently best-in-class.
Primary Recommendation: Voyage AI voyage-3
- Dimensions: 1024
- MTEB Score: 69.0 (highest as of Dec 2025)
- Cost: $$$ but 9.74% better than OpenAI
- Use for: Production systems requiring best retrieval quality
Cost-Effective Alternative: OpenAI text-embedding-3-small
- Dimensions: 1536
- MTEB Score: 62.3
- Cost: $ (5x cheaper than voyage-3)
- Use for: Development, prototyping, cost-sensitive applications
Implementation:
from langchain_voyageai import VoyageAIEmbeddings
from langchain_openai import OpenAIEmbeddings
# Production (best quality)
embeddings = VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key="your-api-key"
)
# Development (cost-effective)
embeddings = OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key="your-api-key"
)See: references/embedding-strategies.md for complete provider comparison
RAGAS Evaluation Metrics
Traditional metrics (BLEU, ROUGE) don't measure RAG quality. RAGAS provides LLM-as-judge evaluation.
4 Core Metrics:
| Metric | Measures | Good Score |
|---|---|---|
| Faithfulness | Factual consistency with retrieved context | > 0.8 |
| Answer Relevancy | Does answer address the user's question? | > 0.7 |
| Context Precision | Are retrieved chunks actually relevant? | > 0.6 |
| Context Recall | Were all necessary chunks retrieved? | > 0.7 |
Quick evaluation script:
# Run RAGAS evaluation (TOKEN-FREE script execution)
python scripts/evaluate_rag.py --dataset eval_data.json --output results.jsonManual implementation:
from ragas import evaluate
from ragas.metrics import faithfulness, answer_relevancy
dataset = {
"question": ["What is the capital of France?"],
"answer": ["Paris is the capital of France."],
"contexts": [["France's capital is Paris."]],
"ground_truth": ["Paris"]
}
result = evaluate(dataset, metrics=[faithfulness, answer_relevancy])
print(f"Faithfulness: {result['faithfulness']}")
print(f"Answer Relevancy: {result['answer_relevancy']}")See: references/evaluation-metrics.md for complete RAGAS implementation guide
Feature Stores
Feature stores solve the "training-serving skew" problem by providing consistent feature computation.
Primary Recommendation: Feast - Open source, works with any backend (PostgreSQL, Redis, DynamoDB, S3, BigQuery, Snowflake)
Basic usage:
from feast import FeatureStore
store = FeatureStore(repo_path="feature_repo/")
# Online serving (low-latency)
features = store.get_online_features(
features=["user_features:total_orders"],
entity_rows=[{"user_id": 1001}]
).to_dict()See: references/feature-stores.md for complete Feast setup and alternatives (Tecton, Hopsworks)
LangChain Orchestration
LangChain is the primary framework for LLM orchestration with the largest ecosystem (24,215+ API reference snippets).
Context7 Library ID: /websites/langchain_oss_python_langchain (Trust: High, Snippets: 435)
Basic RAG Chain:
from langchain_core.prompts import ChatPromptTemplate
from langchain_qdrant import QdrantVectorStore
from langchain_voyageai import VoyageAIEmbeddings
# Setup retriever
vectorstore = QdrantVectorStore(
client=qdrant_client,
embedding=VoyageAIEmbeddings(model="voyage-3")
)
retriever = vectorstore.as_retriever(search_type="mmr", search_kwargs={"k": 5})
# Build chain
prompt = ChatPromptTemplate.from_template(
"Answer based on context:\n{context}\n\nQuestion: {question}"
)
chain = {"context": retriever, "question": lambda x: x} | prompt | ChatOpenAI() | StrOutputParser()
# Stream response
for chunk in chain.stream("What is the capital of France?"):
print(chunk, end="", flush=True)See: references/langchain-patterns.md - Complete LangChain 0.3+ patterns with streaming and hybrid search
Orchestration Tools
Modern AI pipelines require workflow orchestration beyond cron jobs.
Primary Recommendation: Dagster (for ML/AI pipelines) - Asset-centric design, best lineage tracking, perfect for RAG
Example: Embedding Pipeline
from dagster import asset
from langchain_voyageai import VoyageAIEmbeddings
@asset
def raw_documents():
"""Load documents from S3."""
return documents
@asset
def chunked_documents(raw_documents):
"""Split into 512-token chunks with 50-token overlap."""
from langchain.text_splitter import RecursiveCharacterTextSplitter
splitter = RecursiveCharacterTextSplitter(chunk_size=512, chunk_overlap=50)
return splitter.split_documents(raw_documents)
@asset
def embedded_documents(chunked_documents):
"""Generate embeddings with Voyage AI."""
embeddings = VoyageAIEmbeddings(model="voyage-3")
return embeddings.embed_documents([doc.page_content for doc in chunked_documents])See: references/orchestration-tools.md for complete Dagster patterns and alternatives (Prefect, Airflow 3.0, dbt)
Integration with Frontend Skills
ai-chat Skill → RAG Backend
The ai-chat skill consumes RAG pipeline outputs for streaming responses.
Backend API (FastAPI):
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
@app.post("/api/rag/stream")
async def stream_rag(query: str):
async def generate():
chain = RetrievalQA.from_chain_type(llm=OpenAI(streaming=True), retriever=vectorstore.as_retriever())
async for chunk in chain.astream(query):
yield chunk
return StreamingResponse(generate(), media_type="text/plain")See: references/rag-architecture.md for complete frontend integration patterns
search-filter Skill → Semantic Search
The search-filter skill uses semantic search backends for vector similarity.
Backend (Qdrant + Voyage AI):
from qdrant_client import QdrantClient
from langchain_voyageai import VoyageAIEmbeddings
@app.post("/api/search/semantic")
async def semantic_search(query: str, filters: dict):
query_vector = VoyageAIEmbeddings(model="voyage-3").embed_query(query)
results = QdrantClient().search(
collection_name="documents",
query_vector=query_vector,
query_filter=filters,
limit=10
)
return {"results": results}Data Versioning
Primary Recommendation: LakeFS (acquired DVC team November 2025)
Git-like operations on data lakes: branch, commit, merge, time travel. Works with S3/Azure/GCS.
import lakefs
branch = lakefs.Branch("main").create("experiment-voyage-3")
branch.commit("Updated embeddings to voyage-3")
branch.merge_into("main")See: references/data-versioning.md for complete LakeFS setup
Quick Start Workflow
1. Set up vector database:
# Run Qdrant setup script (TOKEN-FREE execution)
python scripts/setup_qdrant.py --collection docs --dimension 10242. Chunk and embed documents:
# Chunk documents (TOKEN-FREE execution)
python scripts/chunk_documents.py \
--input data/documents/ \
--chunk-size 512 \
--overlap 50 \
--output data/chunks/3. Implement RAG pipeline:
See examples/langchain-rag/basic_rag.py for complete working example.
4. Evaluate with RAGAS:
# Run evaluation (TOKEN-FREE execution)
python scripts/evaluate_rag.py \
--dataset data/eval_qa.json \
--output results/ragas_metrics.json5. Deploy with orchestration:
See examples/dagster-pipelines/embedding_pipeline.py for production deployment.
Dependencies
Required Python packages:
# Core RAG
pip install langchain langchain-core langchain-openai langchain-voyageai langchain-qdrant
# Vector database
pip install qdrant-client
# Evaluation
pip install ragas datasets
# Feature stores
pip install feast
# Orchestration
pip install dagster dagster-webserver
# Data versioning
pip install lakefs-clientOptional for alternatives:
# LlamaIndex (alternative to LangChain)
pip install llama-index
# dbt (SQL transformations)
pip install dbt-core dbt-postgres
# Prefect (alternative orchestration)
pip install prefectTroubleshooting
Common Issues:
1. Poor retrieval quality - Check chunk size (try 512 tokens), increase overlap (50-100), try hybrid search, re-rank with Cohere
2. Slow embedding generation - Batch documents (100-1000), use async APIs, cache with Redis, use smaller model for dev
3. High LLM costs - Reduce retrieved chunks (k=3), use cheaper re-ranking models, cache frequent queries
See: references/rag-architecture.md for complete troubleshooting guide
Best Practices
Chunking: Default to 512 tokens with 50-token overlap. Use semantic chunking for complex documents. Preserve code structure for source code.
Embeddings: Use Voyage AI voyage-3 for production, OpenAI text-embedding-3-small for development. Never mix embedding models (re-embed everything if changing).
Evaluation: Run RAGAS metrics on every pipeline change. Maintain test dataset of 50+ question-answer pairs. Track metrics over time.
Orchestration: Use Dagster for ML/AI pipelines, dbt for SQL transformations only. Version control all pipeline code.
Frontend Integration: Always stream LLM responses. Implement retry logic. Show citations/sources to users. Handle empty results gracefully.
Additional Resources
Reference Documentation:
references/rag-architecture.md- Complete RAG pipeline guidereferences/chunking-strategies.md- Decision framework for chunkingreferences/embedding-strategies.md- Embedding model comparisonreferences/langchain-patterns.md- LangChain 0.3+ patternsreferences/feature-stores.md- Feast setup and alternativesreferences/evaluation-metrics.md- RAGAS implementation guide
Working Examples:
examples/langchain-rag/basic_rag.py- Simple RAG chainexamples/langchain-rag/streaming_rag.py- Streaming responsesexamples/langchain-rag/hybrid_search.py- Vector + BM25examples/llamaindex-agents/query_engine.py- LlamaIndex alternativeexamples/feast-features/- Complete feature store setupexamples/dagster-pipelines/embedding_pipeline.py- Production pipeline
Executable Scripts (TOKEN-FREE):
scripts/evaluate_rag.py- RAGAS evaluation runnerscripts/chunk_documents.py- Document chunking utilityscripts/benchmark_retrieval.py- Retrieval quality benchmarkscripts/setup_qdrant.py- Qdrant collection setup
"""
Dagster Pipeline for Embedding Generation
This example demonstrates production-grade orchestration for RAG pipelines:
1. Document ingestion from various sources
2. Chunking with configurable strategies
3. Embedding generation with batching
4. Vector database indexing
5. Asset lineage tracking and versioning
Use Dagster when:
- Building production RAG systems requiring orchestration
- Need to track data lineage and dependencies
- Want to schedule and monitor pipeline runs
- Require asset versioning and recomputation
"""
import os
from typing import List, Dict, Any
from datetime import datetime
from pathlib import Path
from dagster import (
asset,
AssetExecutionContext,
MaterializeResult,
MetadataValue,
Config,
OpExecutionContext,
Definitions,
define_asset_job,
ScheduleDefinition,
AssetSelection,
)
from dotenv import load_dotenv
# LangChain imports
from langchain_core.documents import Document
from langchain_voyageai import VoyageAIEmbeddings
from langchain_openai import OpenAIEmbeddings
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain_qdrant import QdrantVectorStore
from qdrant_client import QdrantClient, models
# Load environment variables
load_dotenv()
class RAGPipelineConfig(Config):
"""Configuration for RAG pipeline assets."""
# Source configuration
source_directory: str = "data/raw_documents"
file_extensions: List[str] = [".txt", ".md", ".pdf"]
# Chunking configuration
chunk_size: int = 512
chunk_overlap: int = 50
# Embedding configuration
embedding_model: str = "voyage-3" # or "text-embedding-3-small"
batch_size: int = 100
# Qdrant configuration
qdrant_url: str = "http://localhost:6333"
collection_name: str = "production_docs"
@asset(
description="Raw documents loaded from source directory",
metadata={
"dagster/priority": "high",
"owner": "data-team"
}
)
def raw_documents(context: AssetExecutionContext, config: RAGPipelineConfig) -> List[Dict[str, Any]]:
"""
Load raw documents from source directory.
Supports multiple file formats and tracks metadata for lineage.
Returns:
List of document dictionaries with content and metadata
"""
source_path = Path(config.source_directory)
if not source_path.exists():
context.log.warning(f"Source directory {source_path} does not exist. Creating sample documents...")
source_path.mkdir(parents=True, exist_ok=True)
_create_sample_documents(source_path)
documents = []
file_count = 0
for ext in config.file_extensions:
for file_path in source_path.rglob(f"*{ext}"):
try:
with open(file_path, 'r', encoding='utf-8') as f:
content = f.read()
documents.append({
"content": content,
"metadata": {
"source": str(file_path),
"file_name": file_path.name,
"file_type": ext,
"loaded_at": datetime.now().isoformat()
}
})
file_count += 1
except Exception as e:
context.log.error(f"Failed to load {file_path}: {e}")
context.log.info(f"Loaded {len(documents)} documents from {file_count} files")
# Track metadata for Dagster UI
return documents
@asset(
description="Documents chunked into 512-token segments with 50-token overlap",
deps=[raw_documents]
)
def chunked_documents(
context: AssetExecutionContext,
config: RAGPipelineConfig,
raw_documents: List[Dict[str, Any]]
) -> List[Document]:
"""
Split raw documents into chunks for embedding generation.
Uses recursive character splitting with configurable size and overlap.
Args:
raw_documents: Output from raw_documents asset
Returns:
List of LangChain Document objects
"""
context.log.info(f"Chunking {len(raw_documents)} documents...")
# Initialize splitter
splitter = RecursiveCharacterTextSplitter(
chunk_size=config.chunk_size,
chunk_overlap=config.chunk_overlap,
length_function=len,
separators=["\n\n", "\n", ". ", " ", ""]
)
# Convert to LangChain documents and split
all_chunks = []
for doc in raw_documents:
langchain_doc = Document(
page_content=doc["content"],
metadata=doc["metadata"]
)
chunks = splitter.split_documents([langchain_doc])
all_chunks.extend(chunks)
context.log.info(f"Created {len(all_chunks)} chunks")
# Log metadata for Dagster UI
context.log_event(
MaterializeResult(
metadata={
"num_chunks": len(all_chunks),
"avg_chunk_size": sum(len(c.page_content) for c in all_chunks) / len(all_chunks),
"chunk_size": config.chunk_size,
"chunk_overlap": config.chunk_overlap,
"preview": MetadataValue.md(all_chunks[0].page_content[:200] + "..." if all_chunks else "No chunks")
}
)
)
return all_chunks
@asset(
description="Embeddings generated using Voyage AI or OpenAI",
deps=[chunked_documents]
)
def document_embeddings(
context: AssetExecutionContext,
config: RAGPipelineConfig,
chunked_documents: List[Document]
) -> List[List[float]]:
"""
Generate embeddings for all document chunks.
Implements batching for efficient API usage and cost control.
Args:
chunked_documents: Output from chunked_documents asset
Returns:
List of embedding vectors
"""
context.log.info(f"Generating embeddings for {len(chunked_documents)} chunks...")
# Select embedding model
if config.embedding_model == "voyage-3":
if not os.getenv("VOYAGE_API_KEY"):
context.log.error("VOYAGE_API_KEY not set, falling back to OpenAI")
config.embedding_model = "text-embedding-3-small"
if config.embedding_model == "voyage-3":
context.log.info("Using Voyage AI voyage-3 (MTEB: 69.0)")
embeddings_model = VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY"),
batch_size=config.batch_size
)
else:
context.log.info("Using OpenAI text-embedding-3-small (MTEB: 62.3)")
embeddings_model = OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY"),
chunk_size=config.batch_size
)
# Generate embeddings in batches
texts = [doc.page_content for doc in chunked_documents]
try:
embeddings = embeddings_model.embed_documents(texts)
context.log.info(f"Successfully generated {len(embeddings)} embeddings")
# Log metadata
context.log_event(
MaterializeResult(
metadata={
"num_embeddings": len(embeddings),
"embedding_dimension": len(embeddings[0]),
"model": config.embedding_model,
"batch_size": config.batch_size
}
)
)
return embeddings
except Exception as e:
context.log.error(f"Embedding generation failed: {e}")
raise
@asset(
description="Embeddings indexed in Qdrant vector database",
deps=[chunked_documents, document_embeddings]
)
def vector_index(
context: AssetExecutionContext,
config: RAGPipelineConfig,
chunked_documents: List[Document],
document_embeddings: List[List[float]]
) -> Dict[str, Any]:
"""
Index embeddings in Qdrant vector database.
Creates collection if needed and uploads all vectors with metadata.
Args:
chunked_documents: Document chunks
document_embeddings: Generated embeddings
Returns:
Indexing statistics and metadata
"""
context.log.info("Indexing embeddings in Qdrant...")
# Initialize Qdrant client
client = QdrantClient(url=config.qdrant_url)
# Get embedding dimension
dimension = len(document_embeddings[0])
# Create collection if doesn't exist
collections = [c.name for c in client.get_collections().collections]
if config.collection_name not in collections:
context.log.info(f"Creating collection '{config.collection_name}' with dimension {dimension}")
client.create_collection(
collection_name=config.collection_name,
vectors_config=models.VectorParams(
size=dimension,
distance=models.Distance.COSINE
)
)
else:
context.log.info(f"Collection '{config.collection_name}' already exists")
# Index documents using LangChain wrapper
if config.embedding_model == "voyage-3":
embeddings_model = VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY")
)
else:
embeddings_model = OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY")
)
vectorstore = QdrantVectorStore(
client=client,
collection_name=config.collection_name,
embedding=embeddings_model
)
# Add documents with embeddings
try:
vectorstore.add_documents(chunked_documents)
context.log.info(f"Successfully indexed {len(chunked_documents)} documents")
# Get collection info
collection_info = client.get_collection(config.collection_name)
stats = {
"collection_name": config.collection_name,
"total_vectors": collection_info.points_count,
"dimension": dimension,
"indexed_at": datetime.now().isoformat()
}
# Log metadata
context.log_event(
MaterializeResult(
metadata={
"collection_name": config.collection_name,
"total_vectors": collection_info.points_count,
"dimension": dimension,
"distance_metric": "cosine"
}
)
)
return stats
except Exception as e:
context.log.error(f"Indexing failed: {e}")
raise
# Define the complete pipeline job
embedding_pipeline_job = define_asset_job(
name="embedding_pipeline_job",
selection=AssetSelection.all(),
description="Complete RAG embedding pipeline from documents to vector index"
)
# Define a schedule (daily at 2 AM)
embedding_pipeline_schedule = ScheduleDefinition(
name="daily_embedding_pipeline",
job=embedding_pipeline_job,
cron_schedule="0 2 * * *", # Daily at 2 AM
description="Run embedding pipeline daily to index new documents"
)
# Group all assets and jobs
defs = Definitions(
assets=[raw_documents, chunked_documents, document_embeddings, vector_index],
jobs=[embedding_pipeline_job],
schedules=[embedding_pipeline_schedule]
)
def _create_sample_documents(directory: Path):
"""
Create sample documents for demonstration.
Args:
directory: Directory to create sample documents in
"""
samples = {
"rag_overview.txt": """Retrieval-Augmented Generation (RAG) Overview
RAG is a technique that enhances LLM responses by retrieving relevant information from a knowledge base before generation. This approach significantly reduces hallucinations and enables LLMs to work with proprietary or up-to-date information.
The RAG pipeline consists of five stages:
1. Ingestion - Load documents from various sources
2. Indexing - Chunk and embed documents
3. Retrieval - Find relevant chunks for queries
4. Generation - Inject context into LLM prompts
5. Evaluation - Measure quality with RAGAS metrics
Benefits include improved accuracy, transparency (ability to cite sources), and easy knowledge updates without retraining.""",
"embedding_models.txt": """Embedding Model Comparison
Choosing the right embedding model impacts retrieval quality and cost:
Voyage AI voyage-3:
- Dimensions: 1024
- MTEB Score: 69.0 (best-in-class)
- Cost: $$$ ($0.12 per 1M tokens)
- Use case: Production systems requiring highest quality
OpenAI text-embedding-3-small:
- Dimensions: 1536
- MTEB Score: 62.3
- Cost: $ ($0.02 per 1M tokens)
- Use case: Development, prototyping, cost-sensitive apps
Never mix embedding models in the same index. If changing models, re-embed all documents.""",
"dagster_orchestration.md": """# Dagster for RAG Pipelines
Dagster is an asset-centric orchestration framework ideal for ML/AI pipelines.
## Key Features
1. **Asset Lineage**: Track dependencies between pipeline stages
2. **Versioning**: Automatic versioning of pipeline outputs
3. **Observability**: Rich UI for monitoring pipeline runs
4. **Type Safety**: Strong typing for asset dependencies
5. **Scheduling**: Cron-based or sensor-based execution
## Why Dagster for RAG?
- Clear data lineage from documents → chunks → embeddings → index
- Easy recomputation when upstream assets change
- Built-in metadata tracking for each asset
- Perfect for production ML systems
## Alternative: Prefect
Prefect is more workflow-oriented. Choose Dagster for data-centric pipelines, Prefect for workflow orchestration."""
}
for filename, content in samples.items():
file_path = directory / filename
with open(file_path, 'w') as f:
f.write(content)
if __name__ == "__main__":
"""
Run the pipeline locally for testing.
For production deployment:
1. Install Dagster: pip install dagster dagster-webserver
2. Start UI: dagster dev -f embedding_pipeline.py
3. Access UI: http://localhost:3000
4. Materialize assets from UI or CLI
"""
print("""
=== Dagster Embedding Pipeline ===
To run this pipeline:
1. Install dependencies:
pip install dagster dagster-webserver langchain-qdrant langchain-voyageai
2. Start Dagster UI:
dagster dev -f embedding_pipeline.py
3. Open browser:
http://localhost:3000
4. Materialize assets:
- Click "Materialize all" in UI
- Or use CLI: dagster asset materialize -f embedding_pipeline.py
5. View lineage:
- Asset graph shows dependencies
- Metadata shows stats for each asset
- Logs show execution details
Asset Pipeline:
raw_documents → chunked_documents → document_embeddings → vector_index
Schedule:
Daily at 2 AM (configurable in embedding_pipeline_schedule)
Configuration:
Edit RAGPipelineConfig class to customize:
- Source directory
- Chunk size/overlap
- Embedding model
- Qdrant settings
""")
Feast Feature Store Example
Complete feature store implementation with Feast for ML feature serving.
Setup
pip install -r requirements.txtInitialize Feast Repository
# Create feature repository
feast init feature_repo
cd feature_repoRun Example
python setup_features.pyThis will: 1. Define user features 2. Generate sample data 3. Apply features to registry 4. Materialize to online store 5. Demonstrate online/offline serving
File Descriptions
setup_features.py- Complete feature store setuprequirements.txt- DependenciesREADME.md- This file
Key Concepts
Feature View: Definition of features and their sources Entity: Join key (e.g., user_id) Online Store: Low-latency serving (Redis, DynamoDB) Offline Store: Training data (S3, BigQuery)
Production Deployment
See setup_features.py for production configuration with:
- Redis online store
- S3 offline store
- Point-in-time joins
- Feature versioning
# Feast Feature Store Example Dependencies
# Core Feast
feast==0.40.1
# Data processing
pandas==2.2.3
numpy==2.1.3
# Database backends (optional)
redis==5.2.1
psycopg2-binary==2.9.10
# Utilities
python-dotenv==1.0.1
"""
Feast Feature Store Complete Example
Demonstrates:
1. Feature definition
2. Data generation
3. Feature materialization
4. Online serving (low-latency)
5. Offline serving (training)
"""
import os
from datetime import datetime, timedelta
import pandas as pd
import numpy as np
from feast import FeatureStore, Entity, FeatureView, Field, FileSource
from feast.types import Float32, Int64, String
from feast.value_type import ValueType
def generate_sample_data():
"""Generate sample user features data."""
# Sample user IDs
user_ids = list(range(1001, 1101)) # 100 users
# Generate features
data = {
"user_id": user_ids,
"total_orders": np.random.randint(0, 100, size=len(user_ids)),
"avg_order_value": np.random.uniform(10.0, 500.0, size=len(user_ids)),
"lifetime_value": np.random.uniform(100.0, 5000.0, size=len(user_ids)),
"days_since_last_order": np.random.randint(0, 365, size=len(user_ids)),
"event_timestamp": [datetime.now() for _ in user_ids]
}
df = pd.DataFrame(data)
# Save to parquet
os.makedirs("data", exist_ok=True)
df.to_parquet("data/user_features.parquet")
print(f"Generated {len(df)} user feature records")
return df
def setup_feature_store():
"""Initialize Feast feature repository."""
# Check if feature_store.yaml exists
if not os.path.exists("feature_store.yaml"):
print("\nCreating feature repository...")
os.system("feast init feature_repo")
os.chdir("feature_repo")
print("Feature repository created!")
else:
print("Feature repository already exists")
def define_features():
"""Define feature views and entities."""
# Entity: User
user = Entity(
name="user",
join_keys=["user_id"],
description="User entity for ML features"
)
# Data source
user_source = FileSource(
path="../data/user_features.parquet",
timestamp_field="event_timestamp"
)
# Feature View: User Features
user_features = FeatureView(
name="user_features",
entities=[user],
schema=[
Field(name="total_orders", dtype=Int64),
Field(name="avg_order_value", dtype=Float32),
Field(name="lifetime_value", dtype=Float32),
Field(name="days_since_last_order", dtype=Int64)
],
source=user_source,
ttl=timedelta(days=1) # Features valid for 1 day
)
# Write feature definitions to features.py
features_code = f"""
from feast import Entity, FeatureView, Field, FileSource
from feast.types import Float32, Int64
from datetime import timedelta
# Entity
user = Entity(
name="user",
join_keys=["user_id"],
description="User entity for ML features"
)
# Data source
user_source = FileSource(
path="../data/user_features.parquet",
timestamp_field="event_timestamp"
)
# Feature View
user_features = FeatureView(
name="user_features",
entities=[user],
schema=[
Field(name="total_orders", dtype=Int64),
Field(name="avg_order_value", dtype=Float32),
Field(name="lifetime_value", dtype=Float32),
Field(name="days_since_last_order", dtype=Int64)
],
source=user_source,
ttl=timedelta(days=1)
)
"""
with open("features.py", "w") as f:
f.write(features_code)
print("Feature definitions created in features.py")
def apply_features():
"""Apply feature definitions to registry."""
print("\nApplying features to registry...")
os.system("feast apply")
print("Features applied!")
def materialize_features():
"""Materialize features to online store."""
print("\nMaterializing features to online store...")
store = FeatureStore(repo_path=".")
# Materialize from start_date to now
start_date = datetime.now() - timedelta(days=1)
end_date = datetime.now()
store.materialize(
start_date=start_date,
end_date=end_date
)
print("Features materialized!")
def demo_online_serving():
"""Demonstrate online feature serving."""
print("\n=== Online Serving Demo ===\n")
store = FeatureStore(repo_path=".")
# Get features for a single user (low-latency)
features = store.get_online_features(
features=[
"user_features:total_orders",
"user_features:avg_order_value",
"user_features:lifetime_value",
"user_features:days_since_last_order"
],
entity_rows=[{"user_id": 1001}]
).to_dict()
print("Online Features for user_id=1001:")
for key, value in features.items():
if key != "user_id":
print(f" {key}: {value[0]}")
def demo_offline_serving():
"""Demonstrate offline feature serving for training."""
print("\n=== Offline Serving Demo ===\n")
store = FeatureStore(repo_path=".")
# Entity dataframe (users and timestamps for training)
entity_df = pd.DataFrame({
"user_id": [1001, 1002, 1003, 1004, 1005],
"event_timestamp": [
datetime.now() for _ in range(5)
]
})
# Get historical features (point-in-time correct)
training_data = store.get_historical_features(
entity_df=entity_df,
features=[
"user_features:total_orders",
"user_features:avg_order_value",
"user_features:lifetime_value",
"user_features:days_since_last_order"
]
).to_df()
print("Training Data (Historical Features):")
print(training_data.head())
# Simulate training
print("\nSimulated Model Training:")
X = training_data[["total_orders", "avg_order_value", "lifetime_value"]]
print(f"Training features shape: {X.shape}")
print("Model would be trained here...")
def demo_prediction_serving():
"""Demonstrate using features for prediction."""
print("\n=== Prediction Serving Demo ===\n")
store = FeatureStore(repo_path=".")
# Get features for prediction
user_id = 1010
features = store.get_online_features(
features=[
"user_features:total_orders",
"user_features:avg_order_value",
"user_features:lifetime_value"
],
entity_rows=[{"user_id": user_id}]
).to_dict()
# Create feature vector
feature_vector = [
features["total_orders"][0],
features["avg_order_value"][0],
features["lifetime_value"][0]
]
print(f"Feature vector for user_id={user_id}:")
print(f" {feature_vector}")
# Simulate prediction
prediction = sum(feature_vector) / len(feature_vector) # Dummy prediction
print(f"\nSimulated prediction: {prediction:.2f}")
def main():
"""Main Feast feature store demo."""
print("=== Feast Feature Store Demo ===\n")
# Step 1: Generate sample data
print("Step 1: Generating sample data...")
generate_sample_data()
# Step 2: Setup feature store
print("\nStep 2: Setting up feature store...")
setup_feature_store()
# Step 3: Define features
print("\nStep 3: Defining features...")
define_features()
# Step 4: Apply to registry
print("\nStep 4: Applying features...")
apply_features()
# Step 5: Materialize to online store
print("\nStep 5: Materializing features...")
materialize_features()
# Step 6: Demo online serving
demo_online_serving()
# Step 7: Demo offline serving
demo_offline_serving()
# Step 8: Demo prediction serving
demo_prediction_serving()
print("\n=== Demo Complete ===")
print("\nKey Takeaways:")
print("1. Features defined ONCE in features.py")
print("2. Same features used for training (offline) and serving (online)")
print("3. No training-serving skew!")
print("4. Point-in-time correctness for training data")
print("5. Low-latency online serving for predictions")
if __name__ == "__main__":
main()
"""
Basic RAG Implementation with LangChain
This example demonstrates the fundamental RAG pattern:
1. Document loading and chunking
2. Embedding generation
3. Vector storage
4. Retrieval-augmented generation
Perfect starting point for understanding RAG architecture.
"""
import os
from typing import List, Dict, Any
from dotenv import load_dotenv
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough
from langchain_core.documents import Document
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_voyageai import VoyageAIEmbeddings
from langchain_qdrant import QdrantVectorStore
from langchain.text_splitter import RecursiveCharacterTextSplitter
from qdrant_client import QdrantClient, models
# Load environment variables
load_dotenv()
# Configuration
QDRANT_URL = os.getenv("QDRANT_URL", "http://localhost:6333")
COLLECTION_NAME = "basic_rag_demo"
CHUNK_SIZE = 512
CHUNK_OVERLAP = 50
def get_embeddings():
"""
Get embedding model based on available API keys.
Prioritizes Voyage AI (best quality) but falls back to OpenAI.
Returns:
Embeddings model instance
"""
if os.getenv("VOYAGE_API_KEY"):
print("Using Voyage AI voyage-3 (MTEB: 69.0)")
return VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY")
)
else:
print("Using OpenAI text-embedding-3-small (MTEB: 62.3)")
return OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY")
)
def setup_vector_store() -> QdrantVectorStore:
"""
Initialize Qdrant vector store with proper collection setup.
Creates collection if it doesn't exist with appropriate vector dimensions.
Returns:
Configured QdrantVectorStore instance
"""
client = QdrantClient(url=QDRANT_URL)
embeddings = get_embeddings()
# Detect embedding dimension
test_embedding = embeddings.embed_query("test")
dimension = len(test_embedding)
# Create collection if needed
collections = [c.name for c in client.get_collections().collections]
if COLLECTION_NAME not in collections:
print(f"Creating collection '{COLLECTION_NAME}' (dim={dimension})")
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=models.VectorParams(
size=dimension,
distance=models.Distance.COSINE
)
)
return QdrantVectorStore(
client=client,
collection_name=COLLECTION_NAME,
embedding=embeddings
)
def chunk_documents(texts: List[str]) -> List[Document]:
"""
Split texts into 512-token chunks with 50-token overlap.
This is the recommended default chunking strategy that balances:
- Context preservation (not too small)
- Relevance precision (not too large)
- Token limit handling (fits in most contexts)
Args:
texts: List of text strings to chunk
Returns:
List of Document objects with chunked content
"""
splitter = RecursiveCharacterTextSplitter(
chunk_size=CHUNK_SIZE,
chunk_overlap=CHUNK_OVERLAP,
length_function=len,
separators=["\n\n", "\n", ". ", " ", ""]
)
documents = [Document(page_content=text) for text in texts]
chunks = splitter.split_documents(documents)
print(f"Created {len(chunks)} chunks from {len(texts)} documents")
return chunks
def create_rag_chain(vectorstore: QdrantVectorStore):
"""
Build a basic RAG chain using LCEL (LangChain Expression Language).
Chain flow:
1. Retrieve relevant chunks from vector store
2. Format into prompt with context
3. Generate response with LLM
4. Parse output as string
Args:
vectorstore: Configured vector store for retrieval
Returns:
Runnable chain for RAG queries
"""
# Retriever: Get top 5 most relevant chunks
retriever = vectorstore.as_retriever(
search_type="similarity",
search_kwargs={"k": 5}
)
# Prompt template with context injection
template = """Answer the question based only on the following context:
{context}
Question: {question}
Provide a clear, concise answer. If the answer is not in the context, say "I don't have enough information to answer this question."
Answer:"""
prompt = ChatPromptTemplate.from_template(template)
# LLM with zero temperature for factual responses
llm = ChatOpenAI(
model="gpt-4o",
temperature=0.0,
openai_api_key=os.getenv("OPENAI_API_KEY")
)
# Build chain using LCEL pipe operator
chain = (
{"context": retriever, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
return chain
def main():
"""
Demonstrate basic RAG workflow with sample documents.
Shows complete pipeline from document ingestion to question answering.
"""
print("=== Basic RAG Pipeline Demo ===\n")
# Sample documents about AI/ML concepts
sample_texts = [
"""Retrieval-Augmented Generation (RAG) is a technique that enhances large language models by retrieving relevant information from a knowledge base before generating responses. This approach combines the benefits of retrieval systems with generative models.
RAG works in three steps:
1. Convert user query to embedding
2. Retrieve most similar documents from vector database
3. Inject retrieved context into LLM prompt
Benefits include reduced hallucinations, ability to cite sources, and easy knowledge updates without retraining.""",
"""Vector databases store data as high-dimensional embeddings, enabling semantic similarity search. Popular vector databases include Qdrant, Pinecone, Weaviate, and pgvector.
Key features:
- Fast approximate nearest neighbor (ANN) search
- Filtering by metadata
- Horizontal scalability
- HNSW (Hierarchical Navigable Small World) indexing
Qdrant is recommended for production RAG systems due to its performance and filtering capabilities.""",
"""Chunking strategy is critical for RAG quality. The recommended default is 512 tokens with 50-token overlap.
Why 512 tokens:
- Too small (<256): Loses context, requires more retrievals
- Too large (>1024): Includes irrelevant content, hits token limits
- 512 is the sweet spot for most use cases
Overlap prevents information loss at chunk boundaries and ensures continuity."""
]
# Step 1: Chunk documents
print("Step 1: Chunking documents...")
chunks = chunk_documents(sample_texts)
# Step 2: Setup vector store and index
print("\nStep 2: Setting up vector store...")
vectorstore = setup_vector_store()
print("Step 3: Indexing chunks...")
vectorstore.add_documents(chunks)
print(f"Indexed {len(chunks)} chunks successfully")
# Step 4: Create RAG chain
print("\nStep 4: Creating RAG chain...")
chain = create_rag_chain(vectorstore)
# Step 5: Query the system
print("\n=== RAG Query Examples ===\n")
questions = [
"What is RAG and how does it work?",
"Why are vector databases important for RAG?",
"What is the recommended chunk size and why?"
]
for question in questions:
print(f"Q: {question}")
answer = chain.invoke(question)
print(f"A: {answer}\n")
print("Demo complete! Basic RAG pipeline working successfully.")
if __name__ == "__main__":
main()
"""
Hybrid Search RAG Implementation (Semantic + Keyword)
This example demonstrates hybrid search combining:
1. Vector similarity search (semantic understanding)
2. BM25 keyword search (exact term matching)
3. Reciprocal Rank Fusion (RRF) for result merging
Use hybrid search when:
- Documents contain specific technical terms or codes
- Users query with exact keywords (product names, error codes)
- Pure semantic search misses important exact matches
"""
import os
from typing import List, Dict, Any
from dotenv import load_dotenv
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough
from langchain_core.documents import Document
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_voyageai import VoyageAIEmbeddings
from langchain_qdrant import QdrantVectorStore
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.retrievers import EnsembleRetriever
from langchain_community.retrievers import BM25Retriever
from qdrant_client import QdrantClient, models
# Load environment variables
load_dotenv()
# Configuration
QDRANT_URL = os.getenv("QDRANT_URL", "http://localhost:6333")
COLLECTION_NAME = "hybrid_search_demo"
CHUNK_SIZE = 512
CHUNK_OVERLAP = 50
def get_embeddings():
"""
Get embedding model for semantic search.
Returns:
Embeddings model instance
"""
if os.getenv("VOYAGE_API_KEY"):
print("Using Voyage AI voyage-3 for semantic search")
return VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY")
)
else:
print("Using OpenAI text-embedding-3-small for semantic search")
return OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY")
)
def setup_vector_store() -> QdrantVectorStore:
"""
Initialize Qdrant vector store for semantic search.
Returns:
Configured QdrantVectorStore instance
"""
client = QdrantClient(url=QDRANT_URL)
embeddings = get_embeddings()
# Detect embedding dimension
test_embedding = embeddings.embed_query("test")
dimension = len(test_embedding)
# Create collection if needed
collections = [c.name for c in client.get_collections().collections]
if COLLECTION_NAME not in collections:
print(f"Creating collection '{COLLECTION_NAME}' (dim={dimension})")
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=models.VectorParams(
size=dimension,
distance=models.Distance.COSINE
)
)
return QdrantVectorStore(
client=client,
collection_name=COLLECTION_NAME,
embedding=embeddings
)
def chunk_documents(texts: List[str]) -> List[Document]:
"""
Split texts into chunks for both semantic and keyword indexing.
Args:
texts: List of text strings to chunk
Returns:
List of Document objects with chunked content
"""
splitter = RecursiveCharacterTextSplitter(
chunk_size=CHUNK_SIZE,
chunk_overlap=CHUNK_OVERLAP,
length_function=len,
separators=["\n\n", "\n", ". ", " ", ""]
)
documents = [Document(page_content=text) for text in texts]
chunks = splitter.split_documents(documents)
print(f"Created {len(chunks)} chunks from {len(texts)} documents")
return chunks
def create_hybrid_retriever(
vectorstore: QdrantVectorStore,
chunks: List[Document],
k: int = 5
):
"""
Create hybrid retriever combining semantic and keyword search.
Uses Reciprocal Rank Fusion (RRF) to merge results from:
1. Vector similarity search (semantic understanding)
2. BM25 keyword search (exact term matching)
RRF formula: score = sum(1 / (rank + k)) for k=60 (standard)
Args:
vectorstore: Vector store for semantic search
chunks: Document chunks for BM25 indexing
k: Number of results to retrieve
Returns:
EnsembleRetriever combining both search methods
"""
print("Creating hybrid retriever (semantic + keyword)...")
# Semantic retriever (vector similarity)
semantic_retriever = vectorstore.as_retriever(
search_type="similarity",
search_kwargs={"k": k}
)
# Keyword retriever (BM25)
keyword_retriever = BM25Retriever.from_documents(chunks)
keyword_retriever.k = k
# Combine with RRF (weights: 50% semantic, 50% keyword)
# Adjust weights based on your use case:
# - Higher semantic weight (0.7) for conceptual queries
# - Higher keyword weight (0.7) for exact term searches
hybrid_retriever = EnsembleRetriever(
retrievers=[semantic_retriever, keyword_retriever],
weights=[0.5, 0.5], # Equal weighting
c=60 # RRF constant (standard value)
)
print("Hybrid retriever ready (RRF fusion enabled)")
return hybrid_retriever
def create_hybrid_rag_chain(retriever):
"""
Build RAG chain using hybrid retriever.
Args:
retriever: Hybrid retriever instance
Returns:
Runnable chain for hybrid RAG queries
"""
# Prompt template
template = """Answer the question based only on the following context:
{context}
Question: {question}
Provide a clear, concise answer. If the answer is not in the context, say "I don't have enough information to answer this question."
Answer:"""
prompt = ChatPromptTemplate.from_template(template)
# LLM
llm = ChatOpenAI(
model="gpt-4o",
temperature=0.0,
openai_api_key=os.getenv("OPENAI_API_KEY")
)
# Build chain
def format_docs(docs):
"""Format retrieved documents for context."""
return "\n\n".join(doc.page_content for doc in docs)
chain = (
{"context": retriever | format_docs, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
return chain
def main():
"""
Demonstrate hybrid search RAG with technical documentation.
Shows how hybrid search handles both semantic and exact keyword queries.
"""
print("=== Hybrid Search RAG Demo ===\n")
# Sample technical documents with specific terms and codes
sample_texts = [
"""Error Code E-RAG-001: Vector dimension mismatch
This error occurs when the embedding model dimensions don't match the Qdrant collection configuration.
Solution:
1. Check embedding model dimension (voyage-3 = 1024, text-embedding-3-small = 1536)
2. Verify Qdrant collection dimension matches
3. Recreate collection if dimensions differ
Example code:
```python
embeddings = VoyageAIEmbeddings(model="voyage-3") # 1024 dimensions
test_vec = embeddings.embed_query("test")
print(f"Dimension: {len(test_vec)}") # Should match collection
```""",
"""BM25 Algorithm for Keyword Search
BM25 (Best Match 25) is a ranking function used for keyword-based retrieval. It improves upon TF-IDF by considering document length normalization.
BM25 formula:
score(D,Q) = sum(IDF(qi) * (f(qi,D) * (k1 + 1)) / (f(qi,D) + k1 * (1 - b + b * |D| / avgdl)))
Where:
- f(qi,D): term frequency in document
- |D|: document length
- avgdl: average document length
- k1: term saturation parameter (typically 1.2-2.0)
- b: length normalization (typically 0.75)
BM25 works well for exact term matching, making it ideal for hybrid search.""",
"""Reciprocal Rank Fusion (RRF) Merging Strategy
RRF combines rankings from multiple retrieval systems. For each document, RRF calculates:
RRF_score = sum(1 / (k + rank_i))
Where:
- k: constant (typically 60)
- rank_i: rank from retrieval system i
Benefits:
- No score normalization needed
- Handles different scoring scales
- Simple and effective
Used in hybrid search to merge semantic and keyword results.""",
"""Product SKU-VEC-1024: Voyage AI Embedding Model
Specifications:
- Model: voyage-3
- Dimensions: 1024
- MTEB Score: 69.0
- Context length: 32,000 tokens
- Pricing: $0.12 per 1M tokens
Use cases:
- Production RAG systems
- Semantic search
- Document clustering
Alternative: SKU-VEC-1536 (OpenAI text-embedding-3-small)"""
]
# Step 1: Chunk documents
print("Step 1: Chunking documents...")
chunks = chunk_documents(sample_texts)
# Step 2: Setup vector store
print("\nStep 2: Setting up vector store...")
vectorstore = setup_vector_store()
# Step 3: Index documents
print("Step 3: Indexing chunks...")
vectorstore.add_documents(chunks)
print(f"Indexed {len(chunks)} chunks")
# Step 4: Create hybrid retriever
print("\nStep 4: Creating hybrid retriever...")
hybrid_retriever = create_hybrid_retriever(vectorstore, chunks, k=3)
# Step 5: Create RAG chain
print("Step 5: Creating hybrid RAG chain...")
chain = create_hybrid_rag_chain(hybrid_retriever)
# Step 6: Test with different query types
print("\n=== Hybrid Search Query Examples ===\n")
queries = [
# Exact keyword query (BM25 should excel)
{
"query": "What is error code E-RAG-001?",
"type": "Keyword (exact match)"
},
# Semantic query (vector search should excel)
{
"query": "How do I merge results from different search methods?",
"type": "Semantic (conceptual)"
},
# Product code query (hybrid should excel)
{
"query": "Tell me about SKU-VEC-1024",
"type": "Hybrid (code + context)"
},
# Mixed query
{
"query": "What algorithm does BM25 use for ranking?",
"type": "Hybrid (term + concept)"
}
]
for item in queries:
print(f"Query Type: {item['type']}")
print(f"Q: {item['query']}")
answer = chain.invoke(item['query'])
print(f"A: {answer}\n")
print("-" * 80 + "\n")
print("Demo complete! Hybrid search combines best of semantic + keyword retrieval.")
if __name__ == "__main__":
main()
"""
Complete RAG Pipeline Example with LangChain 0.3+
This example demonstrates:
1. Document loading
2. Chunking (512 tokens, 50 overlap)
3. Embedding generation (Voyage AI)
4. Vector storage (Qdrant)
5. Retrieval + generation
6. Streaming responses
"""
import os
from dotenv import load_dotenv
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_voyageai import VoyageAIEmbeddings
from langchain_qdrant import QdrantVectorStore
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.document_loaders import TextLoader
from qdrant_client import QdrantClient, models
# Load environment variables
load_dotenv()
# Configuration
QDRANT_URL = os.getenv("QDRANT_URL", "http://localhost:6333")
COLLECTION_NAME = "rag_demo_docs"
CHUNK_SIZE = 512
CHUNK_OVERLAP = 50
def get_embeddings():
"""Get embedding model (Voyage AI if available, else OpenAI)."""
if os.getenv("VOYAGE_API_KEY"):
print("Using Voyage AI embeddings (best quality)")
return VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY"),
batch_size=128
)
else:
print("Using OpenAI embeddings (VOYAGE_API_KEY not set)")
return OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY")
)
def setup_qdrant():
"""Initialize Qdrant client and collection."""
client = QdrantClient(url=QDRANT_URL)
# Get embedding dimension
embeddings = get_embeddings()
test_embedding = embeddings.embed_query("test")
dimension = len(test_embedding)
# Create collection if doesn't exist
collections = [c.name for c in client.get_collections().collections]
if COLLECTION_NAME not in collections:
print(f"Creating collection '{COLLECTION_NAME}' with dimension {dimension}")
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=models.VectorParams(
size=dimension,
distance=models.Distance.COSINE
)
)
else:
print(f"Collection '{COLLECTION_NAME}' already exists")
return client
def load_and_chunk_documents(file_paths: list[str]):
"""Load documents and split into chunks."""
all_chunks = []
for file_path in file_paths:
print(f"Loading: {file_path}")
loader = TextLoader(file_path)
documents = loader.load()
# Split into chunks
splitter = RecursiveCharacterTextSplitter(
chunk_size=CHUNK_SIZE,
chunk_overlap=CHUNK_OVERLAP,
length_function=len,
separators=["\n\n", "\n", " ", ""]
)
chunks = splitter.split_documents(documents)
all_chunks.extend(chunks)
print(f"Created {len(all_chunks)} chunks")
return all_chunks
def index_documents(chunks):
"""Index documents in Qdrant."""
client = setup_qdrant()
embeddings = get_embeddings()
vectorstore = QdrantVectorStore(
client=client,
collection_name=COLLECTION_NAME,
embedding=embeddings
)
print(f"Indexing {len(chunks)} chunks...")
vectorstore.add_documents(chunks)
print("Indexing complete!")
return vectorstore
def create_rag_chain(vectorstore):
"""Create RAG chain with retrieval + generation."""
# Retriever (use MMR for diversity)
retriever = vectorstore.as_retriever(
search_type="mmr",
search_kwargs={
"k": 5,
"fetch_k": 20,
"lambda_mult": 0.5
}
)
# Prompt template
template = """Answer the question based only on the following context:
{context}
Question: {question}
Provide a clear, concise answer. If the answer is not in the context, say "I don't have enough information to answer this question."
Answer:"""
prompt = ChatPromptTemplate.from_template(template)
# LLM (streaming for better UX)
llm = ChatOpenAI(
model="gpt-4o",
temperature=0.0, # Factual responses
streaming=True,
openai_api_key=os.getenv("OPENAI_API_KEY")
)
# Build chain using LCEL
chain = (
{"context": retriever, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
return chain
def main():
"""Main RAG pipeline demonstration."""
print("=== LangChain RAG Pipeline Demo ===\n")
# Step 1: Create sample documents
sample_docs = [
{
"path": "ml_basics.txt",
"content": """Machine Learning Basics
Machine learning is a subset of artificial intelligence that enables computers to learn from data without being explicitly programmed. It uses algorithms to identify patterns in data and make predictions or decisions.
Types of Machine Learning:
1. Supervised Learning - Training with labeled data
2. Unsupervised Learning - Finding patterns in unlabeled data
3. Reinforcement Learning - Learning through trial and error
Common applications include image recognition, natural language processing, and recommendation systems."""
},
{
"path": "deep_learning.txt",
"content": """Deep Learning Overview
Deep learning is a subset of machine learning that uses neural networks with multiple layers (hence "deep"). These networks can automatically learn hierarchical representations of data.
Neural networks consist of:
- Input layer: Receives data
- Hidden layers: Process and transform data
- Output layer: Produces predictions
Popular architectures include CNNs for images, RNNs for sequences, and Transformers for language tasks."""
}
]
# Write sample documents
for doc in sample_docs:
with open(doc["path"], "w") as f:
f.write(doc["content"])
# Step 2: Load and chunk documents
chunks = load_and_chunk_documents([doc["path"] for doc in sample_docs])
# Step 3: Index in Qdrant
vectorstore = index_documents(chunks)
# Step 4: Create RAG chain
chain = create_rag_chain(vectorstore)
# Step 5: Query the system
questions = [
"What is machine learning?",
"What are the types of machine learning?",
"How do neural networks work?",
"What are transformers used for?"
]
print("\n=== RAG Q&A Demo ===\n")
for question in questions:
print(f"Q: {question}")
print("A: ", end="", flush=True)
# Stream response
for chunk in chain.stream(question):
print(chunk, end="", flush=True)
print("\n")
# Cleanup
for doc in sample_docs:
os.remove(doc["path"])
print("\nDemo complete!")
if __name__ == "__main__":
main()
LangChain RAG Examples
Complete working examples of RAG pipelines using LangChain 0.3+.
Setup
1. Install dependencies:
pip install -r requirements.txt2. Set environment variables:
Create .env file:
# Required
OPENAI_API_KEY=your-openai-key
VOYAGE_API_KEY=your-voyage-key # Optional, for best embedding quality
# Optional (defaults to localhost)
QDRANT_URL=http://localhost:63333. Start Qdrant (local):
docker run -p 6333:6333 qdrant/qdrantExamples
1. Basic RAG (basic_rag.py)
Simplest production RAG implementation:
python basic_rag.pyFeatures:
- Document loading
- Chunking (512 tokens, 50 overlap)
- Embedding generation (Voyage AI or OpenAI)
- Vector storage (Qdrant)
- Retrieval + generation
2. Streaming RAG (streaming_rag.py)
RAG with streaming responses:
python streaming_rag.pyFeatures:
- All basic RAG features
- Streaming LLM responses
- Better UX for long answers
- FastAPI endpoint included
3. Hybrid Search (hybrid_search.py)
Combines vector and keyword search:
python hybrid_search.pyFeatures:
- Vector search (semantic)
- BM25 search (keyword)
- Ensemble retrieval (70% vector, 30% keyword)
- Better recall and precision
Usage Patterns
Quick Start
from langchain_qdrant import QdrantVectorStore
from langchain_voyageai import VoyageAIEmbeddings
from qdrant_client import QdrantClient
# Initialize
client = QdrantClient(url="http://localhost:6333")
embeddings = VoyageAIEmbeddings(model="voyage-3")
vectorstore = QdrantVectorStore(
client=client,
collection_name="docs",
embedding=embeddings
)
# Add documents
vectorstore.add_documents(documents)
# Query
results = vectorstore.similarity_search("What is machine learning?", k=5)Production Deployment
See streaming_rag.py for FastAPI integration:
# Start server
uvicorn streaming_rag:app --reload
# Query
curl -X POST http://localhost:8000/api/rag/stream \
-H "Content-Type: application/json" \
-d '{"question": "What is machine learning?"}'File Descriptions
basic_rag.py- Complete RAG implementation (150 lines)streaming_rag.py- Streaming + FastAPI (200 lines)hybrid_search.py- Vector + BM25 hybrid (180 lines)requirements.txt- All dependenciesREADME.md- This file
Troubleshooting
Error: "Collection not found"
- Run
python basic_rag.pyfirst to create collection - Or create manually in Qdrant
Error: "No API key provided"
- Set
OPENAI_API_KEYandVOYAGE_API_KEYin.env
Slow responses:
- Use
text-embedding-3-smallinstead ofvoyage-3for development - Reduce
k(retrieve fewer chunks)
Poor quality:
- Increase chunk overlap (50 → 100)
- Try hybrid search
- Add re-ranking
# LangChain RAG Example Dependencies
# Core LangChain
langchain-core==0.3.28
langchain==0.3.12
# LLM Providers
langchain-openai==0.2.14
langchain-anthropic==0.2.9
langchain-voyageai==0.2.0
# Vector Stores
langchain-qdrant==0.2.0
qdrant-client==1.12.1
# Utilities
python-dotenv==1.0.1
fastapi==0.115.6
uvicorn==0.34.0
"""
Streaming RAG Implementation with LangChain
This example demonstrates streaming responses for better UX:
1. Token-by-token streaming from LLM
2. Async/await patterns for non-blocking I/O
3. Server-Sent Events (SSE) integration pattern
4. Progress indicators during retrieval
Use streaming when:
- Building chat interfaces (better perceived performance)
- Long-form answers (users see progress immediately)
- Production web applications (prevents timeout issues)
"""
import os
import asyncio
from typing import List, AsyncIterator
from dotenv import load_dotenv
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough
from langchain_core.documents import Document
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_voyageai import VoyageAIEmbeddings
from langchain_qdrant import QdrantVectorStore
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.callbacks.base import AsyncCallbackHandler
from qdrant_client import QdrantClient, models
# Load environment variables
load_dotenv()
# Configuration
QDRANT_URL = os.getenv("QDRANT_URL", "http://localhost:6333")
COLLECTION_NAME = "streaming_rag_demo"
CHUNK_SIZE = 512
CHUNK_OVERLAP = 50
class StreamingCallbackHandler(AsyncCallbackHandler):
"""
Callback handler for tracking streaming events.
Useful for debugging and monitoring streaming behavior.
"""
async def on_llm_start(self, serialized: dict, prompts: List[str], **kwargs):
"""Called when LLM starts generating."""
print("\n[Streaming] LLM generation started...")
async def on_llm_new_token(self, token: str, **kwargs):
"""Called for each new token. Can be used for custom processing."""
# Token is already printed by the main loop
pass
async def on_llm_end(self, response, **kwargs):
"""Called when LLM finishes generating."""
print("\n[Streaming] LLM generation complete")
def get_embeddings():
"""
Get embedding model for semantic search.
Returns:
Embeddings model instance
"""
if os.getenv("VOYAGE_API_KEY"):
return VoyageAIEmbeddings(
model="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY")
)
else:
return OpenAIEmbeddings(
model="text-embedding-3-small",
openai_api_key=os.getenv("OPENAI_API_KEY")
)
def setup_vector_store() -> QdrantVectorStore:
"""
Initialize Qdrant vector store.
Returns:
Configured QdrantVectorStore instance
"""
client = QdrantClient(url=QDRANT_URL)
embeddings = get_embeddings()
# Detect embedding dimension
test_embedding = embeddings.embed_query("test")
dimension = len(test_embedding)
# Create collection if needed
collections = [c.name for c in client.get_collections().collections]
if COLLECTION_NAME not in collections:
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=models.VectorParams(
size=dimension,
distance=models.Distance.COSINE
)
)
return QdrantVectorStore(
client=client,
collection_name=COLLECTION_NAME,
embedding=embeddings
)
def chunk_documents(texts: List[str]) -> List[Document]:
"""
Split texts into chunks.
Args:
texts: List of text strings to chunk
Returns:
List of Document objects with chunked content
"""
splitter = RecursiveCharacterTextSplitter(
chunk_size=CHUNK_SIZE,
chunk_overlap=CHUNK_OVERLAP,
length_function=len,
separators=["\n\n", "\n", ". ", " ", ""]
)
documents = [Document(page_content=text) for text in texts]
chunks = splitter.split_documents(documents)
return chunks
def create_streaming_rag_chain(vectorstore: QdrantVectorStore):
"""
Build streaming RAG chain with LCEL.
Key difference from basic RAG: streaming=True on ChatOpenAI
Args:
vectorstore: Configured vector store for retrieval
Returns:
Runnable chain that supports streaming
"""
# Retriever
retriever = vectorstore.as_retriever(
search_type="similarity",
search_kwargs={"k": 5}
)
# Prompt template
template = """Answer the question based only on the following context:
{context}
Question: {question}
Provide a detailed, informative answer. If the answer is not in the context, say "I don't have enough information to answer this question."
Answer:"""
prompt = ChatPromptTemplate.from_template(template)
# LLM with streaming enabled
llm = ChatOpenAI(
model="gpt-4o",
temperature=0.0,
streaming=True, # Enable streaming
openai_api_key=os.getenv("OPENAI_API_KEY")
)
# Build chain
chain = (
{"context": retriever, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)
return chain
async def stream_response(chain, question: str) -> AsyncIterator[str]:
"""
Stream response from RAG chain asynchronously.
This is the pattern to use in FastAPI endpoints or async contexts.
Args:
chain: RAG chain instance
question: User question
Yields:
Token chunks as they're generated
"""
async for chunk in chain.astream(question):
yield chunk
async def async_streaming_example(chain, question: str):
"""
Demonstrate async streaming (for web servers).
Use this pattern in FastAPI/Starlette apps for streaming endpoints.
Args:
chain: RAG chain instance
question: User question
"""
print(f"\n[Async] Q: {question}")
print("[Async] A: ", end="", flush=True)
async for chunk in stream_response(chain, question):
print(chunk, end="", flush=True)
print() # Newline after streaming completes
def sync_streaming_example(chain, question: str):
"""
Demonstrate sync streaming (simpler, for scripts).
Use this pattern in CLI tools or simple scripts.
Args:
chain: RAG chain instance
question: User question
"""
print(f"\n[Sync] Q: {question}")
print("[Sync] A: ", end="", flush=True)
for chunk in chain.stream(question):
print(chunk, end="", flush=True)
print() # Newline after streaming completes
async def fastapi_integration_example():
"""
Show how to integrate streaming RAG with FastAPI.
This is a complete example you can copy into your FastAPI app.
"""
print("\n=== FastAPI Integration Example ===\n")
print("""
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain_qdrant import QdrantVectorStore
app = FastAPI()
# Initialize once at startup
vectorstore = setup_vector_store()
chain = create_streaming_rag_chain(vectorstore)
@app.post("/api/rag/stream")
async def stream_rag(question: str):
\"\"\"
Stream RAG responses using Server-Sent Events.
Usage:
curl -N -X POST "http://localhost:8000/api/rag/stream" \\
-H "Content-Type: application/json" \\
-d '{"question": "What is RAG?"}'
\"\"\"
async def generate():
async for chunk in chain.astream(question):
# SSE format: data: {chunk}\\n\\n
yield f"data: {chunk}\\n\\n"
return StreamingResponse(
generate(),
media_type="text/event-stream"
)
@app.post("/api/rag/complete")
async def complete_rag(question: str):
\"\"\"
Non-streaming endpoint (waits for full response).
Use when you need the complete response before sending.
\"\"\"
response = await chain.ainvoke(question)
return {"answer": response}
# Run with: uvicorn main:app --reload
""")
def main():
"""
Demonstrate streaming RAG patterns.
Shows both sync and async streaming approaches.
"""
print("=== Streaming RAG Demo ===\n")
# Sample documents
sample_texts = [
"""Streaming in Large Language Models
Streaming responses improves user experience by providing immediate feedback. Instead of waiting for the complete response, users see tokens appear in real-time, similar to ChatGPT.
Benefits:
1. Better perceived performance (users see progress)
2. Prevents timeout issues (long responses)
3. Lower memory usage (process tokens incrementally)
4. Early cancellation possible (stop generating if not relevant)
Technical implementation uses Server-Sent Events (SSE) or WebSockets.""",
"""Server-Sent Events (SSE) for Streaming
SSE is a standard for server-to-client streaming over HTTP. It's simpler than WebSockets for one-way communication.
SSE format:
data: {message}\\n\\n
Python (FastAPI) example:
```python
async def generate():
for chunk in stream:
yield f"data: {chunk}\\n\\n"
return StreamingResponse(generate(), media_type="text/event-stream")
```
JavaScript client:
```javascript
const eventSource = new EventSource('/api/stream');
eventSource.onmessage = (event) => {
console.log(event.data);
};
```""",
"""Async vs Sync Streaming in Python
Async streaming (asyncio):
- Non-blocking I/O
- Handles multiple concurrent requests
- Required for FastAPI/Starlette
- Uses: astream(), ainvoke()
Sync streaming:
- Blocking I/O
- Simpler code
- Good for CLI tools, scripts
- Uses: stream(), invoke()
LangChain supports both patterns. Choose based on your use case."""
]
# Setup
print("Setting up vector store...")
chunks = chunk_documents(sample_texts)
vectorstore = setup_vector_store()
vectorstore.add_documents(chunks)
print(f"Indexed {len(chunks)} chunks\n")
# Create streaming chain
chain = create_streaming_rag_chain(vectorstore)
# Example 1: Sync streaming (simple)
print("=== Example 1: Sync Streaming (for CLI tools) ===")
sync_streaming_example(
chain,
"How does streaming improve user experience?"
)
# Example 2: Async streaming (for web servers)
print("\n=== Example 2: Async Streaming (for FastAPI) ===")
asyncio.run(async_streaming_example(
chain,
"What is the difference between SSE and WebSockets?"
))
# Example 3: Show FastAPI integration
asyncio.run(fastapi_integration_example())
print("\n=== Key Takeaways ===")
print("""
1. Enable streaming with: ChatOpenAI(streaming=True)
2. Use chain.stream() for sync contexts (CLI tools)
3. Use chain.astream() for async contexts (FastAPI)
4. SSE format: "data: {chunk}\\n\\n"
5. Always flush output for real-time display
For production: Combine streaming with error handling, retry logic,
and rate limiting for robust RAG applications.
""")
if __name__ == "__main__":
main()
"""
LlamaIndex RAG Example
Simpler API compared to LangChain for RAG-focused applications.
"""
import os
from dotenv import load_dotenv
from llama_index.core import (
VectorStoreIndex,
Document,
Settings,
StorageContext
)
from llama_index.llms.openai import OpenAI
from llama_index.embeddings.openai import OpenAIEmbedding
from llama_index.embeddings.voyageai import VoyageAIEmbedding
from llama_index.vector_stores.qdrant import QdrantVectorStore
from qdrant_client import QdrantClient
# Load environment variables
load_dotenv()
def setup_llama_index():
"""Configure LlamaIndex settings."""
# LLM
Settings.llm = OpenAI(
model="gpt-4o",
temperature=0.0,
api_key=os.getenv("OPENAI_API_KEY")
)
# Embeddings (Voyage AI if available)
if os.getenv("VOYAGE_API_KEY"):
print("Using Voyage AI embeddings")
Settings.embed_model = VoyageAIEmbedding(
model_name="voyage-3",
voyage_api_key=os.getenv("VOYAGE_API_KEY")
)
else:
print("Using OpenAI embeddings")
Settings.embed_model = OpenAIEmbedding(
model="text-embedding-3-small",
api_key=os.getenv("OPENAI_API_KEY")
)
# Chunking settings
Settings.chunk_size = 512
Settings.chunk_overlap = 50
def create_index():
"""Create vector index with Qdrant."""
# Sample documents
documents = [
Document(
text="""Machine Learning Basics
Machine learning is a subset of artificial intelligence that enables computers to learn from data without being explicitly programmed. It uses algorithms to identify patterns in data and make predictions or decisions.
Types of Machine Learning:
1. Supervised Learning - Training with labeled data
2. Unsupervised Learning - Finding patterns in unlabeled data
3. Reinforcement Learning - Learning through trial and error
Common applications include image recognition, natural language processing, and recommendation systems.""",
metadata={"source": "ml_basics", "category": "fundamentals"}
),
Document(
text="""Deep Learning Overview
Deep learning is a subset of machine learning that uses neural networks with multiple layers (hence "deep"). These networks can automatically learn hierarchical representations of data.
Neural networks consist of:
- Input layer: Receives data
- Hidden layers: Process and transform data
- Output layer: Produces predictions
Popular architectures include CNNs for images, RNNs for sequences, and Transformers for language tasks.""",
metadata={"source": "deep_learning", "category": "advanced"}
)
]
# Setup Qdrant
client = QdrantClient(url=os.getenv("QDRANT_URL", "http://localhost:6333"))
vector_store = QdrantVectorStore(
client=client,
collection_name="llamaindex_demo"
)
storage_context = StorageContext.from_defaults(vector_store=vector_store)
# Create index (automatically chunks and embeds)
print("Creating index...")
index = VectorStoreIndex.from_documents(
documents,
storage_context=storage_context,
show_progress=True
)
print("Index created!")
return index
def query_index(index):
"""Query the index."""
# Create query engine
query_engine = index.as_query_engine(
similarity_top_k=5,
streaming=True # Stream responses
)
# Sample queries
questions = [
"What is machine learning?",
"What are the types of machine learning?",
"How do neural networks work?",
"What are transformers used for?"
]
print("\n=== LlamaIndex RAG Q&A ===\n")
for question in questions:
print(f"Q: {question}")
print("A: ", end="", flush=True)
# Stream response
response = query_engine.query(question)
for text in response.response_gen:
print(text, end="", flush=True)
print("\n")
def main():
"""Main LlamaIndex RAG demo."""
print("=== LlamaIndex RAG Demo ===\n")
# Setup
setup_llama_index()
# Create index
index = create_index()
# Query
query_index(index)
print("\nDemo complete!")
if __name__ == "__main__":
main()
LlamaIndex RAG Example
LlamaIndex alternative to LangChain with simpler API for RAG-focused applications.
Setup
pip install -r requirements.txtCreate .env:
OPENAI_API_KEY=your-key
VOYAGE_API_KEY=your-key # Optional
QDRANT_URL=http://localhost:6333When to Use LlamaIndex vs LangChain
Use LlamaIndex when:
- Building pure RAG applications
- Want simpler API
- Don't need complex agent orchestration
- Prefer opinionated defaults
Use LangChain when:
- Need general LLM orchestration
- Building complex agent workflows
- Want more control and flexibility
- Need extensive integrations
Examples
Query Engine (query_engine.py)
Simple RAG query engine:
python query_engine.pyFeatures:
- Automatic chunking
- Simplified indexing
- Built-in query optimization
File Descriptions
query_engine.py- Basic LlamaIndex RAG (100 lines)requirements.txt- DependenciesREADME.md- This file
# LlamaIndex Agent Example Dependencies
# Core LlamaIndex
llama-index==0.12.6
llama-index-core==0.12.6
# LLM Providers
llama-index-llms-openai==0.3.18
llama-index-embeddings-openai==0.3.7
llama-index-embeddings-voyageai==0.3.2
# Vector Stores
llama-index-vector-stores-qdrant==0.5.1
qdrant-client==1.12.1
# Utilities
python-dotenv==1.0.1
skill: "ai-data-engineering"
version: "1.0"
domain: "backend"
# Base outputs required for all AI data engineering projects
base_outputs:
- path: "pipelines/"
must_contain: []
reason: "Data pipeline code for RAG, embeddings, and feature engineering"
- path: "data/"
must_contain: ["documents/", "chunks/"]
reason: "Data storage for source documents and processed chunks"
- path: "config/"
must_contain: []
reason: "Configuration files for embedding models, vector databases, and pipelines"
- path: "tests/"
must_contain: []
reason: "RAGAS evaluation tests and data quality checks"
# Conditional outputs based on configuration
conditional_outputs:
maturity:
starter:
- path: "pipelines/basic_rag.py"
must_contain: ["VectorStore", "Retriever", "ChatPromptTemplate"]
reason: "Basic RAG pipeline with retrieval and generation"
- path: "data/chunks/"
must_contain: []
reason: "Chunked documents (512 tokens, 50-token overlap)"
- path: "config/embeddings.yaml"
must_contain: ["model:", "dimension:"]
reason: "Embedding model configuration (OpenAI text-embedding-3-small for starter)"
- path: "tests/eval_qa.json"
must_contain: ["question", "answer", "contexts"]
reason: "Basic RAGAS evaluation dataset (10-20 Q&A pairs)"
- path: "requirements.txt"
must_contain: ["langchain", "qdrant-client"]
reason: "Core dependencies for RAG pipeline"
intermediate:
- path: "pipelines/streaming_rag.py"
must_contain: ["stream", "StreamingResponse", "async"]
reason: "Streaming RAG responses for real-time UX"
- path: "pipelines/hybrid_search.py"
must_contain: ["BM25", "vector_search", "rerank"]
reason: "Hybrid search combining vector and keyword retrieval"
- path: "feature_repo/"
must_contain: ["feature_store.yaml", "features/"]
reason: "Feast feature store for ML feature serving"
- path: "data/embeddings/"
must_contain: []
reason: "Cached embeddings for performance"
- path: "tests/ragas_evaluation.py"
must_contain: ["faithfulness", "answer_relevancy", "context_precision"]
reason: "Complete RAGAS metrics evaluation script"
- path: "orchestration/dagster/"
must_contain: ["@asset"]
reason: "Dagster pipeline orchestration for embedding generation"
- path: "config/vector_db.yaml"
must_contain: ["collection_name:", "dimension:", "distance:"]
reason: "Vector database configuration (Qdrant)"
advanced:
- path: "pipelines/multi_modal_rag.py"
must_contain: ["ImageEmbeddings", "MultiVectorRetriever"]
reason: "Multi-modal RAG for images, text, and tables"
- path: "pipelines/agent_rag.py"
must_contain: ["Agent", "Tool", "ReAct"]
reason: "Agentic RAG with tool calling and reasoning"
- path: "feature_repo/"
must_contain: ["feature_store.yaml", "features/", "data_sources/"]
reason: "Complete Feast setup with online/offline stores"
- path: "orchestration/dagster/"
must_contain: ["@asset", "@job", "@sensor", "@schedule"]
reason: "Production Dagster pipelines with scheduling and monitoring"
- path: "data_versioning/lakefs/"
must_contain: ["lakectl.yaml"]
reason: "LakeFS data versioning for experiment tracking"
- path: "tests/benchmark_retrieval.py"
must_contain: ["NDCG", "MRR", "Recall@K"]
reason: "Advanced retrieval quality benchmarks"
- path: "monitoring/ragas_metrics.py"
must_contain: ["faithfulness", "answer_relevancy", "context_precision", "context_recall"]
reason: "Production RAGAS monitoring with alerting"
- path: "pipelines/reranking.py"
must_contain: ["Cohere", "rerank", "cross-encoder"]
reason: "Re-ranking pipeline for improved retrieval precision"
model_type:
rag:
- path: "pipelines/rag_chain.py"
must_contain: ["Retriever", "ChatPromptTemplate", "StrOutputParser"]
reason: "RAG chain with retrieval and generation"
- path: "data/chunks/"
must_contain: []
reason: "Document chunks (512 tokens default)"
- path: "config/chunking.yaml"
must_contain: ["chunk_size:", "overlap:", "strategy:"]
reason: "Chunking strategy configuration"
- path: "tests/ragas_eval.json"
must_contain: ["question", "answer", "contexts", "ground_truth"]
reason: "RAGAS evaluation dataset"
embeddings:
- path: "pipelines/embedding_generation.py"
must_contain: ["VoyageAIEmbeddings", "embed_documents", "embed_query"]
reason: "Embedding generation pipeline (Voyage AI voyage-3)"
- path: "data/embeddings/"
must_contain: []
reason: "Generated embeddings cache"
- path: "config/embeddings.yaml"
must_contain: ["model:", "dimension:", "batch_size:"]
reason: "Embedding model configuration"
- path: "tests/embedding_quality.py"
must_contain: ["cosine_similarity", "MTEB"]
reason: "Embedding quality tests"
feature_engineering:
- path: "feature_repo/features/"
must_contain: ["*.py"]
reason: "Feast feature definitions"
- path: "feature_repo/feature_store.yaml"
must_contain: ["project:", "provider:", "online_store:"]
reason: "Feast feature store configuration"
- path: "pipelines/feature_pipeline.py"
must_contain: ["FeatureStore", "get_online_features", "materialize"]
reason: "Feature materialization pipeline"
- path: "tests/feature_validation.py"
must_contain: ["assert", "feature_values"]
reason: "Feature value validation tests"
semantic_search:
- path: "pipelines/semantic_search.py"
must_contain: ["embed_query", "search", "query_vector"]
reason: "Semantic search pipeline with vector similarity"
- path: "pipelines/hybrid_search.py"
must_contain: ["BM25", "vector_search", "fusion"]
reason: "Hybrid search combining semantic and keyword"
- path: "config/search.yaml"
must_contain: ["top_k:", "filters:", "score_threshold:"]
reason: "Search configuration and tuning"
vector_database:
qdrant:
- path: "config/qdrant.yaml"
must_contain: ["collection_name:", "dimension:", "distance:"]
reason: "Qdrant vector database configuration"
- path: "pipelines/qdrant_setup.py"
must_contain: ["QdrantClient", "create_collection", "VectorParams"]
reason: "Qdrant collection setup script"
- path: "requirements.txt"
must_contain: ["qdrant-client"]
reason: "Qdrant client dependency"
pinecone:
- path: "config/pinecone.yaml"
must_contain: ["index_name:", "dimension:", "metric:"]
reason: "Pinecone vector database configuration"
- path: "pipelines/pinecone_setup.py"
must_contain: ["Pinecone", "create_index"]
reason: "Pinecone index setup script"
- path: "requirements.txt"
must_contain: ["pinecone-client"]
reason: "Pinecone client dependency"
weaviate:
- path: "config/weaviate.yaml"
must_contain: ["class_name:", "vectorizer:"]
reason: "Weaviate schema configuration"
- path: "pipelines/weaviate_setup.py"
must_contain: ["Client", "schema.create"]
reason: "Weaviate schema setup script"
- path: "requirements.txt"
must_contain: ["weaviate-client"]
reason: "Weaviate client dependency"
orchestration:
dagster:
- path: "orchestration/dagster/assets/"
must_contain: ["@asset"]
reason: "Dagster asset definitions for data lineage"
- path: "orchestration/dagster/jobs.py"
must_contain: ["@job", "@schedule"]
reason: "Dagster job scheduling and orchestration"
- path: "orchestration/dagster/sensors.py"
must_contain: ["@sensor"]
reason: "Dagster sensors for event-driven pipelines"
- path: "requirements.txt"
must_contain: ["dagster", "dagster-webserver"]
reason: "Dagster dependencies"
prefect:
- path: "orchestration/prefect/flows/"
must_contain: ["@flow", "@task"]
reason: "Prefect flow and task definitions"
- path: "orchestration/prefect/deployments.yaml"
must_contain: ["name:", "schedule:"]
reason: "Prefect deployment configuration"
- path: "requirements.txt"
must_contain: ["prefect"]
reason: "Prefect dependency"
airflow:
- path: "orchestration/airflow/dags/"
must_contain: ["DAG(", "schedule_interval"]
reason: "Airflow DAG definitions"
- path: "orchestration/airflow/dags/"
must_contain: ["*.py"]
reason: "Airflow pipeline orchestration"
- path: "requirements.txt"
must_contain: ["apache-airflow"]
reason: "Airflow dependency"
framework:
langchain:
- path: "pipelines/"
must_contain: ["langchain", "ChatPromptTemplate", "Retriever"]
reason: "LangChain orchestration framework"
- path: "requirements.txt"
must_contain: ["langchain", "langchain-core"]
reason: "LangChain dependencies"
llamaindex:
- path: "pipelines/"
must_contain: ["llama_index", "VectorStoreIndex", "QueryEngine"]
reason: "LlamaIndex orchestration framework"
- path: "requirements.txt"
must_contain: ["llama-index"]
reason: "LlamaIndex dependency"
# Scaffolding files that should be created as starting points
scaffolding:
- path: "pipelines/basic_rag.py"
reason: "Starter RAG pipeline template"
- path: "data/documents/.gitkeep"
reason: "Initialize source documents directory"
- path: "data/chunks/.gitkeep"
reason: "Initialize chunks directory"
- path: "data/embeddings/.gitkeep"
reason: "Initialize embeddings cache directory"
- path: "config/embeddings.yaml"
reason: "Embedding model configuration template"
- path: "config/vector_db.yaml"
reason: "Vector database configuration template"
- path: "config/chunking.yaml"
reason: "Chunking strategy configuration template"
- path: "tests/eval_qa.json"
reason: "RAGAS evaluation dataset template"
- path: "tests/README.md"
reason: "Document evaluation strategy and metrics"
- path: "orchestration/README.md"
reason: "Document orchestration approach (Dagster/Prefect/Airflow)"
- path: "requirements.txt"
reason: "Python dependencies for AI data pipelines"
- path: ".env.example"
reason: "Environment variables template (API keys, database URLs)"
- path: ".gitignore"
reason: "Ignore embeddings cache, API keys, and temporary files"
- path: "README.md"
reason: "Document RAG architecture, setup, and evaluation approach"
# Metadata
metadata:
primary_blueprints: ["rag-pipeline", "ml-pipeline"]
contributes_to:
- "RAG (Retrieval-Augmented Generation) pipelines"
- "Embedding generation and semantic search"
- "ML feature stores and serving"
- "AI/ML data orchestration"
- "Vector database integration"
- "Data versioning for ML experiments"
common_patterns:
- "RAG 5-stage pipeline (Ingestion → Indexing → Retrieval → Generation → Evaluation)"
- "512-token chunks with 50-token overlap (default chunking strategy)"
- "Voyage AI voyage-3 embeddings for production (MTEB 69.0)"
- "RAGAS evaluation metrics (faithfulness, answer_relevancy, context_precision, context_recall)"
- "Hybrid search (vector + BM25 keyword search)"
- "Feast feature store for online/offline ML features"
- "Dagster asset-centric orchestration for data lineage"
- "LakeFS data versioning for experiment tracking"
integration_points:
frontend: "Integrates with ai-chat skill for streaming RAG responses"
search: "Provides semantic search backend for search-filter skill"
databases: "Connects to vector databases (Qdrant, Pinecone, Weaviate)"
orchestration: "Scheduled by Dagster/Prefect/Airflow"
monitoring: "RAGAS metrics tracking and alerting"
versioning: "LakeFS for data/model versioning"
typical_directory_structure: |
project/
├── pipelines/
│ ├── basic_rag.py # Simple RAG chain
│ ├── streaming_rag.py # Streaming responses
│ ├── hybrid_search.py # Vector + BM25
│ ├── embedding_generation.py
│ └── reranking.py # Cohere re-ranking
├── data/
│ ├── documents/ # Source PDFs, DOCX, Markdown
│ ├── chunks/ # Chunked documents
│ └── embeddings/ # Cached embeddings
├── config/
│ ├── embeddings.yaml # Model: voyage-3, dimension: 1024
│ ├── vector_db.yaml # Qdrant configuration
│ └── chunking.yaml # chunk_size: 512, overlap: 50
├── tests/
│ ├── eval_qa.json # RAGAS evaluation dataset
│ ├── ragas_evaluation.py # Metrics runner
│ └── benchmark_retrieval.py # NDCG, MRR, Recall@K
├── orchestration/
│ ├── dagster/
│ │ ├── assets/ # @asset definitions
│ │ ├── jobs.py # @job, @schedule
│ │ └── sensors.py # @sensor
│ └── prefect/
│ └── flows/ # @flow, @task
├── feature_repo/ # Feast feature store
│ ├── features/
│ └── feature_store.yaml
├── monitoring/
│ └── ragas_metrics.py # Production metrics tracking
└── requirements.txt
Document Chunking Strategies for RAG
Comprehensive guide to chunking documents for optimal retrieval quality in RAG systems.
Table of Contents
- Why Chunking Matters
- Default Recommendation
- Chunking Methods
- 1. Fixed Token-Based (Default)
- 2. Semantic Chunking
- 3. Code-Aware Chunking
- 4. Markdown-Aware Chunking
- 5. Sentence-Based Chunking
- Content-Type Specific Strategies
- API Documentation
- Research Papers (PDF)
- Chat Logs / Conversations
- Code Repositories
- Advanced Chunking Patterns
- Sliding Window
- Parent-Child Chunking
- Contextual Compression
- Chunking Strategy Decision Tree
- Chunking Metadata
- Validation and Testing
- Test Chunk Quality
- Evaluate Retrieval Quality
- Common Pitfalls
- Best Practices
- Resources
Why Chunking Matters
Chunking is the most critical decision for RAG quality. Poor chunking causes:
- Lost context (chunks too small)
- Irrelevant retrieval (chunks too large)
- Information fragmentation (bad split points)
- Poor semantic coherence
Impact: Good chunking can improve RAG recall by 20-30%.
Default Recommendation
For most RAG systems:
- Chunk size: 512 tokens
- Overlap: 50 tokens (10%)
- Method: Fixed token-based with recursive splitting
Why these values:
- 512 tokens balances context vs precision
- 50 token overlap prevents information loss at boundaries
- Fits within most embedding model limits (8K+ context)
Chunking Methods
1. Fixed Token-Based (Default)
from langchain.text_splitter import RecursiveCharacterTextSplitter
splitter = RecursiveCharacterTextSplitter(
chunk_size=512, # Tokens, not characters
chunk_overlap=50,
length_function=len, # Use tiktoken for accurate token count
separators=["\n\n", "\n", ". ", " ", ""], # Try splits in order
)
chunks = splitter.split_text(document_text)Use when: General purpose, mixed content types
2. Semantic Chunking
from langchain.text_splitter import SemanticChunker
from langchain_voyageai import VoyageAIEmbeddings
semantic_splitter = SemanticChunker(
embeddings=VoyageAIEmbeddings(model="voyage-3"),
breakpoint_threshold_type="percentile", # Split at semantic boundaries
breakpoint_threshold_amount=90, # Top 10% similarity drops
)
chunks = semantic_splitter.split_text(document_text)Use when: Narrative content (articles, books), need semantic coherence Trade-off: Slower (requires embeddings), variable chunk sizes
3. Code-Aware Chunking
from langchain.text_splitter import RecursiveCharacterTextSplitter
code_splitter = RecursiveCharacterTextSplitter.from_language(
language="python", # python, javascript, rust, go, etc.
chunk_size=512,
chunk_overlap=50,
)
chunks = code_splitter.split_text(source_code)Splits on:
- Function/class boundaries
- Logical blocks (if/for/while)
- Comment sections
- Import statements
Use when: Code documentation, API references, technical docs
4. Markdown-Aware Chunking
from langchain.text_splitter import MarkdownHeaderTextSplitter
markdown_splitter = MarkdownHeaderTextSplitter(
headers_to_split_on=[
("#", "h1"),
("##", "h2"),
("###", "h3"),
]
)
chunks = markdown_splitter.split_text(markdown_text)Use when: Markdown documentation, README files, blog posts
5. Sentence-Based Chunking
import nltk
nltk.download('punkt')
from langchain.text_splitter import NLTKTextSplitter
sentence_splitter = NLTKTextSplitter(
chunk_size=512,
chunk_overlap=50,
)
chunks = sentence_splitter.split_text(text)Use when: Need to preserve sentence boundaries, conversational content
Content-Type Specific Strategies
API Documentation
# Chunk by endpoint + method
# Include endpoint metadata in each chunk
chunk_metadata = {
"endpoint": "/api/users",
"method": "POST",
"section": "Authentication",
}Chunk size: 256-512 tokens (API docs are dense) Split on: Endpoint boundaries, parameter sections
Research Papers (PDF)
# Chunk by section hierarchy
# Preserve abstract, introduction, methods separately
chunk_metadata = {
"section": "Methods",
"page": 5,
"paper_title": "...",
}Chunk size: 512-1024 tokens (academic writing needs context) Split on: Section headers, paragraph boundaries
Chat Logs / Conversations
# Chunk by conversation turn or time window
# Include speaker metadata
chunk_metadata = {
"speaker": "user_123",
"timestamp": "2025-12-03T10:00:00Z",
"conversation_id": "conv_456",
}Chunk size: 256-512 tokens Split on: Speaker turns, time windows (5-10 messages)
Code Repositories
# Chunk by file or function
chunk_metadata = {
"file_path": "src/api/users.py",
"function": "create_user",
"language": "python",
}Chunk size: 256-512 tokens (functions are self-contained) Split on: Function/class boundaries, file boundaries for small files
Advanced Chunking Patterns
Sliding Window
def sliding_window_chunks(text, window_size=512, stride=256):
"""Create overlapping chunks with sliding window"""
chunks = []
for i in range(0, len(text), stride):
chunk = text[i:i + window_size]
if len(chunk) < window_size / 2: # Skip small final chunk
break
chunks.append(chunk)
return chunksUse when: Dense documents where context is critical Trade-off: More chunks = higher storage and query costs
Parent-Child Chunking
# Store both small chunks (for retrieval) and large parent chunks (for context)
small_chunks = split_text(document, chunk_size=128)
large_chunks = split_text(document, chunk_size=512)
# Retrieve small, return large
for small_chunk in retrieved_chunks:
parent_chunk = find_parent_chunk(small_chunk)
context.append(parent_chunk)Use when: Need precise retrieval but broad context Trade-off: Complex indexing, higher storage
Contextual Compression
from langchain.retrievers import ContextualCompressionRetriever
from langchain.retrievers.document_compressors import LLMChainExtractor
# Retrieve large chunks, compress to relevant portions
compressor = LLMChainExtractor.from_llm(llm)
compression_retriever = ContextualCompressionRetriever(
base_compressor=compressor,
base_retriever=vector_store.as_retriever(),
)Use when: Want to reduce irrelevant content in context Trade-off: Extra LLM call per query
Chunking Strategy Decision Tree
CONTENT TYPE?
├─ CODE
│ └─ Code-aware chunking (256-512 tokens, split on functions)
│
├─ STRUCTURED DOCS (Markdown, API docs)
│ └─ Header/section-based (256-512 tokens)
│
├─ NARRATIVE (Articles, books)
│ └─ Semantic chunking (variable size, meaning-based)
│
├─ CONVERSATIONAL (Chat logs)
│ └─ Turn-based or time-window (256-512 tokens)
│
└─ MIXED CONTENT
└─ Fixed token with overlap (512 tokens, 50 overlap)Chunking Metadata
Essential metadata for production:
chunk_metadata = {
# Source tracking
"source": "docs/api-reference.md",
"source_type": "documentation",
"last_updated": "2025-12-03",
# Hierarchical context
"section": "Authentication",
"subsection": "OAuth 2.1",
"heading_hierarchy": ["API", "Auth", "OAuth 2.1"],
# Chunking info
"chunk_index": 3,
"total_chunks": 12,
"chunk_method": "fixed_token",
# Content classification
"content_type": "code_example",
"programming_language": "python",
}Validation and Testing
Test Chunk Quality
def validate_chunks(chunks: List[str]) -> Dict:
"""Validate chunk quality"""
stats = {
"total_chunks": len(chunks),
"avg_length": sum(len(c) for c in chunks) / len(chunks),
"min_length": min(len(c) for c in chunks),
"max_length": max(len(c) for c in chunks),
"empty_chunks": sum(1 for c in chunks if not c.strip()),
}
# Check for issues
issues = []
if stats["empty_chunks"] > 0:
issues.append(f"{stats['empty_chunks']} empty chunks")
if stats["min_length"] < 50:
issues.append("Very small chunks detected")
if stats["max_length"] > 2000:
issues.append("Very large chunks detected")
return {"stats": stats, "issues": issues}Evaluate Retrieval Quality
from ragas.metrics import context_recall
# Test if chunking preserves retrievable information
test_data = {
"question": ["How do I refresh OAuth tokens?"],
"contexts": [retrieved_chunks],
"ground_truth": ["Full answer from original doc"],
}
result = evaluate(test_data, metrics=[context_recall])
# Target: >0.80 recallCommon Pitfalls
❌ Chunking by character count instead of tokens
# Wrong: Characters != tokens
chunks = [text[i:i+2000] for i in range(0, len(text), 2000)]
# Right: Use token-aware splitting
import tiktoken
tokenizer = tiktoken.get_encoding("cl100k_base")❌ No overlap between chunks
# Loses information at boundaries
chunks = split_text(text, chunk_size=512, chunk_overlap=0)
# Better: 10% overlap
chunks = split_text(text, chunk_size=512, chunk_overlap=50)❌ One-size-fits-all chunking
# Code and prose need different strategies
chunks = split_text(mixed_content, chunk_size=512) # Suboptimal
# Better: Detect content type and adjust
if is_code(content):
chunks = code_aware_split(content)
else:
chunks = semantic_split(content)Best Practices
1. Start with 512/50 - Optimize later based on metrics 2. Measure retrieval quality - Use RAGAS context_recall 3. Preserve hierarchy - Include heading paths in metadata 4. Test edge cases - Very long docs, mixed content 5. Monitor chunk sizes - Track distribution over time 6. A/B test strategies - Compare retrieval quality 7. Document your choice - Explain why you chose specific parameters 8. Version your chunking - Re-chunk if strategy changes
Resources
- LangChain Text Splitters: https://python.langchain.com/docs/modules/data_connection/document_transformers/
- Semantic Chunking Paper: https://arxiv.org/abs/2312.06648
- LTTB Downsampling: https://github.com/sveinn-steinarsson/flot-downsample
Data Versioning for AI/ML Pipelines
Version control for datasets, embeddings, and model artifacts using LakeFS and DVC.
Table of Contents
- Why Data Versioning
- LakeFS (Recommended)
- Key Features
- Installation
- Basic Workflow
- RAG Pipeline Integration
- Compare Experiments
- DVC (Alternative)
- Installation
- Basic Workflow
- DVC Pipelines
- Comparison: LakeFS vs DVC
- Embedding Versioning Best Practices
- 1. Version Embedding Model + Data Together
- 2. Immutable Embeddings
- 3. Tag Important Versions
- Experiment Tracking Integration
- MLflow + LakeFS
- Reproducibility Checklist
- Resources
Why Data Versioning
The Problem:
- "Which embedding model generated these vectors?"
- "What data was the model trained on?"
- "How do I reproduce this RAG result?"
- "Can I A/B test different chunking strategies?"
The Solution: Git-like operations for data (branch, commit, merge, rollback).
LakeFS (Recommended)
Git for data lakes. Acquired DVC team (November 2025) - now the unified standard.
Key Features
- Branch, commit, merge operations on data
- Time travel (access historical versions)
- Zero-copy branching (instant, no duplication)
- Works with S3, Azure Blob, GCS
- Iceberg/Delta Lake integration
Installation
# Docker
docker run -p 8000:8000 treeverse/lakefs:latest
# Access UI: http://localhost:8000
# Default credentials: AKIAIOSFOLQUICKSTART / wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEYBasic Workflow
import lakefs
# Initialize client
client = lakefs.Client(
host="http://localhost:8000",
username="AKIAIOSFOLQUICKSTART",
password="wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY",
)
# Create repository
repo = client.repositories.create(
name="my-rag-pipeline",
storage_namespace="s3://my-bucket/lakefs/",
)
# Create branch for experiment
branch = repo.branch("main").create("experiment-voyage-3-embeddings")
# Upload data to branch
branch.object("embeddings/voyage-3.parquet").upload("local-file.parquet")
# Commit changes
branch.commit(message="Switch to Voyage AI embeddings")
# Merge to main after validation
branch.merge_into("main")RAG Pipeline Integration
# Version embeddings pipeline
def generate_embeddings_versioned(documents, branch_name="main"):
# 1. Create experimental branch
branch = repo.branch("main").create(f"embeddings-{branch_name}")
# 2. Generate embeddings
embeddings = voyage_ai.embed(documents)
# 3. Save to branch
branch.object("embeddings/vectors.parquet").upload_dataframe(embeddings)
# 4. Run evaluation
metrics = evaluate_rag(embeddings)
# 5. Commit with metrics
branch.commit(
message=f"Generated embeddings using Voyage AI",
metadata={"recall": metrics.recall, "precision": metrics.precision}
)
# 6. Merge if metrics improved
if metrics.recall > baseline_recall:
branch.merge_into("main")
print("✓ Embeddings improved, merged to main")
else:
print("✗ No improvement, branch preserved for analysis")Compare Experiments
# Compare two chunking strategies
main_branch = repo.branch("main")
exp_branch = repo.branch("chunking-semantic")
# Get diff
diff = exp_branch.diff(other=main_branch)
print(f"Changed objects: {len(diff.results)}")
# Compare metrics
main_metrics = main_branch.object("metrics.json").read()
exp_metrics = exp_branch.object("metrics.json").read()
print(f"Main recall: {main_metrics['recall']}")
print(f"Experiment recall: {exp_metrics['recall']}")DVC (Alternative)
Data Version Control - Git-like tool specifically for ML datasets.
Note: DVC team acquired by LakeFS (Nov 2025). LakeFS now recommended for new projects.
Installation
pip install dvc[s3] # or [gs], [azure]Basic Workflow
# Initialize DVC
dvc init
# Track data file
dvc add data/embeddings.parquet
# Commit tracking file (.dvc)
git add data/embeddings.parquet.dvc data/.gitignore
git commit -m "Add embeddings"
# Push data to remote storage
dvc push
# Pull data on another machine
dvc pullDVC Pipelines
# dvc.yaml
stages:
chunk_documents:
cmd: python scripts/chunk.py
deps:
- data/raw/documents.pdf
params:
- chunk_size
- overlap
outs:
- data/chunks/
generate_embeddings:
cmd: python scripts/embed.py
deps:
- data/chunks/
outs:
- data/embeddings/vectors.parquet
evaluate:
cmd: python scripts/evaluate.py
deps:
- data/embeddings/vectors.parquet
metrics:
- metrics.json# Run pipeline
dvc repro
# View metrics across experiments
dvc metrics show
# Compare experiments
dvc metrics diff experiment-1 experiment-2Comparison: LakeFS vs DVC
| Feature | LakeFS | DVC |
|---|---|---|
| Branching | Zero-copy (instant) | Full copy (slow) |
| Scale | Petabytes | Gigabytes-Terabytes |
| Format | Any (object storage) | Any (file-based) |
| Integration | S3 API compatible | Git-based |
| UI | Web UI included | CLI only |
| Team | Same team (acquired DVC) | Merged into LakeFS |
| Best for | Production data lakes | Small ML projects |
Recommendation: Use LakeFS for new projects (unified future direction).
Embedding Versioning Best Practices
1. Version Embedding Model + Data Together
embedding_metadata = {
"model": "voyage-3",
"model_version": "1.0",
"dimensions": 1024,
"input_data_version": "main@abc123",
"chunking_strategy": "fixed-512-50",
"generated_at": "2025-12-03T10:00:00Z",
}2. Immutable Embeddings
Never update existing embeddings in-place. Create new versions.
# ❌ Don't update
embeddings/vectors.parquet # Update in place
# ✅ Do version
embeddings/v1/vectors.parquet
embeddings/v2/vectors.parquet # New version3. Tag Important Versions
# LakeFS tags
lakefs tag create my-repo production-v1.0 --ref main
# DVC tags
git tag -a embeddings-v1.0 -m "Voyage AI production embeddings"
dvc pushExperiment Tracking Integration
MLflow + LakeFS
import mlflow
with mlflow.start_run():
# Log parameters
mlflow.log_param("chunk_size", 512)
mlflow.log_param("embedding_model", "voyage-3")
mlflow.log_param("data_version", branch.head.id)
# Generate embeddings
embeddings = generate_embeddings()
# Log metrics
metrics = evaluate_rag(embeddings)
mlflow.log_metrics(metrics)
# Log LakeFS commit
mlflow.log_param("lakefs_commit", branch.head.id)Reproducibility Checklist
- [ ] Data version tracked (LakeFS/DVC)
- [ ] Embedding model version recorded
- [ ] Chunking parameters documented
- [ ] Random seeds set (for reproducible experiments)
- [ ] Dependencies pinned (requirements.txt)
- [ ] Evaluation metrics logged
- [ ] Commit hash recorded
- [ ] Environment documented (Python version, OS)
Resources
- LakeFS Docs: https://docs.lakefs.io/
- DVC Docs: https://dvc.org/doc
- MLflow: https://mlflow.org/docs/
Data Orchestration Tools for AI Pipelines
Comparison of Dagster, Prefect, Airflow 3.0, and dbt for orchestrating AI/ML data workflows.
Table of Contents
- Quick Comparison
- Dagster (Recommended for AI Pipelines)
- Why Dagster for AI?
- Basic Pipeline
- Run Pipeline
- Prefect (Developer-Friendly)
- Basic Workflow
- Deployments
- Airflow 3.0
- DAG Definition
- dbt (SQL Transformations Only)
- Use in AI Pipelines
- Decision Framework
- Scheduling Patterns
- Dagster Schedules
- Prefect Schedules
- Airflow Schedules
- Monitoring and Alerting
- Dagster Sensors
- Prefect Automations
- Performance Optimization
- Parallel Task Execution
- Best Practices
- AI Pipeline Example Comparison
- Dagster (Asset-Centric)
- Prefect (Task-Centric)
- Airflow (Operator-Centric)
- Resources
Quick Comparison
| Tool | Best For | Language | Complexity | Learning Curve |
|---|---|---|---|---|
| Dagster | ML/AI pipelines, data lineage | Python | Medium | Medium |
| Prefect | General workflows, developer UX | Python | Low-Medium | Low |
| Airflow 3.0 | Enterprise, battle-tested | Python | High | High |
| dbt | SQL transformations only | SQL | Low | Low |
Dagster (Recommended for AI Pipelines)
Best for: RAG pipelines, embedding generation, ML feature engineering
Why Dagster for AI?
- Asset-centric design - Think in terms of data assets, not tasks
- Best data lineage - Visualize dependencies across pipeline
- Type system - Catch errors before runtime
- Testing built-in - Unit test pipelines easily
- Partitions - Process data in time windows or logical partitions
Basic Pipeline
from dagster import asset, Definitions
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain_voyageai import VoyageAIEmbeddings
from qdrant_client import QdrantClient
@asset
def raw_documents():
"""Load raw documents from S3"""
return load_from_s3("s3://bucket/docs/")
@asset
def chunked_documents(raw_documents):
"""Split documents into 512-token chunks"""
splitter = RecursiveCharacterTextSplitter(
chunk_size=512,
chunk_overlap=50,
)
return splitter.split_documents(raw_documents)
@asset
def embeddings(chunked_documents):
"""Generate embeddings with Voyage AI"""
embedder = VoyageAIEmbeddings(model="voyage-3")
vectors = embedder.embed_documents([doc.page_content for doc in chunked_documents])
return vectors
@asset
def vector_database(chunked_documents, embeddings):
"""Index into Qdrant"""
client = QdrantClient("localhost", port=6333)
client.upsert(
collection_name="documents",
points=[
{"id": i, "vector": vec, "payload": {"text": doc.page_content}}
for i, (doc, vec) in enumerate(zip(chunked_documents, embeddings))
]
)
return {"indexed": len(embeddings)}
defs = Definitions(assets=[raw_documents, chunked_documents, embeddings, vector_database])Run Pipeline
# Install
pip install dagster dagster-webserver
# Start UI
dagster dev
# Materialize assets
dagster asset materialize --select raw_documents+Access UI: http://localhost:3000
---
Prefect (Developer-Friendly)
Best for: General-purpose workflows, simpler setup than Airflow
Basic Workflow
from prefect import flow, task
from prefect.tasks import exponential_backoff
from typing import List
@task(retries=3, retry_delay_seconds=exponential_backoff(backoff_factor=2))
def fetch_documents() -> List[str]:
"""Fetch documents with automatic retries"""
return ["doc1", "doc2", "doc3"]
@task
def chunk_document(doc: str) -> List[str]:
"""Chunk single document"""
return doc.split() # Simplified
@task
def generate_embedding(chunk: str) -> List[float]:
"""Generate embedding for chunk"""
return [0.1, 0.2, 0.3] # Mock embedding
@flow(name="embedding-pipeline")
def embedding_pipeline():
"""Main embedding pipeline flow"""
docs = fetch_documents()
all_chunks = []
for doc in docs:
chunks = chunk_document(doc)
all_chunks.extend(chunks)
embeddings = [generate_embedding(chunk) for chunk in all_chunks]
return {"total_embeddings": len(embeddings)}
# Run
if __name__ == "__main__":
embedding_pipeline()Deployments
# Start Prefect server
prefect server start
# Deploy flow
prefect deploy---
Airflow 3.0
Best for: Enterprise deployments, existing Airflow users, complex dependencies
DAG Definition
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data-team',
'retries': 3,
'retry_delay': timedelta(minutes=5),
}
def chunk_documents(**context):
"""Chunk documents task"""
docs = context['ti'].xcom_pull(task_ids='fetch_documents')
chunks = []
for doc in docs:
chunks.extend(chunk_text(doc, chunk_size=512))
return chunks
def generate_embeddings(**context):
"""Generate embeddings task"""
chunks = context['ti'].xcom_pull(task_ids='chunk_documents')
embeddings = voyage_ai.embed(chunks)
return embeddings
with DAG(
'rag_pipeline',
default_args=default_args,
description='RAG embedding pipeline',
schedule_interval='@daily',
start_date=datetime(2025, 12, 1),
catchup=False,
) as dag:
fetch = PythonOperator(
task_id='fetch_documents',
python_callable=fetch_documents,
)
chunk = PythonOperator(
task_id='chunk_documents',
python_callable=chunk_documents,
)
embed = PythonOperator(
task_id='generate_embeddings',
python_callable=generate_embeddings,
)
index = PythonOperator(
task_id='index_to_qdrant',
python_callable=index_to_qdrant,
)
fetch >> chunk >> embed >> index---
dbt (SQL Transformations Only)
Best for: Data warehouse transformations, NOT for embedding generation
Use in AI Pipelines
dbt is best for pre-processing data BEFORE embedding generation:
-- models/cleaned_documents.sql
SELECT
id,
title,
REGEXP_REPLACE(content, '[^a-zA-Z0-9\s]', '') AS cleaned_content,
category,
last_updated
FROM raw_documents
WHERE last_updated > CURRENT_DATE - INTERVAL '30 days'
AND content IS NOT NULLThen: Use Dagster/Prefect for embedding generation.
---
Decision Framework
WORKFLOW COMPLEXITY?
├─ SIMPLE (10 steps, SQL + Python)
│ └─ Prefect (lowest overhead)
│
├─ AI/ML FOCUS (embeddings, features, model training)
│ └─ Dagster (asset-centric, best lineage)
│
├─ ENTERPRISE (100+ DAGs, complex dependencies)
│ └─ Airflow 3.0 (battle-tested, mature)
│
└─ SQL-ONLY (data warehouse transformations)
└─ dbt (specialized tool)Scheduling Patterns
Dagster Schedules
from dagster import ScheduleDefinition
daily_embedding_schedule = ScheduleDefinition(
name="daily_embeddings",
cron_schedule="0 2 * * *", # 2 AM daily
job=embedding_job,
)Prefect Schedules
from prefect.schedules import CronSchedule
schedule = CronSchedule(cron="0 2 * * *", timezone="UTC")
@flow(schedule=schedule)
def scheduled_pipeline():
# Pipeline logic
passAirflow Schedules
# schedule_interval in DAG definition
schedule_interval='@daily' # or '0 2 * * *'Monitoring and Alerting
Dagster Sensors
from dagster import sensor, RunRequest
@sensor(job=embedding_job)
def s3_document_sensor():
"""Trigger pipeline when new documents uploaded to S3"""
new_files = check_s3_for_new_files()
if new_files:
return RunRequest(
run_config={"s3_files": new_files}
)Prefect Automations
from prefect.blocks.notifications import SlackWebhook
slack = SlackWebhook.load("my-slack-webhook")
@flow(on_failure=[slack.notify])
def pipeline_with_alerts():
# Pipeline that alerts on failure
passPerformance Optimization
Parallel Task Execution
Dagster:
from dagster import OpExecutionContext
@asset
def embeddings(context: OpExecutionContext, chunks):
# Dagster automatically parallelizes across workers
return [generate_embedding(chunk) for chunk in chunks]Prefect:
from prefect import flow, task
@task
def generate_embedding(chunk):
return embed(chunk)
@flow
def parallel_embeddings(chunks):
# Prefect automatically parallelizes
futures = [generate_embedding.submit(chunk) for chunk in chunks]
return [f.result() for f in futures]Best Practices
1. Start simple - Use Prefect for initial prototypes 2. Idempotent tasks - Safe to retry and re-run 3. Partition large datasets - Process in time windows 4. Monitor failures - Set up alerting 5. Version pipeline code - Git + semantic versioning 6. Test locally - Don't test in production 7. Log metadata - Track model versions, parameters 8. Graceful degradation - Handle missing/delayed data
AI Pipeline Example Comparison
Dagster (Asset-Centric)
@asset
def cleaned_text(): ...
@asset(deps=[cleaned_text])
def chunks(): ...
@asset(deps=[chunks])
def embeddings(): ...Prefect (Task-Centric)
@flow
def pipeline():
text = clean_text()
chunks = create_chunks(text)
embeddings = generate_embeddings(chunks)Airflow (Operator-Centric)
fetch >> clean >> chunk >> embed >> indexResources
- Dagster Docs: https://docs.dagster.io/
- Prefect Docs: https://docs.prefect.io/
- Airflow Docs: https://airflow.apache.org/docs/
- dbt Docs: https://docs.getdbt.com/
- LakeFS Docs: https://docs.lakefs.io/
#!/usr/bin/env python3
"""
Document Chunking Script (TOKEN-FREE)
Chunks documents with configurable size and overlap without loading into context.
Usage:
python chunk_documents.py --input docs/ --output chunks/ --chunk-size 512 --overlap 50
Dependencies:
pip install langchain langchain-text-splitters
"""
import argparse
import json
import sys
from pathlib import Path
from typing import List, Dict
def chunk_file(
file_path: Path,
chunk_size: int,
chunk_overlap: int,
file_type: str = "text"
) -> List[Dict]:
"""Chunk a single file."""
try:
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.document_loaders import TextLoader
except ImportError:
print("Error: Missing dependencies. Run: pip install langchain langchain-text-splitters", file=sys.stderr)
sys.exit(1)
# Load document
try:
if file_type == "text":
loader = TextLoader(str(file_path))
else:
print(f"Warning: Unsupported file type '{file_type}', treating as text", file=sys.stderr)
loader = TextLoader(str(file_path))
documents = loader.load()
except Exception as e:
print(f"Error loading {file_path}: {e}", file=sys.stderr)
return []
# Create splitter
splitter = RecursiveCharacterTextSplitter(
chunk_size=chunk_size,
chunk_overlap=chunk_overlap,
length_function=len,
separators=["\n\n", "\n", " ", ""]
)
# Split documents
chunks = splitter.split_documents(documents)
# Convert to dict format
chunk_dicts = []
for i, chunk in enumerate(chunks):
chunk_dicts.append({
"chunk_id": i,
"source_file": str(file_path),
"content": chunk.page_content,
"metadata": chunk.metadata,
"chunk_size": len(chunk.page_content)
})
return chunk_dicts
def chunk_directory(
input_dir: Path,
output_dir: Path,
chunk_size: int,
chunk_overlap: int,
pattern: str = "*.txt"
):
"""Chunk all files in directory."""
# Find files
files = list(input_dir.glob(pattern))
if not files:
print(f"Warning: No files matching '{pattern}' found in {input_dir}", file=sys.stderr)
return
print(f"Found {len(files)} files")
# Create output directory
output_dir.mkdir(parents=True, exist_ok=True)
# Process each file
total_chunks = 0
for file_path in files:
print(f"Chunking: {file_path.name}...", end=" ")
chunks = chunk_file(file_path, chunk_size, chunk_overlap)
if chunks:
# Save chunks
output_file = output_dir / f"{file_path.stem}_chunks.json"
with open(output_file, "w") as f:
json.dump(chunks, f, indent=2)
print(f"{len(chunks)} chunks")
total_chunks += len(chunks)
else:
print("FAILED")
# Save summary
summary = {
"input_directory": str(input_dir),
"output_directory": str(output_dir),
"chunk_size": chunk_size,
"chunk_overlap": chunk_overlap,
"total_files": len(files),
"total_chunks": total_chunks,
"avg_chunks_per_file": total_chunks / len(files) if files else 0
}
summary_file = output_dir / "chunking_summary.json"
with open(summary_file, "w") as f:
json.dump(summary, f, indent=2)
print(f"\nTotal: {total_chunks} chunks from {len(files)} files")
print(f"Average: {summary['avg_chunks_per_file']:.1f} chunks per file")
print(f"\nSummary saved to: {summary_file}")
def main():
parser = argparse.ArgumentParser(
description="Chunk documents for RAG pipelines"
)
parser.add_argument(
"--input",
required=True,
help="Input directory containing documents"
)
parser.add_argument(
"--output",
required=True,
help="Output directory for chunks"
)
parser.add_argument(
"--chunk-size",
type=int,
default=512,
help="Chunk size in characters (default: 512)"
)
parser.add_argument(
"--overlap",
type=int,
default=50,
help="Chunk overlap in characters (default: 50)"
)
parser.add_argument(
"--pattern",
default="*.txt",
help="File pattern to match (default: *.txt)"
)
args = parser.parse_args()
input_dir = Path(args.input)
output_dir = Path(args.output)
if not input_dir.exists():
print(f"Error: Input directory not found: {input_dir}", file=sys.stderr)
return 1
chunk_directory(
input_dir,
output_dir,
args.chunk_size,
args.overlap,
args.pattern
)
return 0
if __name__ == "__main__":
sys.exit(main())
Related skills
FAQ
What chunk size does it recommend?
512 tokens with 50-100 tokens of overlap using fixed token-based chunking as the default.
How does it measure RAG quality?
With RAGAS LLM-as-judge metrics: faithfulness, answer relevancy, context precision, and context recall.