Building a RAG Pipeline | Kamiwaza Docs

Version: 1.0.1 (Latest)

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:

Prerequisites

Before starting, ensure you have:

Architecture Overview

A RAG pipeline consists of four main components:

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_sdk import KamiwazaClient

client = KamiwazaClient(base_url="https://your-kamiwaza.example/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_sdk import KamiwazaClient

client = KamiwazaClient(base_url="https://your-kamiwaza.example/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]}")
print("Downloading model...")
client.models.initiate_model_download(model_repo)
client.models.wait_for_download(model_repo)
print("✅ Model download complete")
print("Deploying model...")
deployment_id = client.serving.deploy_model(repo_id=model_repo)
print(f"✅ Model deployed with ID: {deployment_id}")
# Get OpenAI-compatible client for the deployed model
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_sdk import KamiwazaClient

class RAGPipeline:
    def __init__(self, base_url="https://your-kamiwaza.example/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}")
        print(f"   - Size: {len(content)} characters")
        chunks = self.embedder.chunk_text(
            text=content,
            max_length=1024,
            overlap=102
        )
        print(f"   - Created {len(chunks)} chunks")
        embeddings = self.embedder.embed_chunks(chunks)
        print(f"   - Generated {len(embeddings)} embeddings")
        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_sdk import KamiwazaClient
class RAGQuery:
    def __init__(self, base_url="https://your-kamiwaza.example/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 hasattr(result, 'metadata') 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 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}\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.metadata if hasattr(result, 'metadata') else {}
            sources.append({
                'filename': metadata.get('filename', 'Unknown'),
                'document_title': metadata.get('document_title', ''),
                'chunk_index': metadata.get('chunk_index', 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="https://your-kamiwaza.example/api")
    deployments = client.serving.list_active_deployments()
    for deployment in deployments:
        print(f"   - {deployment.m_name}: {deployment.status}")
    collections = client.vectordb.list_collections()

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="https://your-kamiwaza.example/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

Document Processing

Vector Search Optimization

LLM Integration

Troubleshooting

Common Issues

Poor Retrieval Quality

Slow Query Performance

Inaccurate Responses

Next Steps

Now that you have a working RAG pipeline with Kamiwaza-deployed models:

Key Benefits of This SDK-Based Approach

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.