Building a RAG Pipeline | Kamiwaza Docs
Documentation for Kamiwaza 0.9.0
This is documentation for Kamiwaza 0.9.0, which is no longer actively maintained. For the current GA release, see 1.0.1.
Version: 0.9.0
Overview
Retrieval-Augmented Generation (RAG) combines the power of large language models with your own documents and data to provide accurate, contextual responses. This guide walks you through building a complete 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
The embedding model will be automatically loaded when you create an embedder - no manual deployment needed:
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" # Small efficient model
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)
print("✅ Model download complete")
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" # Use a proven working model
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)
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]:
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['metadata'] if isinstance(result, dict) and 'metadata' in result else {}
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 "[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.
Context:
{context}
Question: {query}
Answer:"""
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['metadata'] if isinstance(result, dict) and 'metadata' in result else {}
score = result.get('score', 0.0)
sources.append({
'filename': metadata.get('filename', 'Unknown'),
'document_title': metadata.get('document_title', ''),
'chunk_index': metadata.get('chunk_index', 0),
'score': score
})
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:
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
# Monitor system resources and manage deployments
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()
print(f"📚 Vector Collections: {collections}")
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)
print(f"✅ Stopped deployment for {chat_model_repo}") if success 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
Slow Query Performance
- Optimize vector index parameters
- Reduce
top_kin retrieval - Use GPU acceleration for embeddings
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.