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:

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_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

Troubleshooting

Common Issues

Poor Retrieval Quality

Slow Query Performance

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.