Building a RAG Pipeline | Kamiwaza Docs
Documentation for Kamiwaza 0.5.1
This documentation provides guidance on building a Retrieval-Augmented Generation (RAG) pipeline using Kamiwaza's core services.
What You'll Build
By the end of this guide, you'll have:
- Document ingestion system that processes various file formats
- Embedding pipeline that converts text to vector representations
- Vector search system for finding relevant context
- LLM integration that generates responses using retrieved context
- Web interface for querying your documents
Prerequisites
Before starting, ensure you have:
- Kamiwaza installed and running ( Installation Guide)
- At least 16GB of available RAM
- Sample documents (markdown format) to process
- Basic familiarity with Python (for SDK examples)
Architecture Overview
A RAG pipeline consists of four main components:
- Documents
- Ingestion & Chunking
- Embedding Generation
- Vector Storage
- Similarity Search
- LLM Generation
- Response
Step 1: Deploy Required Models
First, we'll deploy an embedding model for vectorizing text and a language model for generating responses.
Deploy an Embedding Model
from kamiwaza_client import KamiwazaClient
client = KamiwazaClient(base_url="http://localhost:7777/api/")
embedder = client.embedding.get_embedder(
model="BAAI/bge-base-en-v1.5",
provider_type="huggingface_embedding"
)
print("✅ Embedding model ready for use")
Deploy a Language Model
Deploy a language model using Kamiwaza for response generation:
from kamiwaza_client import KamiwazaClient
client = KamiwazaClient(base_url="http://localhost:7777/api/")
model_repo = "Qwen/Qwen3-0.6B-GGUF"
models = client.models.search_models(model_repo, exact=True)
print(f"Found model: {models[0]}")
client.models.initiate_model_download(model_repo)
client.models.wait_for_download(model_repo)
deployment_id = client.serving.deploy_model(repo_id=model_repo)
print(f"✅ Model deployed with ID: {deployment_id}")
openai_client = client.openai.get_client(repo_id=model_repo)
print("✅ OpenAI-compatible client ready")
Check Deployment Status
# List active deployments to verify
deployments = client.serving.list_active_deployments()
for deployment in deployments:
print(f"✅ {deployment.m_name} is {deployment.status}")
print(f" Endpoint: {deployment.endpoint}")
Step 2: Document Ingestion Pipeline
Now we'll create a pipeline to process documents, chunk them, and generate embeddings.
Document Processing Script
import os
from pathlib import Path
from typing import List, Dict
from kamiwaza_client import KamiwazaClient
class RAGPipeline:
def __init__(self, base_url="http://localhost:7777/api/"):
self.client = KamiwazaClient(base_url=base_url)
self.embedding_model = "BAAI/bge-base-en-v1.5"
self.collection_name = "documents"
self.embedder = self.client.embedding.get_embedder(
model=self.embedding_model,
provider_type="huggingface_embedding"
)
print(f"✅ RAG Pipeline initialized with model: {self.embedding_model}")
def add_documents_to_catalog(self, filepaths: List[str]) -> List:
datasets = []
for filepath in filepaths:
try:
dataset = self.client.catalog.create_dataset(
dataset_name=filepath,
platform="file",
environment="PROD",
description=f"RAG document: {Path(filepath).name}"
)
if dataset.urn:
datasets.append(dataset)
print(f"✅ Added to catalog: {Path(filepath).name}")
except Exception as e:
print(f"❌ Error adding {filepath}: {str(e)}")
return datasets
def process_document(self, file_path: str):
doc_path = Path(file_path)
if not doc_path.exists():
raise FileNotFoundError(f"File not found: {doc_path}")
with open(doc_path, 'r', encoding='utf-8') as f:
content = f.read()
print(f"📄 Processing document: {doc_path.name}")
chunks = self.embedder.chunk_text(
text=content,
max_length=1024,
overlap=102
)
embeddings = self.embedder.embed_chunks(chunks)
metadata_list = []
for i, chunk in enumerate(chunks):
chunk_text = chunk[:900] + "..." if len(chunk) > 900 else chunk
metadata = {
"model_name": self.embedding_model,
"source": str(doc_path),
"offset": i,
"filename": doc_path.name,
"chunk_text": chunk_text,
"chunk_index": i,
"chunk_size": len(chunk),
"document_title": doc_path.stem
}
metadata_list.append(metadata)
field_list = [
("chunk_text", "str"),
("chunk_index", "int"),
("chunk_size", "int"),
("document_title", "str")
]
self.client.vectordb.insert(
vectors=embeddings,
metadata=metadata_list,
collection_name=self.collection_name,
field_list=field_list
)
print(f"✅ Successfully stored {len(chunks)} chunks in collection '{self.collection_name}'")
return len(chunks)
# Usage
pipeline = RAGPipeline()
DOCUMENT_PATHS = [
"./docs/intro.md",
"./docs/models/overview.md",
"./docs/architecture/overview.md",
"./docs/architecture/architecture.md",
"./docs/architecture/components.md"
]
datasets = pipeline.add_documents_to_catalog(DOCUMENT_PATHS)
total_chunks = 0
for doc_path in DOCUMENT_PATHS:
try:
chunks = pipeline.process_document(doc_path)
total_chunks += chunks
except Exception as e:
print(f"❌ Error processing {doc_path}: {str(e)}")
print(f"\n🎉 Total chunks processed: {total_chunks}")
Step 3: Implement Retrieval and Generation
Now we'll create the query interface that retrieves relevant documents and generates responses.
from typing import List, Dict
from kamiwaza_client import KamiwazaClient
class RAGQuery:
def __init__(self, base_url="http://localhost:7777/api/", chat_model_repo="Qwen/Qwen3-0.6B-GGUF"):
self.client = KamiwazaClient(base_url=base_url)
self.embedding_model = "BAAI/bge-base-en-v1.5"
self.chat_model_repo = chat_model_repo
self.collection_name = "documents"
self.embedder = self.client.embedding.get_embedder(
model=self.embedding_model,
provider_type="huggingface_embedding"
)
try:
self.openai_client = self.client.openai.get_client(repo_id=self.chat_model_repo)
print(f"✅ RAG Query system initialized with chat model: {self.chat_model_repo}")
except Exception as e:
print(f"⚠️ Warning: Could not initialize chat model client: {e}")
self.openai_client = None
def semantic_search(self, query: str, limit: int = 5) -> List[Dict]:
print(f"🔍 Searching for: '{query}'")
query_embedding = self.embedder.create_embedding(query).embedding
results = self.client.vectordb.search(
query_vector=query_embedding,
collection_name=self.collection_name,
limit=limit,
output_fields=[
"source", "offset", "filename", "model_name",
"chunk_text", "chunk_index", "chunk_size", "document_title"
]
)
return results
def format_context(self, search_results: List[Dict]) -> str:
context_parts = []
for result in search_results:
metadata = result.get('metadata', {})
chunk_text = metadata.get('chunk_text', '')
document_title = metadata.get('document_title', 'Unknown')
if chunk_text:
context_parts.append(f"Document: {document_title}\n{chunk_text}")
return "\n\n".join(context_parts)
def generate_response(self, query: str, context: str) -> str:
if not self.openai_client:
return f"[Error: Chat model not available. Please deploy {self.chat_model_repo} first]"
prompt = f"""Based on the following context, answer the user's question. If the context doesn't contain enough information to answer the question, say so.\n\nContext:\n{context}\n\nQuestion: {query}\nAnswer:"""
response = self.openai_client.chat.completions.create(
messages=[{"role": "user", "content": prompt}],
model="model",
max_tokens=500,
temperature=0.7,
stream=False
)
return response.choices[0].message.content
def query(self, user_question: str, limit: int = 5) -> Dict:
search_results = self.semantic_search(user_question, limit=limit)
context = self.format_context(search_results)
response = self.generate_response(user_question, context)
sources = []
for result in search_results:
metadata = result.get('metadata', {})
sources.append({
'filename': metadata.get('filename', 'Unknown'),
'document_title': metadata.get('document_title', ''),
'chunk_index': metadata.get('chunk_index', 0),
'score': result.get('score', 0.0)
})
return {
'question': user_question,
'answer': response,
'context': context,
'sources': sources,
'num_results': len(search_results)
}
Step 4: Example Queries
Let's test the RAG system with an example query to demonstrate its capabilities:
rag = RAGQuery()
query_response = rag.query("What is one cool thing about Kamiwaza?")
print(query_response['answer'])
Step 5: Production Considerations
When moving your RAG system to production, consider these key aspects:
Resource Management
def monitor_system_health():
client = KamiwazaClient(base_url="http://localhost:7777/api/")
deployments = client.serving.list_active_deployments()
for deployment in deployments:
print(f" - {deployment.m_name}: {deployment.status}")
collections = client.vectordb.list_collections()
# Run health check
monitor_system_health()
Clean up
When done, stop the model deployment to free resources:
def cleanup_rag_system(chat_model_repo="Qwen/Qwen3-0.6B-GGUF"):
client = KamiwazaClient(base_url="http://localhost:7777/api/")
success = client.serving.stop_deployment(repo_id=chat_model_repo)
if success:
print(f"✅ Stopped deployment for {chat_model_repo}")
else:
print(f"❌ Failed to stop deployment for {chat_model_repo}")
cleanup_rag_system()
Best Practices
- Chunk Size: Keep chunks between 200-800 tokens for optimal retrieval
- Overlap: Add 50-100 token overlap between chunks to preserve context
- Metadata: Include rich metadata (source, date, author) for filtering
- Preprocessing: Clean text, remove headers/footers, handle special characters
Troubleshooting
Common Issues
Poor Retrieval Quality
- Check embedding model performance on your domain
- Adjust chunk size and overlap
- Consider domain-specific fine-tuning
Next Steps
Now that you have a working RAG pipeline with Kamiwaza-deployed models:
- Optimize Retrieval: Experiment with different embedding models, chunk sizes, and similarity thresholds
- Scale Your System: Deploy multiple model instances using Kamiwaza's distributed architecture
- Monitor Performance: Add logging and metrics to track query performance and model usage
Your RAG pipeline is now ready to answer questions using your own documents! The combination of Kamiwaza's SDK with proper document processing creates a robust foundation for production RAG applications.