PART ONE: UNDERSTANDING THE LANDSCAPE - WHAT ARE WE BUILDING AND WHY?
Welcome to your journey into modern application deployment! If you have ever wondered how companies like Netflix, Spotify, or Google run applications that serve millions of users simultaneously without crashing, you are about to discover their secret: containerization and orchestration.
In this tutorial, we will build something genuinely useful and modern: a Retrieval-Augmented Generation chatbot powered by Large Language Models. This is not just a toy example. You will create a production-ready system that combines cutting-edge AI technology with industry-standard deployment practices. By the end of this tutorial, you will understand how to build, containerize, and orchestrate complex applications using Docker and Kubernetes.
THE PROBLEM WE ARE SOLVING
Imagine you are building a chatbot that can answer questions about your company's documentation. A traditional chatbot might give generic responses or hallucinate information. A RAG-based LLM chatbot solves this by retrieving relevant documents first, then using them to generate accurate, context-aware responses. This is the same technology powering modern AI assistants and enterprise knowledge systems.
But here is the challenge: such a system has multiple moving parts. You need a service to process documents, another to store and search through them, a third to run the AI model, and a fourth to coordinate everything. Each part needs to scale independently, recover from failures automatically, and be updated without taking the whole system down. This is where Docker and Kubernetes shine.
WHAT IS DOCKER AND WHY DO WE NEED IT?
Think of Docker as a way to package your application with everything it needs to run: the code, the runtime environment, the libraries, and the configuration files. This package is called a container. Unlike traditional applications that might work on your machine but fail on a server because of different installed libraries or configurations, a Docker container runs identically everywhere.
A container is similar to a shipping container in the physical world. Just as shipping containers standardized global trade by providing a consistent package that works on any ship, truck, or train, Docker containers standardize software deployment by providing a consistent package that works on any computer, server, or cloud platform.
The key difference between containers and virtual machines is efficiency. A virtual machine includes an entire operating system, which can be gigabytes in size and take minutes to start. A container shares the host operating system's kernel and includes only your application and its dependencies, making it lightweight (often just megabytes) and fast to start (often just seconds).
WHAT IS KUBERNETES AND WHY DO WE NEED IT?
If Docker is about packaging applications, Kubernetes is about running them at scale. Kubernetes is an orchestration platform that manages containers across multiple machines. It handles deployment, scaling, networking, and recovery automatically.
Imagine you have deployed your chatbot and suddenly thousands of users start asking questions simultaneously. Without Kubernetes, you would need to manually start more containers, configure load balancing, and monitor everything. With Kubernetes, you simply tell it "I want five instances of this service" and it handles the rest. If a container crashes, Kubernetes automatically restarts it. If a server fails, Kubernetes moves the containers to healthy servers. If traffic increases, Kubernetes can automatically scale up your application.
Kubernetes introduces several key concepts that we will explore throughout this tutorial. A Pod is the smallest deployable unit in Kubernetes, typically containing one or more closely related containers. A Deployment manages Pods and ensures the desired number of replicas are always running. A Service provides a stable network endpoint to access Pods, even as they are created and destroyed. An Ingress manages external access to services, typically HTTP and HTTPS traffic.
THE ARCHITECTURE WE WILL BUILD
Our RAG-based LLM chatbot will consist of four microservices, each running in its own Docker container and managed by Kubernetes. Understanding why we split the application this way is crucial to understanding microservices architecture.
The Document Ingestion Service will be responsible for receiving documents, splitting them into chunks, and preparing them for embedding. This service needs to handle file uploads and text processing, which are CPU-intensive tasks. By isolating this functionality, we can scale it independently when users upload many documents.
The Vector Database Service will store document embeddings and perform similarity searches. This service requires significant memory to hold vector indexes and needs to respond quickly to search queries. By making it a separate service, we can deploy it on memory-optimized hardware and scale it based on query load.
The LLM Inference Service will run the actual language model to generate responses. This is the most resource-intensive component, requiring GPU acceleration for acceptable performance. By isolating it, we can deploy it on GPU-enabled nodes and scale it based on generation requests. We will optimize this service to work with different GPU types including NVIDIA CUDA, AMD ROCm, and Apple Silicon.
The API Gateway Service will coordinate requests between the client and the other services. When a user asks a question, this service retrieves relevant documents from the vector database, constructs a prompt with the context, sends it to the LLM service, and returns the response. This service handles the orchestration logic and can be scaled based on user request volume.
Each of these services will be independently deployable, scalable, and maintainable. This is the essence of microservices architecture: breaking a complex system into smaller, focused components that work together through well-defined interfaces.
PART TWO: DOCKER FUNDAMENTALS - CONTAINERS AND IMAGES
Before we build our microservices, we need to understand Docker's core concepts: images and containers. An image is a blueprint, a read-only template that contains everything needed to run an application. A container is a running instance of an image. You can create multiple containers from the same image, just as you can bake multiple cakes from the same recipe.
YOUR FIRST DOCKERFILE
A Dockerfile is a text file containing instructions to build a Docker image. Let us start with the simplest possible example: a Python application that prints a message. Create a file named hello.py with the following content:
# hello.py
# A simple Python script to demonstrate Docker basics
def main():
print("Hello from inside a Docker container!")
print("This application is completely isolated from your host system.")
if __name__ == "__main__":
main()
Now create a file named Dockerfile in the same directory. The Dockerfile contains instructions that Docker will execute to build your image:
# Dockerfile for a simple Python application
# Each instruction creates a new layer in the image
# FROM specifies the base image to start from
# We use Python 3.11 on Alpine Linux for a small image size
FROM python:3.11-alpine
# WORKDIR sets the working directory inside the container
# All subsequent commands will run in this directory
WORKDIR /app
# COPY transfers files from your host to the container
# The first argument is the source on your host
# The second argument is the destination in the container
COPY hello.py .
# CMD specifies the command to run when the container starts
# This is the entry point of your application
CMD ["python", "hello.py"]
Let us examine each instruction carefully. The FROM instruction specifies the base image. Instead of starting from scratch, we build on top of an existing image that already has Python installed. The python:3.11-alpine image is a minimal Linux distribution with Python 3.11, keeping our final image small.
The WORKDIR instruction sets the working directory. This is where your application files will live inside the container. If the directory does not exist, Docker creates it automatically.
The COPY instruction transfers files from your development machine into the container image. Here we copy hello.py from the current directory on your host into the /app directory in the container.
The CMD instruction defines what command runs when someone starts a container from this image. In this case, we run Python with our hello.py script.
To build this image, open a terminal in the directory containing your Dockerfile and run:
docker build -t my-first-app .
The -t flag tags the image with a name (my-first-app) so you can reference it easily. The dot at the end tells Docker to look for the Dockerfile in the current directory. Docker will execute each instruction in the Dockerfile, creating a new layer for each one. You will see output showing the progress of each step.
Once the build completes, you can run a container from this image:
docker run my-first-app
You should see the output from your Python script. Congratulations! You have just containerized your first application. The container started, executed your code, and stopped. Everything happened in complete isolation from your host system.
UNDERSTANDING LAYERS AND CACHING
Docker images are built in layers, and understanding this is crucial for creating efficient images. Each instruction in your Dockerfile creates a new layer. Docker caches these layers, so if you rebuild an image and nothing has changed in a layer, Docker reuses the cached version instead of rebuilding it.
This has important implications for how you structure your Dockerfile. Consider a more realistic Python application that has dependencies. Create a file named requirements.txt:
flask==3.0.0
requests==2.31.0
Now create a simple Flask web application in app.py:
# app.py
# A simple Flask web server to demonstrate Docker networking
from flask import Flask, jsonify
import os
# Create a Flask application instance
app = Flask(__name__)
# Define a route that responds to HTTP GET requests
@app.route('/health')
def health():
"""Health check endpoint for monitoring"""
return jsonify({
'status': 'healthy',
'service': 'demo-service',
'container_id': os.environ.get('HOSTNAME', 'unknown')
})
@app.route('/')
def home():
"""Main endpoint that returns a welcome message"""
return jsonify({
'message': 'Welcome to the containerized Flask app!',
'endpoints': ['/health', '/']
})
if __name__ == '__main__':
# Run the Flask development server
# host='0.0.0.0' makes it accessible from outside the container
# port=5000 is the standard Flask port
app.run(host='0.0.0.0', port=5000, debug=True)
Here is an inefficient Dockerfile for this application:
# Inefficient Dockerfile - DO NOT USE THIS APPROACH
FROM python:3.11-alpine
WORKDIR /app
# This copies everything first
COPY . .
# Then installs dependencies
RUN pip install --no-cache-dir -r requirements.txt
CMD ["python", "app.py"]
The problem with this approach is that every time you change any file in your project, Docker invalidates the cache for the COPY instruction and all subsequent instructions. This means Docker will reinstall all your dependencies even if requirements.txt has not changed, which can take a long time for projects with many dependencies.
Here is a better approach:
# Efficient Dockerfile using layer caching
FROM python:3.11-alpine
WORKDIR /app
# Copy only the requirements file first
COPY requirements.txt .
# Install dependencies in a separate layer
# This layer will be cached unless requirements.txt changes
RUN pip install --no-cache-dir -r requirements.txt
# Copy the application code
# Changes to app.py won't invalidate the dependency layer
COPY app.py .
# Expose the port the app runs on
# This is documentation; it doesn't actually publish the port
EXPOSE 5000
CMD ["python", "app.py"]
By copying requirements.txt and installing dependencies before copying the application code, we ensure that the dependency installation layer is only rebuilt when requirements.txt changes. Changes to app.py will only require rebuilding the final COPY layer, making builds much faster during development.
Build this image:
docker build -t flask-demo .
Run it with port mapping so you can access the web server from your host:
docker run -p 5000:5000 flask-demo
The -p flag maps port 5000 on your host to port 5000 in the container. Now you can open a web browser and navigate to http://localhost:5000 to see your containerized web application running. Press Ctrl+C to stop the container.
MULTI-STAGE BUILDS FOR PRODUCTION
For production applications, we want the smallest possible images to reduce attack surface, download time, and storage costs. Multi-stage builds allow you to use one image for building your application and a different, smaller image for running it.
Consider a Go application that needs to be compiled. Create a file named main.go:
// main.go
// A simple HTTP server in Go to demonstrate multi-stage builds
package main
import (
"encoding/json"
"fmt"
"log"
"net/http"
)
// Response structure for JSON output
type Response struct {
Message string `json:"message"`
Version string `json:"version"`
}
// Handler for the root endpoint
func homeHandler(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
response := Response{
Message: "Hello from a multi-stage Docker build!",
Version: "1.0.0",
}
json.NewEncoder(w).Encode(response)
}
func main() {
http.HandleFunc("/", homeHandler)
fmt.Println("Server starting on port 8080...")
log.Fatal(http.ListenAndServe(":8080", nil))
}
Here is a multi-stage Dockerfile for this Go application:
# Multi-stage Dockerfile for Go application
# Stage 1: Build stage
# Use the full Go image which includes the compiler and build tools
FROM golang:1.21-alpine AS builder
# Install any build dependencies
RUN apk add --no-cache git
WORKDIR /build
# Copy go module files and download dependencies
# In a real project, you would have go.mod and go.sum files
COPY main.go .
# Build the application
# CGO_ENABLED=0 creates a statically linked binary
# -ldflags="-w -s" strips debug information to reduce size
RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-w -s" -o server main.go
# Stage 2: Runtime stage
# Use a minimal base image for the final container
FROM alpine:latest
# Add ca-certificates for HTTPS requests
RUN apk --no-cache add ca-certificates
WORKDIR /app
# Copy only the compiled binary from the builder stage
# This is the key to multi-stage builds: we leave behind all build tools
COPY --from=builder /build/server .
# Create a non-root user for security
RUN addgroup -g 1000 appuser && \
adduser -D -u 1000 -G appuser appuser && \
chown -R appuser:appuser /app
# Switch to the non-root user
USER appuser
EXPOSE 8080
CMD ["./server"]
This Dockerfile has two FROM instructions, creating two stages. The first stage uses the full golang:1.21-alpine image which includes the Go compiler and all build tools. We compile our application in this stage. The second stage uses the minimal alpine:latest image and copies only the compiled binary from the first stage using COPY --from=builder.
The result is a final image that contains only the runtime dependencies and the compiled binary, without any of the build tools. This can reduce image size from hundreds of megabytes to just tens of megabytes. Additionally, we create a non-root user and run the application as that user, following security best practices.
Build and run this image:
docker build -t go-multistage .
docker run -p 8080:8080 go-multistage
You can verify the image size using:
docker images go-multistage
You will see that the final image is significantly smaller than if we had included all the build tools.
PART THREE: BUILDING MICROSERVICES - OUR RAG CHATBOT COMPONENTS
Now that you understand Docker basics, let us build the microservices for our RAG-based LLM chatbot. We will start with the simplest service and progressively add complexity, ensuring you understand each component before moving to the next.
MICROSERVICE ONE: THE DOCUMENT INGESTION SERVICE
The document ingestion service receives documents, splits them into chunks, and prepares them for embedding. This service exposes a REST API that accepts document uploads and processes them into manageable pieces.
Create a directory structure for this service:
document-ingestion/
app.py
requirements.txt
Dockerfile
.dockerignore
The .dockerignore file tells Docker which files to exclude from the build context, similar to .gitignore for Git:
__pycache__
*.pyc
*.pyo
*.pyd
.Python
env/
venv/
.pytest_cache
.coverage
*.log
The requirements.txt file lists our Python dependencies:
flask==3.0.0
werkzeug==3.0.1
python-multipart==0.0.6
nltk==3.8.1
tiktoken==0.5.2
Now let us create the application in app.py. This will be more substantial than our previous examples because it implements real functionality:
# app.py
# Document Ingestion Service for RAG-based LLM Chatbot
# This service receives documents and chunks them for embedding
from flask import Flask, request, jsonify
import os
import logging
from typing import List, Dict
import tiktoken
# Configure logging for production debugging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
# Create Flask application
app = Flask(__name__)
# Configuration from environment variables with sensible defaults
# This allows us to configure the service without rebuilding the container
CHUNK_SIZE = int(os.environ.get('CHUNK_SIZE', '512'))
CHUNK_OVERLAP = int(os.environ.get('CHUNK_OVERLAP', '50'))
MAX_FILE_SIZE = int(os.environ.get('MAX_FILE_SIZE', '10485760')) # 10MB default
class DocumentChunker:
"""
Handles the chunking of documents into smaller pieces.
This is necessary because LLMs have context length limits.
"""
def __init__(self, chunk_size: int, chunk_overlap: int):
"""
Initialize the chunker with size and overlap parameters.
Args:
chunk_size: Maximum number of tokens per chunk
chunk_overlap: Number of tokens to overlap between chunks
"""
self.chunk_size = chunk_size
self.chunk_overlap = chunk_overlap
# Initialize tokenizer for accurate token counting
self.tokenizer = tiktoken.get_encoding("cl100k_base")
logger.info(f"Initialized chunker: size={chunk_size}, overlap={chunk_overlap}")
def chunk_text(self, text: str, metadata: Dict = None) -> List[Dict]:
"""
Split text into overlapping chunks based on token count.
Args:
text: The input text to chunk
metadata: Optional metadata to attach to each chunk
Returns:
List of dictionaries containing chunk text and metadata
"""
if not text or not text.strip():
logger.warning("Received empty text for chunking")
return []
# Tokenize the entire text
tokens = self.tokenizer.encode(text)
total_tokens = len(tokens)
logger.info(f"Processing text with {total_tokens} tokens")
chunks = []
start_idx = 0
chunk_id = 0
while start_idx < total_tokens:
# Calculate end index for this chunk
end_idx = min(start_idx + self.chunk_size, total_tokens)
# Extract tokens for this chunk
chunk_tokens = tokens[start_idx:end_idx]
# Decode tokens back to text
chunk_text = self.tokenizer.decode(chunk_tokens)
# Create chunk object with metadata
chunk = {
'id': chunk_id,
'text': chunk_text,
'start_token': start_idx,
'end_token': end_idx,
'token_count': len(chunk_tokens),
'metadata': metadata or {}
}
chunks.append(chunk)
chunk_id += 1
# Move start index forward, accounting for overlap
start_idx = end_idx - self.chunk_overlap
# Prevent infinite loop if overlap is too large
if start_idx >= end_idx:
break
logger.info(f"Created {len(chunks)} chunks from text")
return chunks
# Initialize the chunker with configuration
chunker = DocumentChunker(CHUNK_SIZE, CHUNK_OVERLAP)
@app.route('/health', methods=['GET'])
def health_check():
"""
Health check endpoint for Kubernetes liveness and readiness probes.
Returns service status and configuration.
"""
return jsonify({
'status': 'healthy',
'service': 'document-ingestion',
'version': '1.0.0',
'config': {
'chunk_size': CHUNK_SIZE,
'chunk_overlap': CHUNK_OVERLAP,
'max_file_size': MAX_FILE_SIZE
}
}), 200
@app.route('/ingest', methods=['POST'])
def ingest_document():
"""
Main endpoint for document ingestion.
Accepts text or file uploads and returns chunked documents.
"""
try:
# Check if request contains JSON data with text
if request.is_json:
data = request.get_json()
text = data.get('text', '')
metadata = data.get('metadata', {})
if not text:
return jsonify({'error': 'No text provided'}), 400
# Process the text into chunks
chunks = chunker.chunk_text(text, metadata)
return jsonify({
'success': True,
'chunks': chunks,
'total_chunks': len(chunks)
}), 200
# Check if request contains file upload
elif 'file' in request.files:
file = request.files['file']
# Validate file size
file.seek(0, os.SEEK_END)
file_size = file.tell()
file.seek(0)
if file_size > MAX_FILE_SIZE:
return jsonify({
'error': f'File too large. Maximum size is {MAX_FILE_SIZE} bytes'
}), 400
# Read file content
text = file.read().decode('utf-8')
# Extract metadata from form data
metadata = {
'filename': file.filename,
'size': file_size
}
# Process the text into chunks
chunks = chunker.chunk_text(text, metadata)
return jsonify({
'success': True,
'chunks': chunks,
'total_chunks': len(chunks)
}), 200
else:
return jsonify({
'error': 'No text or file provided'
}), 400
except Exception as e:
logger.error(f"Error processing document: {str(e)}", exc_info=True)
return jsonify({
'error': 'Internal server error',
'message': str(e)
}), 500
@app.route('/config', methods=['GET'])
def get_config():
"""
Returns current service configuration.
Useful for debugging and monitoring.
"""
return jsonify({
'chunk_size': CHUNK_SIZE,
'chunk_overlap': CHUNK_OVERLAP,
'max_file_size': MAX_FILE_SIZE
}), 200
if __name__ == '__main__':
# Run the Flask development server
# In production, use a proper WSGI server like Gunicorn
port = int(os.environ.get('PORT', 5001))
app.run(host='0.0.0.0', port=port, debug=False)
This service implements several important concepts. First, it uses environment variables for configuration, allowing us to change settings without rebuilding the Docker image. Second, it implements proper error handling and logging, essential for production services. Third, it provides a health check endpoint that Kubernetes will use to monitor the service. Fourth, it chunks documents based on token count rather than character count, which is more accurate for LLM applications.
Now let us create the Dockerfile for this service:
# Dockerfile for Document Ingestion Service
# Uses multi-stage build for security and efficiency
FROM python:3.11-slim AS base
# Set environment variables to prevent Python from writing pyc files
# and to ensure output is sent straight to terminal without buffering
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
# Create a non-root user for security
RUN groupadd -r appuser && useradd -r -g appuser appuser
# Set working directory
WORKDIR /app
# Install dependencies in a separate layer for better caching
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt && \
python -m nltk.downloader punkt
# Copy application code
COPY app.py .
# Change ownership of application files to non-root user
RUN chown -R appuser:appuser /app
# Switch to non-root user
USER appuser
# Expose the application port
EXPOSE 5001
# Health check for Docker and Kubernetes
HEALTHCHECK --interval=30s --timeout=3s --start-period=5s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5001/health')" || exit 1
# Run the application
CMD ["python", "app.py"]
This Dockerfile follows best practices we discussed earlier. It uses a slim base image to reduce size, runs as a non-root user for security, and includes a health check. Build this image:
cd document-ingestion
docker build -t rag-chatbot/document-ingestion:1.0 .
Test the service locally:
docker run -p 5001:5001 rag-chatbot/document-ingestion:1.0
In another terminal, test the service using curl:
curl -X POST http://localhost:5001/ingest \
-H "Content-Type: application/json" \
-d '{"text": "This is a test document. It contains multiple sentences. We will chunk it into smaller pieces for embedding.", "metadata": {"source": "test"}}'
You should receive a JSON response containing the chunked text. This confirms that our first microservice is working correctly.
MICROSERVICE TWO: THE VECTOR DATABASE SERVICE
The vector database service stores document embeddings and performs similarity searches. For this tutorial, we will use a simple in-memory vector store, but in production you would use a dedicated vector database like Milvus, Weaviate, or Pinecone.
Create a directory structure:
vector-database/
app.py
requirements.txt
Dockerfile
.dockerignore
The requirements.txt file:
flask==3.0.0
numpy==1.26.2
scikit-learn==1.3.2
sentence-transformers==2.2.2
The application in app.py:
# app.py
# Vector Database Service for RAG-based LLM Chatbot
# This service stores embeddings and performs similarity search
from flask import Flask, request, jsonify
import numpy as np
from sentence_transformers import SentenceTransformer
from sklearn.metrics.pairwise import cosine_similarity
import logging
import os
from typing import List, Dict, Tuple
import threading
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
# Configuration
EMBEDDING_MODEL = os.environ.get('EMBEDDING_MODEL', 'all-MiniLM-L6-v2')
TOP_K = int(os.environ.get('TOP_K', '5'))
class VectorStore:
"""
In-memory vector store for document embeddings.
In production, use a dedicated vector database.
"""
def __init__(self, model_name: str):
"""
Initialize the vector store with an embedding model.
Args:
model_name: Name of the sentence-transformers model to use
"""
logger.info(f"Loading embedding model: {model_name}")
self.model = SentenceTransformer(model_name)
self.embeddings = []
self.documents = []
self.lock = threading.Lock()
logger.info("Vector store initialized successfully")
def add_documents(self, documents: List[Dict]) -> bool:
"""
Add documents to the vector store and generate embeddings.
Args:
documents: List of document dictionaries with 'text' and 'metadata'
Returns:
True if successful, False otherwise
"""
try:
with self.lock:
# Extract text from documents
texts = [doc['text'] for doc in documents]
# Generate embeddings
logger.info(f"Generating embeddings for {len(texts)} documents")
new_embeddings = self.model.encode(texts, show_progress_bar=False)
# Store embeddings and documents
self.embeddings.extend(new_embeddings)
self.documents.extend(documents)
logger.info(f"Added {len(documents)} documents. Total: {len(self.documents)}")
return True
except Exception as e:
logger.error(f"Error adding documents: {str(e)}", exc_info=True)
return False
def search(self, query: str, top_k: int = 5) -> List[Tuple[Dict, float]]:
"""
Search for documents similar to the query.
Args:
query: The search query text
top_k: Number of top results to return
Returns:
List of tuples containing (document, similarity_score)
"""
try:
with self.lock:
if not self.documents:
logger.warning("Search attempted on empty vector store")
return []
# Generate query embedding
query_embedding = self.model.encode([query], show_progress_bar=False)
# Calculate cosine similarity with all documents
embeddings_array = np.array(self.embeddings)
similarities = cosine_similarity(query_embedding, embeddings_array)[0]
# Get top-k indices
top_indices = np.argsort(similarities)[-top_k:][::-1]
# Prepare results
results = [
(self.documents[idx], float(similarities[idx]))
for idx in top_indices
]
logger.info(f"Search completed. Returned {len(results)} results")
return results
except Exception as e:
logger.error(f"Error during search: {str(e)}", exc_info=True)
return []
def get_stats(self) -> Dict:
"""
Get statistics about the vector store.
Returns:
Dictionary containing store statistics
"""
with self.lock:
return {
'total_documents': len(self.documents),
'embedding_dimension': len(self.embeddings[0]) if self.embeddings else 0,
'model': EMBEDDING_MODEL
}
# Initialize vector store
logger.info("Initializing vector store...")
vector_store = VectorStore(EMBEDDING_MODEL)
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint for Kubernetes probes."""
stats = vector_store.get_stats()
return jsonify({
'status': 'healthy',
'service': 'vector-database',
'version': '1.0.0',
'stats': stats
}), 200
@app.route('/add', methods=['POST'])
def add_documents():
"""
Add documents to the vector store.
Expects JSON with 'documents' array.
"""
try:
data = request.get_json()
if not data or 'documents' not in data:
return jsonify({'error': 'No documents provided'}), 400
documents = data['documents']
if not isinstance(documents, list):
return jsonify({'error': 'Documents must be a list'}), 400
# Validate document structure
for doc in documents:
if 'text' not in doc:
return jsonify({'error': 'Each document must have a text field'}), 400
# Add documents to vector store
success = vector_store.add_documents(documents)
if success:
return jsonify({
'success': True,
'added': len(documents),
'total': len(vector_store.documents)
}), 200
else:
return jsonify({'error': 'Failed to add documents'}), 500
except Exception as e:
logger.error(f"Error in add_documents: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/search', methods=['POST'])
def search():
"""
Search for similar documents.
Expects JSON with 'query' and optional 'top_k'.
"""
try:
data = request.get_json()
if not data or 'query' not in data:
return jsonify({'error': 'No query provided'}), 400
query = data['query']
top_k = data.get('top_k', TOP_K)
# Perform search
results = vector_store.search(query, top_k)
# Format results
formatted_results = [
{
'document': doc,
'similarity': score
}
for doc, score in results
]
return jsonify({
'success': True,
'query': query,
'results': formatted_results,
'count': len(formatted_results)
}), 200
except Exception as e:
logger.error(f"Error in search: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/stats', methods=['GET'])
def get_stats():
"""Get vector store statistics."""
stats = vector_store.get_stats()
return jsonify(stats), 200
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5002))
app.run(host='0.0.0.0', port=port, debug=False)
The Dockerfile for this service:
# Dockerfile for Vector Database Service
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Set ownership
RUN chown -R appuser:appuser /app
USER appuser
EXPOSE 5002
HEALTHCHECK --interval=30s --timeout=3s --start-period=40s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5002/health')" || exit 1
CMD ["python", "app.py"]
Note that we set a longer start-period for the health check because loading the embedding model takes more time than starting a simple web server. Build and test this service:
cd vector-database
docker build -t rag-chatbot/vector-database:1.0 .
docker run -p 5002:5002 rag-chatbot/vector-database:1.0
Test adding documents:
curl -X POST http://localhost:5002/add \
-H "Content-Type: application/json" \
-d '{"documents": [{"text": "Docker is a containerization platform", "metadata": {"topic": "docker"}}, {"text": "Kubernetes orchestrates containers", "metadata": {"topic": "kubernetes"}}]}'
Test searching:
curl -X POST http://localhost:5002/search \
-H "Content-Type: application/json" \
-d '{"query": "What is container orchestration?", "top_k": 2}'
You should receive results showing the most similar documents to your query.
MICROSERVICE THREE: THE LLM INFERENCE SERVICE WITH GPU OPTIMIZATION
The LLM inference service is the most complex component because it needs to handle different GPU types and optimize inference performance. We will create a service that can work with NVIDIA CUDA, AMD ROCm, and Apple Silicon, as well as CPU-only environments.
Create the directory structure:
llm-inference/
app.py
requirements.txt
Dockerfile.cuda
Dockerfile.rocm
Dockerfile.cpu
.dockerignore
We create multiple Dockerfiles because each GPU type requires different base images and dependencies. Let us start with requirements.txt:
flask==3.0.0
transformers==4.36.0
torch==2.1.0
accelerate==0.25.0
bitsandbytes==0.41.3
sentencepiece==0.1.99
The application in app.py:
# app.py
# LLM Inference Service for RAG-based Chatbot
# Supports NVIDIA CUDA, AMD ROCm, Apple Silicon, and CPU
from flask import Flask, request, jsonify
import torch
from transformers import AutoModelForCausalLM, AutoTokenizer, BitsAndBytesConfig
import logging
import os
from typing import Dict, Optional
import platform
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
# Configuration
MODEL_NAME = os.environ.get('MODEL_NAME', 'TinyLlama/TinyLlama-1.1B-Chat-v1.0')
MAX_LENGTH = int(os.environ.get('MAX_LENGTH', '512'))
TEMPERATURE = float(os.environ.get('TEMPERATURE', '0.7'))
USE_QUANTIZATION = os.environ.get('USE_QUANTIZATION', 'true').lower() == 'true'
def detect_device() -> str:
"""
Detect the best available device for inference.
Returns:
Device string: 'cuda', 'mps', 'rocm', or 'cpu'
"""
if torch.cuda.is_available():
device = 'cuda'
gpu_name = torch.cuda.get_device_name(0)
logger.info(f"NVIDIA CUDA available: {gpu_name}")
elif hasattr(torch.backends, 'mps') and torch.backends.mps.is_available():
device = 'mps'
logger.info("Apple Silicon (MPS) available")
elif torch.version.hip:
device = 'cuda' # ROCm uses CUDA API
logger.info("AMD ROCm available")
else:
device = 'cpu'
logger.info("No GPU available, using CPU")
return device
class LLMInferenceEngine:
"""
Handles LLM model loading and inference with GPU optimization.
"""
def __init__(self, model_name: str, device: str, use_quantization: bool):
"""
Initialize the inference engine.
Args:
model_name: HuggingFace model identifier
device: Device to run inference on
use_quantization: Whether to use 4-bit quantization
"""
self.model_name = model_name
self.device = device
self.use_quantization = use_quantization and device in ['cuda', 'mps']
logger.info(f"Loading model: {model_name}")
logger.info(f"Device: {device}")
logger.info(f"Quantization: {self.use_quantization}")
# Load tokenizer
self.tokenizer = AutoTokenizer.from_pretrained(model_name)
# Set padding token if not set
if self.tokenizer.pad_token is None:
self.tokenizer.pad_token = self.tokenizer.eos_token
# Configure quantization for memory efficiency
if self.use_quantization:
quantization_config = BitsAndBytesConfig(
load_in_4bit=True,
bnb_4bit_compute_dtype=torch.float16,
bnb_4bit_use_double_quant=True,
bnb_4bit_quant_type="nf4"
)
self.model = AutoModelForCausalLM.from_pretrained(
model_name,
quantization_config=quantization_config,
device_map="auto",
trust_remote_code=True
)
else:
self.model = AutoModelForCausalLM.from_pretrained(
model_name,
torch_dtype=torch.float16 if device != 'cpu' else torch.float32,
trust_remote_code=True
)
self.model.to(device)
# Set model to evaluation mode
self.model.eval()
logger.info("Model loaded successfully")
def generate(
self,
prompt: str,
max_length: int = 512,
temperature: float = 0.7,
top_p: float = 0.9,
top_k: int = 50
) -> str:
"""
Generate text based on the prompt.
Args:
prompt: Input prompt for generation
max_length: Maximum length of generated text
temperature: Sampling temperature
top_p: Nucleus sampling parameter
top_k: Top-k sampling parameter
Returns:
Generated text
"""
try:
# Tokenize input
inputs = self.tokenizer(
prompt,
return_tensors="pt",
truncation=True,
max_length=max_length
)
# Move inputs to device
inputs = {k: v.to(self.device) for k, v in inputs.items()}
# Generate with optimized parameters
with torch.no_grad():
outputs = self.model.generate(
**inputs,
max_new_tokens=max_length,
temperature=temperature,
top_p=top_p,
top_k=top_k,
do_sample=True,
pad_token_id=self.tokenizer.pad_token_id,
eos_token_id=self.tokenizer.eos_token_id
)
# Decode output
generated_text = self.tokenizer.decode(
outputs[0],
skip_special_tokens=True
)
# Remove the prompt from the output
if generated_text.startswith(prompt):
generated_text = generated_text[len(prompt):].strip()
return generated_text
except Exception as e:
logger.error(f"Error during generation: {str(e)}", exc_info=True)
raise
def get_model_info(self) -> Dict:
"""
Get information about the loaded model.
Returns:
Dictionary with model information
"""
return {
'model_name': self.model_name,
'device': self.device,
'quantization': self.use_quantization,
'parameters': sum(p.numel() for p in self.model.parameters()),
'dtype': str(next(self.model.parameters()).dtype)
}
# Initialize inference engine
logger.info("Initializing LLM inference engine...")
device = detect_device()
inference_engine = LLMInferenceEngine(MODEL_NAME, device, USE_QUANTIZATION)
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint."""
model_info = inference_engine.get_model_info()
return jsonify({
'status': 'healthy',
'service': 'llm-inference',
'version': '1.0.0',
'model': model_info
}), 200
@app.route('/generate', methods=['POST'])
def generate():
"""
Generate text based on input prompt.
Expects JSON with 'prompt' and optional generation parameters.
"""
try:
data = request.get_json()
if not data or 'prompt' not in data:
return jsonify({'error': 'No prompt provided'}), 400
prompt = data['prompt']
max_length = data.get('max_length', MAX_LENGTH)
temperature = data.get('temperature', TEMPERATURE)
top_p = data.get('top_p', 0.9)
top_k = data.get('top_k', 50)
# Generate response
logger.info(f"Generating response for prompt of length {len(prompt)}")
generated_text = inference_engine.generate(
prompt=prompt,
max_length=max_length,
temperature=temperature,
top_p=top_p,
top_k=top_k
)
return jsonify({
'success': True,
'prompt': prompt,
'generated_text': generated_text,
'model': MODEL_NAME
}), 200
except Exception as e:
logger.error(f"Error in generate: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/model-info', methods=['GET'])
def model_info():
"""Get detailed model information."""
info = inference_engine.get_model_info()
return jsonify(info), 200
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5003))
app.run(host='0.0.0.0', port=port, debug=False)
Now let us create the Dockerfiles for different GPU types. First, Dockerfile.cuda for NVIDIA GPUs:
# Dockerfile.cuda
# LLM Inference Service optimized for NVIDIA CUDA GPUs
FROM nvidia/cuda:12.1.0-runtime-ubuntu22.04
# Install Python and system dependencies
RUN apt-get update && apt-get install -y \
python3.11 \
python3-pip \
&& rm -rf /var/lib/apt/lists/*
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
CUDA_VISIBLE_DEVICES=0
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install PyTorch with CUDA support
RUN pip3 install --no-cache-dir \
torch==2.1.0 \
--index-url https://download.pytorch.org/whl/cu121
# Install other dependencies
COPY requirements.txt .
RUN pip3 install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create cache directory for HuggingFace models
RUN mkdir -p /app/.cache && chown -R appuser:appuser /app
USER appuser
EXPOSE 5003
# Longer start period because model loading takes time
HEALTHCHECK --interval=30s --timeout=5s --start-period=120s --retries=3 \
CMD python3 -c "import requests; requests.get('http://localhost:5003/health')" || exit 1
CMD ["python3", "app.py"]
Dockerfile.rocm for AMD GPUs:
# Dockerfile.rocm
# LLM Inference Service optimized for AMD ROCm GPUs
FROM rocm/pytorch:rocm6.0_ubuntu22.04_py3.10_pytorch_2.1.1
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
HSA_OVERRIDE_GFX_VERSION=11.0.0
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create cache directory
RUN mkdir -p /app/.cache && chown -R appuser:appuser /app
USER appuser
EXPOSE 5003
HEALTHCHECK --interval=30s --timeout=5s --start-period=120s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5003/health')" || exit 1
CMD ["python", "app.py"]
Dockerfile.cpu for CPU-only environments:
# Dockerfile.cpu
# LLM Inference Service for CPU-only environments
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create cache directory
RUN mkdir -p /app/.cache && chown -R appuser:appuser /app
USER appuser
EXPOSE 5003
HEALTHCHECK --interval=30s --timeout=5s --start-period=120s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5003/health')" || exit 1
CMD ["python", "app.py"]
Build the appropriate image for your hardware. For NVIDIA GPUs:
cd llm-inference
docker build -f Dockerfile.cuda -t rag-chatbot/llm-inference:1.0-cuda .
For CPU-only:
docker build -f Dockerfile.cpu -t rag-chatbot/llm-inference:1.0-cpu .
Note that running this service requires significant resources. The model download and loading can take several minutes on the first run.
MICROSERVICE FOUR: THE API GATEWAY SERVICE
The API gateway orchestrates requests between all services. It receives user queries, retrieves relevant documents from the vector database, constructs prompts with context, sends them to the LLM service, and returns responses.
Create the directory structure:
api-gateway/
app.py
requirements.txt
Dockerfile
.dockerignore
The requirements.txt file:
flask==3.0.0
requests==2.31.0
flask-cors==4.0.0
The application in app.py:
# app.py
# API Gateway Service for RAG-based LLM Chatbot
# Orchestrates requests between all microservices
from flask import Flask, request, jsonify
from flask_cors import CORS
import requests
import logging
import os
from typing import List, Dict, Optional
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
CORS(app) # Enable CORS for client applications
# Service endpoints from environment variables
INGESTION_SERVICE = os.environ.get('INGESTION_SERVICE', 'http://localhost:5001')
VECTOR_SERVICE = os.environ.get('VECTOR_SERVICE', 'http://localhost:5002')
LLM_SERVICE = os.environ.get('LLM_SERVICE', 'http://localhost:5003')
# Configuration
TOP_K_DOCUMENTS = int(os.environ.get('TOP_K_DOCUMENTS', '3'))
REQUEST_TIMEOUT = int(os.environ.get('REQUEST_TIMEOUT', '30'))
class RAGOrchestrator:
"""
Orchestrates the RAG pipeline across microservices.
"""
def __init__(
self,
ingestion_url: str,
vector_url: str,
llm_url: str,
timeout: int
):
"""
Initialize the orchestrator with service URLs.
Args:
ingestion_url: URL of the document ingestion service
vector_url: URL of the vector database service
llm_url: URL of the LLM inference service
timeout: Request timeout in seconds
"""
self.ingestion_url = ingestion_url
self.vector_url = vector_url
self.llm_url = llm_url
self.timeout = timeout
logger.info(f"Initialized RAG orchestrator")
logger.info(f"Ingestion service: {ingestion_url}")
logger.info(f"Vector service: {vector_url}")
logger.info(f"LLM service: {llm_url}")
def ingest_document(self, text: str, metadata: Dict = None) -> Dict:
"""
Ingest a document through the ingestion service.
Args:
text: Document text
metadata: Optional metadata
Returns:
Response from ingestion service
"""
try:
response = requests.post(
f"{self.ingestion_url}/ingest",
json={'text': text, 'metadata': metadata or {}},
timeout=self.timeout
)
response.raise_for_status()
return response.json()
except requests.exceptions.RequestException as e:
logger.error(f"Error calling ingestion service: {str(e)}")
raise
def add_to_vector_store(self, documents: List[Dict]) -> Dict:
"""
Add documents to the vector store.
Args:
documents: List of document dictionaries
Returns:
Response from vector service
"""
try:
response = requests.post(
f"{self.vector_url}/add",
json={'documents': documents},
timeout=self.timeout
)
response.raise_for_status()
return response.json()
except requests.exceptions.RequestException as e:
logger.error(f"Error calling vector service: {str(e)}")
raise
def search_documents(self, query: str, top_k: int = 3) -> List[Dict]:
"""
Search for relevant documents.
Args:
query: Search query
top_k: Number of documents to retrieve
Returns:
List of relevant documents with similarity scores
"""
try:
response = requests.post(
f"{self.vector_url}/search",
json={'query': query, 'top_k': top_k},
timeout=self.timeout
)
response.raise_for_status()
data = response.json()
return data.get('results', [])
except requests.exceptions.RequestException as e:
logger.error(f"Error calling vector service: {str(e)}")
raise
def generate_response(
self,
prompt: str,
max_length: int = 512,
temperature: float = 0.7
) -> str:
"""
Generate a response using the LLM service.
Args:
prompt: Input prompt
max_length: Maximum generation length
temperature: Sampling temperature
Returns:
Generated text
"""
try:
response = requests.post(
f"{self.llm_url}/generate",
json={
'prompt': prompt,
'max_length': max_length,
'temperature': temperature
},
timeout=self.timeout * 2 # LLM needs more time
)
response.raise_for_status()
data = response.json()
return data.get('generated_text', '')
except requests.exceptions.RequestException as e:
logger.error(f"Error calling LLM service: {str(e)}")
raise
def construct_rag_prompt(self, query: str, context_docs: List[Dict]) -> str:
"""
Construct a prompt for the LLM with retrieved context.
Args:
query: User query
context_docs: Retrieved context documents
Returns:
Formatted prompt string
"""
# Extract text from context documents
context_texts = [
doc['document']['text']
for doc in context_docs
if 'document' in doc and 'text' in doc['document']
]
# Combine context
context = "\n\n".join(context_texts)
# Construct prompt using a template
prompt = f"""You are a helpful assistant. Use the following context to answer the question. If you cannot answer based on the context, say so.
Context: {context}
Question: {query}
Answer:"""
return prompt
def answer_query(
self,
query: str,
top_k: int = 3,
max_length: int = 512,
temperature: float = 0.7
) -> Dict:
"""
Complete RAG pipeline: retrieve documents and generate answer.
Args:
query: User query
top_k: Number of documents to retrieve
max_length: Maximum generation length
temperature: Sampling temperature
Returns:
Dictionary with answer and metadata
"""
try:
# Step 1: Retrieve relevant documents
logger.info(f"Searching for documents relevant to: {query}")
context_docs = self.search_documents(query, top_k)
if not context_docs:
logger.warning("No relevant documents found")
return {
'answer': "I don't have enough information to answer this question.",
'sources': [],
'query': query
}
# Step 2: Construct prompt with context
prompt = self.construct_rag_prompt(query, context_docs)
# Step 3: Generate answer
logger.info("Generating answer with LLM")
answer = self.generate_response(prompt, max_length, temperature)
# Step 4: Format response
sources = [
{
'text': doc['document']['text'][:200] + '...',
'similarity': doc['similarity'],
'metadata': doc['document'].get('metadata', {})
}
for doc in context_docs
]
return {
'answer': answer,
'sources': sources,
'query': query
}
except Exception as e:
logger.error(f"Error in RAG pipeline: {str(e)}", exc_info=True)
raise
# Initialize orchestrator
orchestrator = RAGOrchestrator(
INGESTION_SERVICE,
VECTOR_SERVICE,
LLM_SERVICE,
REQUEST_TIMEOUT
)
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint."""
# Check connectivity to all services
services_status = {}
for service_name, service_url in [
('ingestion', INGESTION_SERVICE),
('vector', VECTOR_SERVICE),
('llm', LLM_SERVICE)
]:
try:
response = requests.get(
f"{service_url}/health",
timeout=5
)
services_status[service_name] = 'healthy' if response.ok else 'unhealthy'
except Exception:
services_status[service_name] = 'unreachable'
overall_status = 'healthy' if all(
status == 'healthy' for status in services_status.values()
) else 'degraded'
return jsonify({
'status': overall_status,
'service': 'api-gateway',
'version': '1.0.0',
'services': services_status
}), 200
@app.route('/ingest', methods=['POST'])
def ingest():
"""
Ingest a document into the system.
Chunks the document and adds it to the vector store.
"""
try:
data = request.get_json()
if not data or 'text' not in data:
return jsonify({'error': 'No text provided'}), 400
text = data['text']
metadata = data.get('metadata', {})
# Step 1: Chunk the document
logger.info("Ingesting document")
ingestion_result = orchestrator.ingest_document(text, metadata)
chunks = ingestion_result.get('chunks', [])
if not chunks:
return jsonify({'error': 'Failed to chunk document'}), 500
# Step 2: Add chunks to vector store
logger.info(f"Adding {len(chunks)} chunks to vector store")
vector_result = orchestrator.add_to_vector_store(chunks)
return jsonify({
'success': True,
'chunks_created': len(chunks),
'chunks_stored': vector_result.get('added', 0)
}), 200
except Exception as e:
logger.error(f"Error in ingest: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/query', methods=['POST'])
def query():
"""
Answer a query using the RAG pipeline.
Retrieves relevant documents and generates an answer.
"""
try:
data = request.get_json()
if not data or 'query' not in data:
return jsonify({'error': 'No query provided'}), 400
user_query = data['query']
top_k = data.get('top_k', TOP_K_DOCUMENTS)
max_length = data.get('max_length', 512)
temperature = data.get('temperature', 0.7)
# Execute RAG pipeline
logger.info(f"Processing query: {user_query}")
result = orchestrator.answer_query(
user_query,
top_k,
max_length,
temperature
)
return jsonify({
'success': True,
**result
}), 200
except Exception as e:
logger.error(f"Error in query: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5000))
app.run(host='0.0.0.0', port=port, debug=False)
The Dockerfile for the API gateway:
# Dockerfile for API Gateway Service
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Set ownership
RUN chown -R appuser:appuser /app
USER appuser
EXPOSE 5000
HEALTHCHECK --interval=30s --timeout=3s --start-period=10s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5000/health')" || exit 1
CMD ["python", "app.py"]
Build this service:
cd api-gateway
docker build -t rag-chatbot/api-gateway:1.0 .
Now we have all four microservices containerized! Each service is independently deployable, scalable, and maintainable. In the next section, we will learn about Kubernetes and how to orchestrate these services.
PART FOUR: KUBERNETES BASICS - ORCHESTRATION CONCEPTS
You have successfully containerized all four microservices for our RAG chatbot. Now we need to deploy and manage them in production. This is where Kubernetes becomes essential. Running containers manually with Docker commands works for development, but it does not scale. You need automatic recovery from failures, load balancing, service discovery, and the ability to update services without downtime. Kubernetes provides all of this and more.
UNDERSTANDING KUBERNETES ARCHITECTURE
Kubernetes follows a master-worker architecture. The master node (called the control plane) manages the cluster, while worker nodes run your applications. You interact with the control plane through the Kubernetes API using a command-line tool called kubectl.
The control plane consists of several components. The API server is the front end of Kubernetes, handling all API requests. The scheduler assigns Pods to nodes based on resource requirements and constraints. The controller manager runs controllers that maintain the desired state of the cluster. The etcd is a distributed key-value store that holds all cluster data.
Worker nodes run your containerized applications. Each node has a kubelet, which communicates with the control plane and manages containers on that node. The kube-proxy handles network routing for services. A container runtime (like Docker or containerd) actually runs the containers.
KUBERNETES RESOURCES: PODS, DEPLOYMENTS, AND SERVICES
Let us understand the key Kubernetes resources by examining how they work together. A Pod is the smallest deployable unit in Kubernetes. It represents one or more containers that share storage and network resources. Typically, a Pod contains a single container, but you might group multiple tightly coupled containers in one Pod.
Here is a simple Pod definition for our document ingestion service:
apiVersion: v1
kind: Pod
metadata:
name: document-ingestion-pod
labels:
app: document-ingestion
spec:
containers:
- name: document-ingestion
image: rag-chatbot/document-ingestion:1.0
ports:
- containerPort: 5001
env:
- name: CHUNK_SIZE
value: "512"
- name: CHUNK_OVERLAP
value: "50"
This YAML file describes a Pod. The apiVersion field specifies which version of the Kubernetes API we are using. The kind field indicates that this is a Pod resource. The metadata section contains the Pod's name and labels. Labels are key-value pairs used to organize and select resources.
The spec section defines the desired state. We specify one container with the name document-ingestion, using our Docker image. We expose port 5001 and set environment variables for configuration.
However, you should almost never create Pods directly. If a Pod crashes or the node it runs on fails, Kubernetes will not automatically recreate it. Instead, you use higher-level resources like Deployments that manage Pods for you.
A Deployment manages a set of identical Pods and ensures the desired number of replicas are always running. If a Pod crashes, the Deployment automatically creates a new one. If you update the Deployment, it performs a rolling update, gradually replacing old Pods with new ones without downtime.
Here is a Deployment for our document ingestion service:
apiVersion: apps/v1
kind: Deployment
metadata:
name: document-ingestion
labels:
app: document-ingestion
spec:
replicas: 3
selector:
matchLabels:
app: document-ingestion
template:
metadata:
labels:
app: document-ingestion
spec:
containers:
- name: document-ingestion
image: rag-chatbot/document-ingestion:1.0
ports:
- containerPort: 5001
env:
- name: CHUNK_SIZE
value: "512"
- name: CHUNK_OVERLAP
value: "50"
resources:
requests:
memory: "256Mi"
cpu: "250m"
limits:
memory: "512Mi"
cpu: "500m"
livenessProbe:
httpGet:
path: /health
port: 5001
initialDelaySeconds: 15
periodSeconds: 20
readinessProbe:
httpGet:
path: /health
port: 5001
initialDelaySeconds: 5
periodSeconds: 10
This Deployment creates three replicas of our document ingestion service. The replicas field specifies how many Pods we want running. The selector field tells the Deployment which Pods it manages, using label matching. The template section is a Pod template that defines what each Pod should look like.
Notice the resources section. The requests specify the minimum resources guaranteed to the container. The limits specify the maximum resources it can use. This helps Kubernetes schedule Pods efficiently and prevents one Pod from consuming all node resources.
The livenessProbe checks if the container is still running. If the probe fails, Kubernetes restarts the container. The readinessProbe checks if the container is ready to serve traffic. If the probe fails, Kubernetes stops sending traffic to that Pod until it becomes ready again. Both probes use our /health endpoint.
Now we have Pods running, but how do other Pods or external clients access them? Pods are ephemeral; they can be created and destroyed at any time, and each Pod gets a different IP address. We need a stable way to access a set of Pods. This is what Services provide.
A Service is an abstraction that defines a logical set of Pods and a policy for accessing them. The Service has a stable IP address and DNS name, even as the underlying Pods change.
Here is a Service for our document ingestion Deployment:
apiVersion: v1
kind: Service
metadata:
name: document-ingestion-service
spec:
selector:
app: document-ingestion
ports:
- protocol: TCP
port: 80
targetPort: 5001
type: ClusterIP
This Service selects all Pods with the label app: document-ingestion. It exposes port 80 and forwards traffic to port 5001 on the selected Pods. The type ClusterIP means this Service is only accessible within the cluster. Other Pods can access it using the DNS name document-ingestion-service.
There are other Service types. NodePort exposes the Service on each node's IP at a static port, making it accessible from outside the cluster. LoadBalancer creates an external load balancer (in cloud environments) that routes traffic to the Service. We will use these types later for external access.
CONFIGMAPS AND SECRETS FOR CONFIGURATION
Hard-coding configuration in Deployment files is inflexible. If you want to change a configuration value, you need to update the Deployment and redeploy. Kubernetes provides ConfigMaps and Secrets for managing configuration data separately from application code.
A ConfigMap stores non-sensitive configuration data as key-value pairs. Here is a ConfigMap for our document ingestion service:
apiVersion: v1
kind: ConfigMap
metadata:
name: document-ingestion-config
data:
CHUNK_SIZE: "512"
CHUNK_OVERLAP: "50"
MAX_FILE_SIZE: "10485760"
You can reference this ConfigMap in your Deployment:
apiVersion: apps/v1
kind: Deployment
metadata:
name: document-ingestion
spec:
replicas: 3
selector:
matchLabels:
app: document-ingestion
template:
metadata:
labels:
app: document-ingestion
spec:
containers:
- name: document-ingestion
image: rag-chatbot/document-ingestion:1.0
ports:
- containerPort: 5001
envFrom:
- configMapRef:
name: document-ingestion-config
The envFrom field loads all key-value pairs from the ConfigMap as environment variables. Now you can update configuration by changing the ConfigMap without modifying the Deployment.
For sensitive data like API keys, passwords, or certificates, use Secrets instead of ConfigMaps. Secrets work similarly to ConfigMaps but are designed for confidential data. Kubernetes stores Secrets in base64 encoding and provides additional security features.
Here is a Secret for storing an API key:
apiVersion: v1
kind: Secret
metadata:
name: llm-api-secret
type: Opaque
data:
API_KEY: YXBpLWtleS12YWx1ZS1oZXJl
The data is base64-encoded. To create this Secret, you can use kubectl:
kubectl create secret generic llm-api-secret --from-literal=API_KEY=your-api-key-here
Reference the Secret in your Deployment:
env:
- name: API_KEY
valueFrom:
secretKeyRef:
name: llm-api-secret
key: API_KEY
This approach keeps sensitive data out of your Deployment files and container images, improving security.
NAMESPACES FOR ORGANIZING RESOURCES
As your cluster grows, you will have many resources. Namespaces provide a way to organize and isolate resources. Think of namespaces as virtual clusters within your physical cluster. They are useful for separating different environments (development, staging, production) or different teams.
Kubernetes creates several default namespaces. The default namespace is where resources are created if you do not specify a namespace. The kube-system namespace contains Kubernetes system components. The kube-public namespace is readable by all users and is typically used for cluster information.
You can create your own namespaces. Here is a namespace for our RAG chatbot application:
apiVersion: v1
kind: Namespace
metadata:
name: rag-chatbot
Create resources in this namespace by adding the namespace field to the metadata:
apiVersion: apps/v1
kind: Deployment
metadata:
name: document-ingestion
namespace: rag-chatbot
Or specify the namespace when using kubectl:
kubectl apply -f deployment.yaml --namespace=rag-chatbot
Namespaces provide isolation for resources and allow you to apply resource quotas and access controls per namespace.
PART FIVE: DEPLOYING TO KUBERNETES - PRACTICAL IMPLEMENTATION
Now that you understand Kubernetes concepts, let us deploy our RAG chatbot to a Kubernetes cluster. We will create all the necessary Kubernetes resources and see how they work together to create a production-ready application.
SETTING UP A LOCAL KUBERNETES CLUSTER
For learning and development, you can run Kubernetes locally using Minikube or Kind (Kubernetes in Docker). These tools create a single-node cluster on your machine. For this tutorial, we will use Minikube because it is beginner-friendly and supports GPU passthrough for our LLM service.
Install Minikube following the official documentation for your operating system. Once installed, start a cluster:
minikube start --cpus=4 --memory=8192 --driver=docker
This creates a cluster with 4 CPUs and 8GB of memory. The driver flag specifies that Minikube should use Docker as the container runtime.
Verify the cluster is running:
kubectl cluster-info
kubectl get nodes
You should see one node in the Ready state. Minikube also configures kubectl to communicate with this cluster automatically.
CREATING THE NAMESPACE AND CONFIGMAPS
Let us start by creating a namespace for our application. Create a file named namespace.yaml:
apiVersion: v1
kind: Namespace
metadata:
name: rag-chatbot
labels:
name: rag-chatbot
environment: development
Apply this to the cluster:
kubectl apply -f namespace.yaml
Verify the namespace was created:
kubectl get namespaces
Now create ConfigMaps for each service. Create a file named configmaps.yaml:
apiVersion: v1
kind: ConfigMap
metadata:
name: document-ingestion-config
namespace: rag-chatbot
data:
CHUNK_SIZE: "512"
CHUNK_OVERLAP: "50"
MAX_FILE_SIZE: "10485760"
PORT: "5001"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: vector-database-config
namespace: rag-chatbot
data:
EMBEDDING_MODEL: "all-MiniLM-L6-v2"
TOP_K: "5"
PORT: "5002"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: llm-inference-config
namespace: rag-chatbot
data:
MODEL_NAME: "TinyLlama/TinyLlama-1.1B-Chat-v1.0"
MAX_LENGTH: "512"
TEMPERATURE: "0.7"
USE_QUANTIZATION: "true"
PORT: "5003"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: api-gateway-config
namespace: rag-chatbot
data:
TOP_K_DOCUMENTS: "3"
REQUEST_TIMEOUT: "30"
PORT: "5000"
INGESTION_SERVICE: "http://document-ingestion-service:80"
VECTOR_SERVICE: "http://vector-database-service:80"
LLM_SERVICE: "http://llm-inference-service:80"
Notice the three dashes separating multiple resources in one file. This is a YAML convention that allows you to define multiple resources in a single file. Apply these ConfigMaps:
kubectl apply -f configmaps.yaml
Verify they were created:
kubectl get configmaps -n rag-chatbot
DEPLOYING THE DOCUMENT INGESTION SERVICE
Now let us deploy our first microservice. Create a file named document-ingestion-deployment.yaml:
apiVersion: apps/v1
kind: Deployment
metadata:
name: document-ingestion
namespace: rag-chatbot
labels:
app: document-ingestion
component: backend
spec:
replicas: 2
selector:
matchLabels:
app: document-ingestion
template:
metadata:
labels:
app: document-ingestion
component: backend
spec:
containers:
- name: document-ingestion
image: rag-chatbot/document-ingestion:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5001
name: http
protocol: TCP
envFrom:
- configMapRef:
name: document-ingestion-config
resources:
requests:
memory: "256Mi"
cpu: "250m"
limits:
memory: "512Mi"
cpu: "500m"
livenessProbe:
httpGet:
path: /health
port: 5001
initialDelaySeconds: 15
periodSeconds: 20
timeoutSeconds: 3
failureThreshold: 3
readinessProbe:
httpGet:
path: /health
port: 5001
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 3
failureThreshold: 3
---
apiVersion: v1
kind: Service
metadata:
name: document-ingestion-service
namespace: rag-chatbot
labels:
app: document-ingestion
spec:
selector:
app: document-ingestion
ports:
- protocol: TCP
port: 80
targetPort: 5001
name: http
type: ClusterIP
This file defines both the Deployment and the Service for the document ingestion microservice. The imagePullPolicy: IfNotPresent tells Kubernetes to use the local image if available, which is important for our locally built images. In production, you would push images to a container registry and use imagePullPolicy: Always.
Before applying this, we need to make our locally built Docker images available to Minikube. Minikube runs in its own Docker environment, so it cannot see images built on your host. Use this command to load images into Minikube:
minikube image load rag-chatbot/document-ingestion:1.0
Now apply the Deployment:
kubectl apply -f document-ingestion-deployment.yaml
Watch the Pods being created:
kubectl get pods -n rag-chatbot -w
The -w flag watches for changes. You should see two Pods (because we specified replicas: 2) transitioning from Pending to ContainerCreating to Running. Press Ctrl+C to stop watching.
Verify the Service was created:
kubectl get services -n rag-chatbot
You should see the document-ingestion-service with a ClusterIP address.
DEPLOYING THE VECTOR DATABASE SERVICE
Create a file named vector-database-deployment.yaml:
apiVersion: apps/v1
kind: Deployment
metadata:
name: vector-database
namespace: rag-chatbot
labels:
app: vector-database
component: backend
spec:
replicas: 1
selector:
matchLabels:
app: vector-database
template:
metadata:
labels:
app: vector-database
component: backend
spec:
containers:
- name: vector-database
image: rag-chatbot/vector-database:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5002
name: http
protocol: TCP
envFrom:
- configMapRef:
name: vector-database-config
resources:
requests:
memory: "1Gi"
cpu: "500m"
limits:
memory: "2Gi"
cpu: "1000m"
livenessProbe:
httpGet:
path: /health
port: 5002
initialDelaySeconds: 40
periodSeconds: 30
timeoutSeconds: 5
failureThreshold: 3
readinessProbe:
httpGet:
path: /health
port: 5002
initialDelaySeconds: 30
periodSeconds: 20
timeoutSeconds: 5
failureThreshold: 3
---
apiVersion: v1
kind: Service
metadata:
name: vector-database-service
namespace: rag-chatbot
labels:
app: vector-database
spec:
selector:
app: vector-database
ports:
- protocol: TCP
port: 80
targetPort: 5002
name: http
type: ClusterIP
Notice we use only one replica for the vector database. In a production system, you would use a distributed vector database that can be scaled horizontally, but our simple in-memory implementation does not support multiple replicas sharing state.
Also notice the higher resource requests and limits. The vector database needs more memory to store embeddings and perform similarity searches. The health check delays are longer because loading the embedding model takes time.
Load the image and deploy:
minikube image load rag-chatbot/vector-database:1.0
kubectl apply -f vector-database-deployment.yaml
Monitor the deployment:
kubectl get pods -n rag-chatbot -l app=vector-database -w
DEPLOYING THE LLM INFERENCE SERVICE WITH GPU SUPPORT
The LLM inference service is the most resource-intensive component. For optimal performance, it should run on GPU-enabled nodes. However, setting up GPU support in Kubernetes requires additional configuration.
For NVIDIA GPUs, you need to install the NVIDIA GPU Operator in your cluster. For Minikube, you can enable GPU passthrough if your host has an NVIDIA GPU:
minikube start --driver=kvm2 --kvm-gpu
For this tutorial, we will deploy the CPU version of the LLM service, but the configuration is similar for GPU versions. Create a file named llm-inference-deployment.yaml:
apiVersion: apps/v1
kind: Deployment
metadata:
name: llm-inference
namespace: rag-chatbot
labels:
app: llm-inference
component: backend
spec:
replicas: 1
selector:
matchLabels:
app: llm-inference
template:
metadata:
labels:
app: llm-inference
component: backend
spec:
containers:
- name: llm-inference
image: rag-chatbot/llm-inference:1.0-cpu
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5003
name: http
protocol: TCP
envFrom:
- configMapRef:
name: llm-inference-config
resources:
requests:
memory: "4Gi"
cpu: "2000m"
limits:
memory: "8Gi"
cpu: "4000m"
livenessProbe:
httpGet:
path: /health
port: 5003
initialDelaySeconds: 120
periodSeconds: 60
timeoutSeconds: 10
failureThreshold: 3
readinessProbe:
httpGet:
path: /health
port: 5003
initialDelaySeconds: 100
periodSeconds: 30
timeoutSeconds: 10
failureThreshold: 3
---
apiVersion: v1
kind: Service
metadata:
name: llm-inference-service
namespace: rag-chatbot
labels:
app: llm-inference
spec:
selector:
app: llm-inference
ports:
- protocol: TCP
port: 80
targetPort: 5003
name: http
type: ClusterIP
For GPU-enabled deployments, you would add resource requests for GPUs. Here is an example for NVIDIA GPUs:
resources:
requests:
memory: "8Gi"
cpu: "4000m"
nvidia.com/gpu: 1
limits:
memory: "16Gi"
cpu: "8000m"
nvidia.com/gpu: 1
The nvidia.com/gpu resource type is provided by the NVIDIA GPU Operator. For AMD ROCm GPUs, you would use amd.com/gpu. For Apple Silicon, GPU access is handled differently through the Metal framework.
Load the image and deploy:
minikube image load rag-chatbot/llm-inference:1.0-cpu
kubectl apply -f llm-inference-deployment.yaml
This deployment will take several minutes because the container needs to download the LLM model on first startup. Monitor the logs to see progress:
kubectl logs -n rag-chatbot -l app=llm-inference -f
The -f flag follows the logs in real-time. You will see messages about downloading and loading the model. Once you see "Model loaded successfully", the service is ready.
DEPLOYING THE API GATEWAY SERVICE
Finally, let us deploy the API gateway that orchestrates all the other services. Create a file named api-gateway-deployment.yaml:
apiVersion: apps/v1
kind: Deployment
metadata:
name: api-gateway
namespace: rag-chatbot
labels:
app: api-gateway
component: frontend
spec:
replicas: 2
selector:
matchLabels:
app: api-gateway
template:
metadata:
labels:
app: api-gateway
component: frontend
spec:
containers:
- name: api-gateway
image: rag-chatbot/api-gateway:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5000
name: http
protocol: TCP
envFrom:
- configMapRef:
name: api-gateway-config
resources:
requests:
memory: "128Mi"
cpu: "100m"
limits:
memory: "256Mi"
cpu: "200m"
livenessProbe:
httpGet:
path: /health
port: 5000
initialDelaySeconds: 10
periodSeconds: 20
timeoutSeconds: 5
failureThreshold: 3
readinessProbe:
httpGet:
path: /health
port: 5000
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 5
failureThreshold: 3
---
apiVersion: v1
kind: Service
metadata:
name: api-gateway-service
namespace: rag-chatbot
labels:
app: api-gateway
spec:
selector:
app: api-gateway
ports:
- protocol: TCP
port: 80
targetPort: 5000
name: http
type: ClusterIP
Load the image and deploy:
minikube image load rag-chatbot/api-gateway:1.0
kubectl apply -f api-gateway-deployment.yaml
Verify all services are running:
kubectl get pods -n rag-chatbot
You should see all Pods in the Running state with health checks passing.
EXPOSING THE APPLICATION WITH INGRESS
Now all our services are running inside the cluster, but we cannot access them from outside. We need to expose the API gateway to external clients. While we could use a LoadBalancer Service, the recommended approach for HTTP/HTTPS traffic is to use an Ingress resource.
An Ingress manages external access to services in a cluster, typically HTTP and HTTPS. It provides load balancing, SSL termination, and name-based virtual hosting. To use Ingress, you need an Ingress Controller running in your cluster.
Enable the NGINX Ingress Controller in Minikube:
minikube addons enable ingress
Verify the Ingress Controller is running:
kubectl get pods -n ingress-nginx
Now create an Ingress resource. Create a file named ingress.yaml:
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: rag-chatbot-ingress
namespace: rag-chatbot
annotations:
nginx.ingress.kubernetes.io/rewrite-target: /
nginx.ingress.kubernetes.io/proxy-body-size: "10m"
nginx.ingress.kubernetes.io/proxy-read-timeout: "300"
nginx.ingress.kubernetes.io/proxy-send-timeout: "300"
spec:
ingressClassName: nginx
rules:
- host: rag-chatbot.local
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: api-gateway-service
port:
number: 80
This Ingress routes all HTTP traffic for rag-chatbot.local to our API gateway service. The annotations configure NGINX-specific settings like maximum request body size and timeout values.
Apply the Ingress:
kubectl apply -f ingress.yaml
Get the Ingress IP address:
kubectl get ingress -n rag-chatbot
For Minikube, you need to run minikube tunnel in a separate terminal to make the Ingress accessible:
minikube tunnel
This command creates a network route from your host to the Minikube cluster. Keep this running.
Add an entry to your /etc/hosts file (or C:\Windows\System32\drivers\etc\hosts on Windows) to resolve rag-chatbot.local:
127.0.0.1 rag-chatbot.local
Now you can access your RAG chatbot from your browser or command line! Test the health endpoint:
curl http://rag-chatbot.local/health
You should receive a JSON response showing the status of all services.
TESTING THE COMPLETE RAG PIPELINE
Let us test the complete RAG pipeline by ingesting a document and then querying it. First, ingest some documentation about Docker:
curl -X POST http://rag-chatbot.local/ingest \
-H "Content-Type: application/json" \
-d '{
"text": "Docker is a platform for developing, shipping, and running applications in containers. Containers are lightweight, standalone, executable packages that include everything needed to run a piece of software, including the code, runtime, system tools, libraries, and settings. Docker containers are isolated from each other and the host system, but they share the operating system kernel, making them more efficient than virtual machines. Docker uses images as templates for creating containers. An image is built from a Dockerfile, which contains instructions for assembling the image. Docker Hub is a registry where you can find and share container images.",
"metadata": {"source": "docker-intro", "topic": "containers"}
}'
You should receive a response indicating the document was chunked and stored successfully. Now query the system:
curl -X POST http://rag-chatbot.local/query \
-H "Content-Type: application/json" \
-d '{
"query": "What is Docker and how does it work?",
"top_k": 3,
"max_length": 200
}'
The system will retrieve relevant document chunks, construct a prompt with context, generate an answer using the LLM, and return the response along with source information. This demonstrates the complete RAG pipeline working across all four microservices orchestrated by Kubernetes!
PART SIX: ADVANCED CONCEPTS - SCALING, MONITORING, AND OPTIMIZATION
You now have a working Kubernetes application with all four microservices deployed and communicating. Let us explore advanced concepts that make your application production-ready.
HORIZONTAL POD AUTOSCALING
One of Kubernetes' most powerful features is automatic scaling based on resource utilization or custom metrics. The Horizontal Pod Autoscaler automatically scales the number of Pods in a Deployment based on observed metrics.
First, ensure the Metrics Server is installed in your cluster. For Minikube:
minikube addons enable metrics-server
Create a HorizontalPodAutoscaler for the API gateway. Create a file named hpa.yaml:
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: api-gateway-hpa
namespace: rag-chatbot
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: api-gateway
minReplicas: 2
maxReplicas: 10
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 80
behavior:
scaleDown:
stabilizationWindowSeconds: 300
policies:
- type: Percent
value: 50
periodSeconds: 60
scaleUp:
stabilizationWindowSeconds: 0
policies:
- type: Percent
value: 100
periodSeconds: 30
- type: Pods
value: 2
periodSeconds: 30
selectPolicy: Max
This HPA monitors CPU and memory utilization. When average CPU usage exceeds seventy percent or memory usage exceeds eighty percent, it scales up the number of Pods. It maintains between two and ten replicas. The behavior section controls how quickly scaling happens, preventing rapid oscillation.
Apply the HPA:
kubectl apply -f hpa.yaml
Monitor the HPA:
kubectl get hpa -n rag-chatbot -w
You will see current and target metrics, along with the current number of replicas. To test autoscaling, you could generate load using a tool like Apache Bench or hey.
RESOURCE QUOTAS AND LIMIT RANGES
In a shared cluster, you want to prevent one namespace from consuming all resources. ResourceQuota limits the total resources that can be consumed in a namespace. Create a file named resource-quota.yaml:
apiVersion: v1
kind: ResourceQuota
metadata:
name: rag-chatbot-quota
namespace: rag-chatbot
spec:
hard:
requests.cpu: "10"
requests.memory: "20Gi"
limits.cpu: "20"
limits.memory: "40Gi"
pods: "50"
services: "10"
persistentvolumeclaims: "5"
This quota limits the namespace to ten CPU cores of requests, twenty CPU cores of limits, twenty gigabytes of memory requests, forty gigabytes of memory limits, fifty Pods, ten Services, and five PersistentVolumeClaims.
A LimitRange sets default resource requests and limits for containers that do not specify them. Create a file named limit-range.yaml:
apiVersion: v1
kind: LimitRange
metadata:
name: rag-chatbot-limits
namespace: rag-chatbot
spec:
limits:
- max:
cpu: "4"
memory: "8Gi"
min:
cpu: "50m"
memory: "64Mi"
default:
cpu: "500m"
memory: "512Mi"
defaultRequest:
cpu: "100m"
memory: "128Mi"
type: Container
Apply both:
kubectl apply -f resource-quota.yaml
kubectl apply -f limit-range.yaml
Now any container created without resource specifications will get the default values, and no container can exceed the maximum values.
PERSISTENT STORAGE FOR STATEFUL SERVICES
Our vector database currently stores embeddings in memory, which means they are lost when the Pod restarts. For production, you need persistent storage. Kubernetes provides PersistentVolumes and PersistentVolumeClaims for this purpose.
A PersistentVolume is a piece of storage in the cluster that has been provisioned by an administrator or dynamically provisioned using Storage Classes. A PersistentVolumeClaim is a request for storage by a user.
Create a PersistentVolumeClaim for the vector database. Create a file named pvc.yaml:
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: vector-database-pvc
namespace: rag-chatbot
spec:
accessModes:
- ReadWriteOnce
resources:
requests:
storage: 10Gi
storageClassName: standard
Apply this:
kubectl apply -f pvc.yaml
Now modify the vector database Deployment to use this volume. Update vector-database-deployment.yaml:
spec:
containers:
- name: vector-database
image: rag-chatbot/vector-database:1.0
volumeMounts:
- name: data
mountPath: /app/data
volumes:
- name: data
persistentVolumeClaim:
claimName: vector-database-pvc
The volumeMounts section mounts the volume at /app/data inside the container. You would need to modify the application code to save and load embeddings from this directory.
IMPLEMENTING LOGGING AND MONITORING
Production applications need comprehensive logging and monitoring. Kubernetes provides several options. The most common approach is the ELK stack (Elasticsearch, Logstash, Kibana) or the more modern Loki stack (Loki, Promtail, Grafana).
For metrics, Prometheus is the standard choice. It scrapes metrics from your applications and stores them in a time-series database. Grafana provides visualization and dashboarding.
To expose metrics from your Python applications, use the prometheus-client library. Add to requirements.txt:
prometheus-client==0.19.0
Modify your Flask applications to expose metrics. Add to app.py:
from prometheus_client import Counter, Histogram, generate_latest, CONTENT_TYPE_LATEST
import time
# Define metrics
REQUEST_COUNT = Counter(
'http_requests_total',
'Total HTTP requests',
['method', 'endpoint', 'status']
)
REQUEST_DURATION = Histogram(
'http_request_duration_seconds',
'HTTP request duration',
['method', 'endpoint']
)
# Middleware to track metrics
@app.before_request
def before_request():
request.start_time = time.time()
@app.after_request
def after_request(response):
duration = time.time() - request.start_time
REQUEST_DURATION.labels(
method=request.method,
endpoint=request.endpoint or 'unknown'
).observe(duration)
REQUEST_COUNT.labels(
method=request.method,
endpoint=request.endpoint or 'unknown',
status=response.status_code
).inc()
return response
# Metrics endpoint
@app.route('/metrics')
def metrics():
return generate_latest(), 200, {'Content-Type': CONTENT_TYPE_LATEST}
Now your services expose metrics at the /metrics endpoint in Prometheus format. You can install Prometheus and Grafana in your cluster using Helm:
helm repo add prometheus-community https://prometheus-community.github.io/helm-charts
helm repo update
helm install prometheus prometheus-community/kube-prometheus-stack -n monitoring --create-namespace
This installs Prometheus, Grafana, and Alertmanager with sensible defaults. Access Grafana:
kubectl port-forward -n monitoring svc/prometheus-grafana 3000:80
Open http://localhost:3000 in your browser. The default credentials are admin/prom-operator. You can create dashboards to visualize your application metrics.
IMPLEMENTING HEALTH CHECKS AND GRACEFUL SHUTDOWN
We already implemented liveness and readiness probes, but for production robustness, you should also implement graceful shutdown. When Kubernetes terminates a Pod, it sends a SIGTERM signal to the container. Your application should catch this signal, stop accepting new requests, finish processing existing requests, and then exit.
Modify your Flask applications to handle graceful shutdown:
import signal
import sys
class GracefulShutdown:
def __init__(self):
self.shutdown_requested = False
signal.signal(signal.SIGTERM, self.handle_sigterm)
signal.signal(signal.SIGINT, self.handle_sigterm)
def handle_sigterm(self, signum, frame):
logger.info("Shutdown signal received, starting graceful shutdown")
self.shutdown_requested = True
def is_shutting_down(self):
return self.shutdown_requested
shutdown_handler = GracefulShutdown()
@app.before_request
def check_shutdown():
if shutdown_handler.is_shutting_down():
return jsonify({'error': 'Service is shutting down'}), 503
Additionally, configure a preStop hook in your Deployment to give the application time to finish processing:
spec:
containers:
- name: api-gateway
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 15"]
This ensures Kubernetes waits fifteen seconds after sending SIGTERM before forcefully killing the container.
ADVANCED GPU OPTIMIZATION TECHNIQUES
For the LLM inference service, optimizing GPU utilization is critical for performance and cost efficiency. Several advanced techniques can significantly improve throughput.
Model quantization reduces the precision of model weights from 32-bit or 16-bit floating point to 8-bit or even 4-bit integers. This reduces memory usage and increases inference speed with minimal accuracy loss. We already implemented 4-bit quantization using the bitsandbytes library, but you can go further with techniques like GPTQ or AWQ quantization.
Batching multiple requests together improves GPU utilization. Modify the LLM inference service to support batched inference:
def generate_batch(
self,
prompts: List[str],
max_length: int = 512,
temperature: float = 0.7
) -> List[str]:
"""Generate responses for multiple prompts in a batch."""
inputs = self.tokenizer(
prompts,
return_tensors="pt",
padding=True,
truncation=True,
max_length=max_length
)
inputs = {k: v.to(self.device) for k, v in inputs.items()}
with torch.no_grad():
outputs = self.model.generate(
**inputs,
max_new_tokens=max_length,
temperature=temperature,
do_sample=True,
pad_token_id=self.tokenizer.pad_token_id
)
results = []
for i, output in enumerate(outputs):
text = self.tokenizer.decode(output, skip_special_tokens=True)
if text.startswith(prompts[i]):
text = text[len(prompts[i]):].strip()
results.append(text)
return results
For NVIDIA GPUs, you can use TensorRT for further optimization. TensorRT is a deep learning inference optimizer and runtime that can provide up to 8x faster inference. Converting your model to TensorRT format requires additional steps but can dramatically improve performance.
For AMD ROCm GPUs, ensure you are using the latest ROCm version and optimized libraries. The vLLM library has excellent ROCm support and provides advanced features like paged attention and continuous batching.
For Apple Silicon, use the MLX framework for optimal performance. MLX is specifically designed for Apple Silicon and leverages the unified memory architecture. You would need to convert your model to MLX format and modify the inference code accordingly.
IMPLEMENTING CIRCUIT BREAKERS AND RETRY LOGIC
In a microservices architecture, services depend on each other. If one service becomes slow or unavailable, it can cascade and affect the entire system. Circuit breakers prevent this by detecting failures and stopping requests to failing services.
Implement a circuit breaker in the API gateway. Add to requirements.txt:
pybreaker==1.0.1
Modify the RAGOrchestrator class:
from pybreaker import CircuitBreaker
class RAGOrchestrator:
def __init__(self, ingestion_url, vector_url, llm_url, timeout):
self.ingestion_url = ingestion_url
self.vector_url = vector_url
self.llm_url = llm_url
self.timeout = timeout
# Create circuit breakers for each service
self.ingestion_breaker = CircuitBreaker(
fail_max=5,
timeout_duration=60,
name='ingestion-service'
)
self.vector_breaker = CircuitBreaker(
fail_max=5,
timeout_duration=60,
name='vector-service'
)
self.llm_breaker = CircuitBreaker(
fail_max=3,
timeout_duration=120,
name='llm-service'
)
def search_documents(self, query, top_k=3):
@self.vector_breaker
def _search():
response = requests.post(
f"{self.vector_url}/search",
json={'query': query, 'top_k': top_k},
timeout=self.timeout
)
response.raise_for_status()
return response.json().get('results', [])
try:
return _search()
except Exception as e:
logger.error(f"Circuit breaker opened for vector service: {str(e)}")
return []
The circuit breaker monitors failures. After five consecutive failures, it opens the circuit and immediately returns an error without calling the service. After sixty seconds, it allows one request through to test if the service has recovered. This prevents overwhelming a struggling service and allows it to recover.
IMPLEMENTING A SERVICE MESH WITH ISTIO
For advanced traffic management, security, and observability, consider implementing a service mesh like Istio. A service mesh provides features like mutual TLS between services, advanced routing, canary deployments, and distributed tracing without modifying application code.
Installing Istio is beyond the scope of this tutorial, but the basic concept is that Istio injects a sidecar proxy into each Pod. This proxy intercepts all network traffic and provides the service mesh features. You can then define traffic policies using Istio resources.
For example, you could implement a canary deployment where ninety percent of traffic goes to the stable version and ten percent goes to the new version:
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
name: api-gateway
namespace: rag-chatbot
spec:
hosts:
- api-gateway-service
http:
- match:
- headers:
canary:
exact: "true"
route:
- destination:
host: api-gateway-service
subset: v2
- route:
- destination:
host: api-gateway-service
subset: v1
weight: 90
- destination:
host: api-gateway-service
subset: v2
weight: 10
This allows you to test new versions with a small percentage of traffic before rolling out to all users.
PART SEVEN: PRODUCTION CONSIDERATIONS - SECURITY AND BEST PRACTICES
Deploying to production requires additional considerations beyond functionality. Security, reliability, and operational excellence are critical.
SECURITY BEST PRACTICES
Security should be built into every layer of your application. Start with container security. Never run containers as root. We already implemented this in our Dockerfiles by creating non-root users. Additionally, scan your images for vulnerabilities using tools like Trivy:
trivy image rag-chatbot/document-ingestion:1.0
Fix any high or critical vulnerabilities by updating dependencies or base images.
Use Pod Security Standards to enforce security policies. Kubernetes provides three levels: Privileged, Baseline, and Restricted. For production, use the Restricted standard. Create a file named pod-security.yaml:
apiVersion: v1
kind: Namespace
metadata:
name: rag-chatbot
labels:
pod-security.kubernetes.io/enforce: restricted
pod-security.kubernetes.io/audit: restricted
pod-security.kubernetes.io/warn: restricted
This enforces the restricted security standard on all Pods in the namespace. You may need to update your Deployments to comply with these restrictions, such as setting securityContext:
spec:
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
seccompProfile:
type: RuntimeDefault
containers:
- name: document-ingestion
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
readOnlyRootFilesystem: true
Implement network policies to control traffic between Pods. By default, all Pods can communicate with all other Pods. Network policies restrict this. Create a file named network-policy.yaml:
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: api-gateway-policy
namespace: rag-chatbot
spec:
podSelector:
matchLabels:
app: api-gateway
policyTypes:
- Ingress
- Egress
ingress:
- from:
- namespaceSelector:
matchLabels:
name: ingress-nginx
ports:
- protocol: TCP
port: 5000
egress:
- to:
- podSelector:
matchLabels:
app: document-ingestion
ports:
- protocol: TCP
port: 5001
- to:
- podSelector:
matchLabels:
app: vector-database
ports:
- protocol: TCP
port: 5002
- to:
- podSelector:
matchLabels:
app: llm-inference
ports:
- protocol: TCP
port: 5003
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 53
This policy allows the API gateway to receive traffic only from the Ingress controller and send traffic only to the other services and DNS. Create similar policies for all services.
For sensitive data like API keys, use Kubernetes Secrets with encryption at rest. Enable encryption at rest in your cluster and use external secret management systems like HashiCorp Vault or AWS Secrets Manager for additional security.
Implement Role-Based Access Control to restrict who can access and modify resources. Create service accounts with minimal permissions for your applications:
apiVersion: v1
kind: ServiceAccount
metadata:
name: api-gateway-sa
namespace: rag-chatbot
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: api-gateway-role
namespace: rag-chatbot
rules:
- apiGroups: [""]
resources: ["configmaps"]
verbs: ["get", "list"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: api-gateway-rolebinding
namespace: rag-chatbot
subjects:
- kind: ServiceAccount
name: api-gateway-sa
namespace: rag-chatbot
roleRef:
kind: Role
name: api-gateway-role
apiGroup: rbac.authorization.k8s.io
Reference this service account in your Deployment:
spec:
serviceAccountName: api-gateway-sa
IMPLEMENTING SSL/TLS WITH CERT-MANAGER
For production, you need HTTPS. Cert-manager automates certificate management. Install cert-manager:
kubectl apply -f https://github.com/cert-manager/cert-manager/releases/download/v1.13.0/cert-manager.yaml
Create a ClusterIssuer for Let's Encrypt:
apiVersion: cert-manager.io/v1
kind: ClusterIssuer
metadata:
name: letsencrypt-prod
spec:
acme:
server: https://acme-v02.api.letsencrypt.org/directory
email: your-email@example.com
privateKeySecretRef:
name: letsencrypt-prod
solvers:
- http01:
ingress:
class: nginx
Update your Ingress to use TLS:
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: rag-chatbot-ingress
namespace: rag-chatbot
annotations:
cert-manager.io/cluster-issuer: letsencrypt-prod
spec:
ingressClassName: nginx
tls:
- hosts:
- rag-chatbot.example.com
secretName: rag-chatbot-tls
rules:
- host: rag-chatbot.example.com
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: api-gateway-service
port:
number: 80
Cert-manager will automatically obtain and renew certificates from Let's Encrypt.
IMPLEMENTING BACKUP AND DISASTER RECOVERY
For stateful services like the vector database, implement regular backups. Use Velero for cluster-level backup and restore:
velero install --provider aws --bucket my-backup-bucket --secret-file ./credentials-velero
Create a backup schedule:
velero schedule create rag-chatbot-daily --schedule="0 2 * * *" --include-namespaces rag-chatbot
This creates daily backups at 2 AM. Test your restore process regularly to ensure backups are working.
IMPLEMENTING CONTINUOUS DEPLOYMENT
For production, implement a CI/CD pipeline that automatically builds, tests, and deploys your application. A typical pipeline includes these stages:
Build stage: Build Docker images and tag them with the commit SHA or version number. Push images to a container registry like Docker Hub, Google Container Registry, or Amazon ECR.
Test stage: Run unit tests, integration tests, and security scans. Only proceed if all tests pass.
Deploy stage: Update Kubernetes manifests with the new image tags and apply them to the cluster. Use tools like kubectl, Helm, or GitOps tools like ArgoCD or Flux.
Here is an example GitHub Actions workflow:
name: CI/CD Pipeline
on:
push:
branches: [main]
jobs:
build-and-deploy:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Build Docker images
run: |
docker build -t myregistry/document-ingestion:${{ github.sha }} ./document-ingestion
docker build -t myregistry/vector-database:${{ github.sha }} ./vector-database
docker build -t myregistry/llm-inference:${{ github.sha }} -f ./llm-inference/Dockerfile.cpu ./llm-inference
docker build -t myregistry/api-gateway:${{ github.sha }} ./api-gateway
- name: Run tests
run: |
pytest tests/
- name: Push images
run: |
echo ${{ secrets.REGISTRY_PASSWORD }} | docker login -u ${{ secrets.REGISTRY_USERNAME }} --password-stdin
docker push myregistry/document-ingestion:${{ github.sha }}
docker push myregistry/vector-database:${{ github.sha }}
docker push myregistry/llm-inference:${{ github.sha }}
docker push myregistry/api-gateway:${{ github.sha }}
- name: Update Kubernetes manifests
run: |
sed -i "s|image: rag-chatbot/document-ingestion:.*|image: myregistry/document-ingestion:${{ github.sha }}|" k8s/document-ingestion-deployment.yaml
sed -i "s|image: rag-chatbot/vector-database:.*|image: myregistry/vector-database:${{ github.sha }}|" k8s/vector-database-deployment.yaml
sed -i "s|image: rag-chatbot/llm-inference:.*|image: myregistry/llm-inference:${{ github.sha }}|" k8s/llm-inference-deployment.yaml
sed -i "s|image: rag-chatbot/api-gateway:.*|image: myregistry/api-gateway:${{ github.sha }}|" k8s/api-gateway-deployment.yaml
- name: Deploy to Kubernetes
run: |
kubectl apply -f k8s/
env:
KUBECONFIG: ${{ secrets.KUBECONFIG }}
This workflow builds images, runs tests, pushes to a registry, updates manifests, and deploys to Kubernetes automatically on every commit to main.
COST OPTIMIZATION STRATEGIES
Running Kubernetes in production can be expensive, especially with GPU workloads. Implement these strategies to optimize costs:
Right-size your resources. Monitor actual resource usage and adjust requests and limits accordingly. Over-provisioning wastes money.
Use cluster autoscaling to automatically add or remove nodes based on demand. Most cloud providers support this.
Use spot instances or preemptible VMs for non-critical workloads. These are significantly cheaper but can be terminated with short notice.
Implement pod disruption budgets to ensure availability during node maintenance:
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
name: api-gateway-pdb
namespace: rag-chatbot
spec:
minAvailable: 1
selector:
matchLabels:
app: api-gateway
For GPU workloads, use time-slicing or multi-instance GPU to share GPUs between multiple workloads when full GPU power is not needed.
Implement request-based autoscaling for the LLM service. Scale to zero when there are no requests and scale up based on queue length.
CONCLUSION: YOUR JOURNEY CONTINUES
Congratulations! You have built a complete, production-ready Kubernetes application from scratch. You started with no knowledge of Docker or Kubernetes and now understand how to containerize applications, orchestrate them with Kubernetes, implement advanced features like autoscaling and monitoring, and follow security best practices.
You built a real-world RAG-based LLM chatbot with four microservices: document ingestion, vector database, LLM inference with GPU optimization, and API gateway. Each service is independently deployable, scalable, and maintainable. You learned how to expose your application to external traffic using Ingress, how to manage configuration with ConfigMaps and Secrets, how to implement health checks and graceful shutdown, and how to monitor and optimize your application.
The concepts you learned apply to any Kubernetes application, not just LLM chatbots. Microservices architecture, containerization, orchestration, and cloud-native practices are fundamental skills for modern software development.
Your journey does not end here. Kubernetes is a vast ecosystem with many more features to explore. Consider learning about StatefulSets for stateful applications, DaemonSets for node-level services, Jobs and CronJobs for batch processing, Custom Resource Definitions for extending Kubernetes, Operators for managing complex applications, and service meshes like Istio or Linkerd for advanced traffic management.
The cloud-native landscape evolves rapidly. Stay current by following the Kubernetes blog, participating in the community, and experimenting with new tools and techniques. Build projects, break things, learn from failures, and share your knowledge with others.
You now have the foundation to build and deploy sophisticated applications that scale to millions of users. The skills you learned are in high demand and will serve you well throughout your career. Keep learning, keep building, and most importantly, keep shipping!
Production-Ready RAG-Based LLM Chatbot with Docker and Kubernetes
Project Structure
rag-chatbot/
├── docker-compose.yml
├── k8s/
│ ├── namespace.yaml
│ ├── configmaps.yaml
│ ├── secrets.yaml
│ ├── document-ingestion/
│ │ ├── deployment.yaml
│ │ └── service.yaml
│ ├── vector-database/
│ │ ├── deployment.yaml
│ │ ├── service.yaml
│ │ └── pvc.yaml
│ ├── llm-inference/
│ │ ├── deployment.yaml
│ │ └── service.yaml
│ ├── api-gateway/
│ │ ├── deployment.yaml
│ │ └── service.yaml
│ ├── ingress.yaml
│ ├── hpa.yaml
│ ├── network-policies.yaml
│ └── monitoring/
│ └── servicemonitor.yaml
├── services/
│ ├── document-ingestion/
│ │ ├── app.py
│ │ ├── requirements.txt
│ │ ├── Dockerfile
│ │ ├── .dockerignore
│ │ └── tests/
│ │ └── test_app.py
│ ├── vector-database/
│ │ ├── app.py
│ │ ├── requirements.txt
│ │ ├── Dockerfile
│ │ ├── .dockerignore
│ │ └── tests/
│ │ └── test_app.py
│ ├── llm-inference/
│ │ ├── app.py
│ │ ├── requirements.txt
│ │ ├── Dockerfile.cpu
│ │ ├── Dockerfile.cuda
│ │ ├── .dockerignore
│ │ └── tests/
│ │ └── test_app.py
│ └── api-gateway/
│ ├── app.py
│ ├── requirements.txt
│ ├── Dockerfile
│ ├── .dockerignore
│ └── tests/
│ └── test_app.py
├── scripts/
│ ├── build-images.sh
│ ├── deploy-local.sh
│ ├── deploy-k8s.sh
│ ├── test-system.sh
│ └── cleanup.sh
├── tests/
│ └── integration/
│ └── test_rag_pipeline.py
├── docs/
│ ├── DEPLOYMENT.md
│ ├── ARCHITECTURE.md
│ └── API.md
├── .github/
│ └── workflows/
│ └── ci-cd.yaml
├── Makefile
└── README.md
1. Document Ingestion Service
services/document-ingestion/app.py
"""
Document Ingestion Service for RAG-based LLM Chatbot
Handles document chunking and preprocessing
"""
from flask import Flask, request, jsonify
from flask_cors import CORS
import os
import logging
import tiktoken
from typing import List, Dict, Optional
from prometheus_client import Counter, Histogram, generate_latest, CONTENT_TYPE_LATEST
import time
import signal
import sys
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
CORS(app)
# Configuration from environment variables
CHUNK_SIZE = int(os.environ.get('CHUNK_SIZE', '512'))
CHUNK_OVERLAP = int(os.environ.get('CHUNK_OVERLAP', '50'))
MAX_FILE_SIZE = int(os.environ.get('MAX_FILE_SIZE', '10485760')) # 10MB
# Prometheus metrics
REQUEST_COUNT = Counter(
'http_requests_total',
'Total HTTP requests',
['method', 'endpoint', 'status']
)
REQUEST_DURATION = Histogram(
'http_request_duration_seconds',
'HTTP request duration',
['method', 'endpoint']
)
CHUNKS_CREATED = Counter(
'chunks_created_total',
'Total number of chunks created'
)
class GracefulShutdown:
"""Handle graceful shutdown on SIGTERM"""
def __init__(self):
self.shutdown_requested = False
signal.signal(signal.SIGTERM, self.handle_sigterm)
signal.signal(signal.SIGINT, self.handle_sigterm)
def handle_sigterm(self, signum, frame):
logger.info("Shutdown signal received, starting graceful shutdown")
self.shutdown_requested = True
def is_shutting_down(self):
return self.shutdown_requested
shutdown_handler = GracefulShutdown()
class DocumentChunker:
"""
Handles the chunking of documents into smaller pieces.
Uses token-based chunking for accurate LLM context management.
"""
def __init__(self, chunk_size: int, chunk_overlap: int):
self.chunk_size = chunk_size
self.chunk_overlap = chunk_overlap
try:
self.tokenizer = tiktoken.get_encoding("cl100k_base")
logger.info(f"Initialized chunker: size={chunk_size}, overlap={chunk_overlap}")
except Exception as e:
logger.error(f"Failed to initialize tokenizer: {e}")
raise
def chunk_text(self, text: str, metadata: Optional[Dict] = None) -> List[Dict]:
"""
Split text into overlapping chunks based on token count.
Args:
text: The input text to chunk
metadata: Optional metadata to attach to each chunk
Returns:
List of dictionaries containing chunk text and metadata
"""
if not text or not text.strip():
logger.warning("Received empty text for chunking")
return []
try:
# Tokenize the entire text
tokens = self.tokenizer.encode(text)
total_tokens = len(tokens)
logger.info(f"Processing text with {total_tokens} tokens")
chunks = []
start_idx = 0
chunk_id = 0
while start_idx < total_tokens:
# Calculate end index for this chunk
end_idx = min(start_idx + self.chunk_size, total_tokens)
# Extract tokens for this chunk
chunk_tokens = tokens[start_idx:end_idx]
# Decode tokens back to text
chunk_text = self.tokenizer.decode(chunk_tokens)
# Create chunk object with metadata
chunk = {
'id': chunk_id,
'text': chunk_text,
'start_token': start_idx,
'end_token': end_idx,
'token_count': len(chunk_tokens),
'metadata': metadata or {}
}
chunks.append(chunk)
chunk_id += 1
CHUNKS_CREATED.inc()
# Move start index forward, accounting for overlap
start_idx = end_idx - self.chunk_overlap
# Prevent infinite loop if overlap is too large
if start_idx >= end_idx:
break
logger.info(f"Created {len(chunks)} chunks from text")
return chunks
except Exception as e:
logger.error(f"Error chunking text: {str(e)}", exc_info=True)
raise
# Initialize the chunker
chunker = DocumentChunker(CHUNK_SIZE, CHUNK_OVERLAP)
@app.before_request
def before_request():
"""Track request start time and check shutdown status"""
if shutdown_handler.is_shutting_down():
return jsonify({'error': 'Service is shutting down'}), 503
request.start_time = time.time()
@app.after_request
def after_request(response):
"""Track request metrics"""
if hasattr(request, 'start_time'):
duration = time.time() - request.start_time
REQUEST_DURATION.labels(
method=request.method,
endpoint=request.endpoint or 'unknown'
).observe(duration)
REQUEST_COUNT.labels(
method=request.method,
endpoint=request.endpoint or 'unknown',
status=response.status_code
).inc()
return response
@app.route('/health', methods=['GET'])
def health_check():
"""
Health check endpoint for Kubernetes liveness and readiness probes.
"""
return jsonify({
'status': 'healthy',
'service': 'document-ingestion',
'version': '1.0.0',
'config': {
'chunk_size': CHUNK_SIZE,
'chunk_overlap': CHUNK_OVERLAP,
'max_file_size': MAX_FILE_SIZE
}
}), 200
@app.route('/ready', methods=['GET'])
def readiness_check():
"""
Readiness check endpoint.
"""
if shutdown_handler.is_shutting_down():
return jsonify({'status': 'not ready', 'reason': 'shutting down'}), 503
return jsonify({'status': 'ready'}), 200
@app.route('/ingest', methods=['POST'])
def ingest_document():
"""
Main endpoint for document ingestion.
Accepts text or file uploads and returns chunked documents.
"""
try:
# Check if request contains JSON data with text
if request.is_json:
data = request.get_json()
text = data.get('text', '')
metadata = data.get('metadata', {})
if not text:
return jsonify({'error': 'No text provided'}), 400
# Process the text into chunks
chunks = chunker.chunk_text(text, metadata)
return jsonify({
'success': True,
'chunks': chunks,
'total_chunks': len(chunks)
}), 200
# Check if request contains file upload
elif 'file' in request.files:
file = request.files['file']
# Validate file size
file.seek(0, os.SEEK_END)
file_size = file.tell()
file.seek(0)
if file_size > MAX_FILE_SIZE:
return jsonify({
'error': f'File too large. Maximum size is {MAX_FILE_SIZE} bytes'
}), 400
# Read file content
try:
text = file.read().decode('utf-8')
except UnicodeDecodeError:
return jsonify({'error': 'File must be UTF-8 encoded text'}), 400
# Extract metadata from form data
metadata = {
'filename': file.filename,
'size': file_size
}
# Process the text into chunks
chunks = chunker.chunk_text(text, metadata)
return jsonify({
'success': True,
'chunks': chunks,
'total_chunks': len(chunks)
}), 200
else:
return jsonify({
'error': 'No text or file provided'
}), 400
except Exception as e:
logger.error(f"Error processing document: {str(e)}", exc_info=True)
return jsonify({
'error': 'Internal server error',
'message': str(e)
}), 500
@app.route('/config', methods=['GET'])
def get_config():
"""
Returns current service configuration.
"""
return jsonify({
'chunk_size': CHUNK_SIZE,
'chunk_overlap': CHUNK_OVERLAP,
'max_file_size': MAX_FILE_SIZE
}), 200
@app.route('/metrics', methods=['GET'])
def metrics():
"""
Prometheus metrics endpoint.
"""
return generate_latest(), 200, {'Content-Type': CONTENT_TYPE_LATEST}
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5001))
logger.info(f"Starting Document Ingestion Service on port {port}")
app.run(host='0.0.0.0', port=port, debug=False)
services/document-ingestion/requirements.txt
flask==3.0.0
flask-cors==4.0.0
werkzeug==3.0.1
tiktoken==0.5.2
prometheus-client==0.19.0
gunicorn==21.2.0
services/document-ingestion/Dockerfile
# Multi-stage build for Document Ingestion Service
FROM python:3.11-slim AS base
# Set environment variables
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
PIP_NO_CACHE_DIR=1 \
PIP_DISABLE_PIP_VERSION_CHECK=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
# Set working directory
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application code
COPY app.py .
# Create necessary directories and set permissions
RUN mkdir -p /app/logs && \
chown -R appuser:appuser /app
# Switch to non-root user
USER appuser
# Expose port
EXPOSE 5001
# Health check
HEALTHCHECK --interval=30s --timeout=3s --start-period=5s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5001/health', timeout=2)" || exit 1
# Run with gunicorn for production
CMD ["gunicorn", "--bind", "0.0.0.0:5001", "--workers", "4", "--timeout", "120", "--access-logfile", "-", "--error-logfile", "-", "app:app"]
services/document-ingestion/.dockerignore
__pycache__
*.pyc
*.pyo
*.pyd
.Python
env/
venv/
.pytest_cache
.coverage
*.log
.git
.gitignore
README.md
tests/
*.md
services/document-ingestion/tests/test_app.py
"""
Unit tests for Document Ingestion Service
"""
import pytest
import json
from app import app, DocumentChunker
@pytest.fixture
def client():
app.config['TESTING'] = True
with app.test_client() as client:
yield client
def test_health_check(client):
"""Test health check endpoint"""
response = client.get('/health')
assert response.status_code == 200
data = json.loads(response.data)
assert data['status'] == 'healthy'
assert data['service'] == 'document-ingestion'
def test_readiness_check(client):
"""Test readiness check endpoint"""
response = client.get('/ready')
assert response.status_code == 200
data = json.loads(response.data)
assert data['status'] == 'ready'
def test_ingest_text(client):
"""Test document ingestion with text"""
payload = {
'text': 'This is a test document. ' * 100,
'metadata': {'source': 'test'}
}
response = client.post('/ingest',
data=json.dumps(payload),
content_type='application/json')
assert response.status_code == 200
data = json.loads(response.data)
assert data['success'] is True
assert data['total_chunks'] > 0
assert len(data['chunks']) > 0
def test_ingest_empty_text(client):
"""Test ingestion with empty text"""
payload = {'text': ''}
response = client.post('/ingest',
data=json.dumps(payload),
content_type='application/json')
assert response.status_code == 400
def test_chunker():
"""Test DocumentChunker directly"""
chunker = DocumentChunker(chunk_size=100, chunk_overlap=10)
text = "This is a test. " * 50
chunks = chunker.chunk_text(text, {'test': 'metadata'})
assert len(chunks) > 0
assert all('text' in chunk for chunk in chunks)
assert all('metadata' in chunk for chunk in chunks)
assert all(chunk['metadata']['test'] == 'metadata' for chunk in chunks)
def test_config_endpoint(client):
"""Test configuration endpoint"""
response = client.get('/config')
assert response.status_code == 200
data = json.loads(response.data)
assert 'chunk_size' in data
assert 'chunk_overlap' in data
2. Vector Database Service
services/vector-database/app.py
"""
Vector Database Service for RAG-based LLM Chatbot
Handles embedding storage and similarity search
"""
from flask import Flask, request, jsonify
from flask_cors import CORS
import numpy as np
from sentence_transformers import SentenceTransformer
from sklearn.metrics.pairwise import cosine_similarity
import logging
import os
from typing import List, Dict, Tuple, Optional
import threading
import pickle
import time
import signal
from prometheus_client import Counter, Histogram, Gauge, generate_latest, CONTENT_TYPE_LATEST
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
CORS(app)
# Configuration
EMBEDDING_MODEL = os.environ.get('EMBEDDING_MODEL', 'all-MiniLM-L6-v2')
TOP_K = int(os.environ.get('TOP_K', '5'))
PERSIST_PATH = os.environ.get('PERSIST_PATH', '/app/data/vectors.pkl')
# Prometheus metrics
REQUEST_COUNT = Counter('http_requests_total', 'Total HTTP requests', ['method', 'endpoint', 'status'])
REQUEST_DURATION = Histogram('http_request_duration_seconds', 'HTTP request duration', ['method', 'endpoint'])
DOCUMENTS_STORED = Gauge('documents_stored_total', 'Total number of documents stored')
SEARCH_QUERIES = Counter('search_queries_total', 'Total number of search queries')
class GracefulShutdown:
"""Handle graceful shutdown"""
def __init__(self):
self.shutdown_requested = False
signal.signal(signal.SIGTERM, self.handle_sigterm)
signal.signal(signal.SIGINT, self.handle_sigterm)
def handle_sigterm(self, signum, frame):
logger.info("Shutdown signal received")
self.shutdown_requested = True
def is_shutting_down(self):
return self.shutdown_requested
shutdown_handler = GracefulShutdown()
class VectorStore:
"""
In-memory vector store with persistence support.
"""
def __init__(self, model_name: str, persist_path: Optional[str] = None):
logger.info(f"Loading embedding model: {model_name}")
self.model = SentenceTransformer(model_name)
self.embeddings = []
self.documents = []
self.lock = threading.Lock()
self.persist_path = persist_path
# Try to load existing data
if persist_path and os.path.exists(persist_path):
self.load()
logger.info("Vector store initialized successfully")
def add_documents(self, documents: List[Dict]) -> bool:
"""Add documents and generate embeddings"""
try:
with self.lock:
# Extract text from documents
texts = [doc['text'] for doc in documents]
# Generate embeddings
logger.info(f"Generating embeddings for {len(texts)} documents")
new_embeddings = self.model.encode(texts, show_progress_bar=False)
# Store embeddings and documents
self.embeddings.extend(new_embeddings)
self.documents.extend(documents)
# Update metrics
DOCUMENTS_STORED.set(len(self.documents))
logger.info(f"Added {len(documents)} documents. Total: {len(self.documents)}")
# Persist to disk
if self.persist_path:
self.save()
return True
except Exception as e:
logger.error(f"Error adding documents: {str(e)}", exc_info=True)
return False
def search(self, query: str, top_k: int = 5) -> List[Tuple[Dict, float]]:
"""Search for similar documents"""
try:
with self.lock:
if not self.documents:
logger.warning("Search attempted on empty vector store")
return []
# Generate query embedding
query_embedding = self.model.encode([query], show_progress_bar=False)
# Calculate cosine similarity
embeddings_array = np.array(self.embeddings)
similarities = cosine_similarity(query_embedding, embeddings_array)[0]
# Get top-k indices
top_k = min(top_k, len(self.documents))
top_indices = np.argsort(similarities)[-top_k:][::-1]
# Prepare results
results = [
(self.documents[idx], float(similarities[idx]))
for idx in top_indices
]
SEARCH_QUERIES.inc()
logger.info(f"Search completed. Returned {len(results)} results")
return results
except Exception as e:
logger.error(f"Error during search: {str(e)}", exc_info=True)
return []
def get_stats(self) -> Dict:
"""Get vector store statistics"""
with self.lock:
return {
'total_documents': len(self.documents),
'embedding_dimension': len(self.embeddings[0]) if self.embeddings else 0,
'model': EMBEDDING_MODEL
}
def save(self):
"""Persist vector store to disk"""
try:
os.makedirs(os.path.dirname(self.persist_path), exist_ok=True)
with open(self.persist_path, 'wb') as f:
pickle.dump({
'embeddings': self.embeddings,
'documents': self.documents
}, f)
logger.info(f"Vector store saved to {self.persist_path}")
except Exception as e:
logger.error(f"Error saving vector store: {str(e)}", exc_info=True)
def load(self):
"""Load vector store from disk"""
try:
with open(self.persist_path, 'rb') as f:
data = pickle.load(f)
self.embeddings = data['embeddings']
self.documents = data['documents']
DOCUMENTS_STORED.set(len(self.documents))
logger.info(f"Loaded {len(self.documents)} documents from {self.persist_path}")
except Exception as e:
logger.error(f"Error loading vector store: {str(e)}", exc_info=True)
def clear(self):
"""Clear all documents"""
with self.lock:
self.embeddings = []
self.documents = []
DOCUMENTS_STORED.set(0)
if self.persist_path and os.path.exists(self.persist_path):
os.remove(self.persist_path)
logger.info("Vector store cleared")
# Initialize vector store
logger.info("Initializing vector store...")
vector_store = VectorStore(EMBEDDING_MODEL, PERSIST_PATH)
@app.before_request
def before_request():
"""Track request start time"""
if shutdown_handler.is_shutting_down():
return jsonify({'error': 'Service is shutting down'}), 503
request.start_time = time.time()
@app.after_request
def after_request(response):
"""Track request metrics"""
if hasattr(request, 'start_time'):
duration = time.time() - request.start_time
REQUEST_DURATION.labels(
method=request.method,
endpoint=request.endpoint or 'unknown'
).observe(duration)
REQUEST_COUNT.labels(
method=request.method,
endpoint=request.endpoint or 'unknown',
status=response.status_code
).inc()
return response
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint"""
stats = vector_store.get_stats()
return jsonify({
'status': 'healthy',
'service': 'vector-database',
'version': '1.0.0',
'stats': stats
}), 200
@app.route('/ready', methods=['GET'])
def readiness_check():
"""Readiness check endpoint"""
if shutdown_handler.is_shutting_down():
return jsonify({'status': 'not ready'}), 503
return jsonify({'status': 'ready'}), 200
@app.route('/add', methods=['POST'])
def add_documents():
"""Add documents to the vector store"""
try:
data = request.get_json()
if not data or 'documents' not in data:
return jsonify({'error': 'No documents provided'}), 400
documents = data['documents']
if not isinstance(documents, list):
return jsonify({'error': 'Documents must be a list'}), 400
# Validate document structure
for doc in documents:
if 'text' not in doc:
return jsonify({'error': 'Each document must have a text field'}), 400
# Add documents to vector store
success = vector_store.add_documents(documents)
if success:
return jsonify({
'success': True,
'added': len(documents),
'total': len(vector_store.documents)
}), 200
else:
return jsonify({'error': 'Failed to add documents'}), 500
except Exception as e:
logger.error(f"Error in add_documents: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/search', methods=['POST'])
def search():
"""Search for similar documents"""
try:
data = request.get_json()
if not data or 'query' not in data:
return jsonify({'error': 'No query provided'}), 400
query = data['query']
top_k = data.get('top_k', TOP_K)
# Perform search
results = vector_store.search(query, top_k)
# Format results
formatted_results = [
{
'document': doc,
'similarity': score
}
for doc, score in results
]
return jsonify({
'success': True,
'query': query,
'results': formatted_results,
'count': len(formatted_results)
}), 200
except Exception as e:
logger.error(f"Error in search: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/stats', methods=['GET'])
def get_stats():
"""Get vector store statistics"""
stats = vector_store.get_stats()
return jsonify(stats), 200
@app.route('/clear', methods=['POST'])
def clear():
"""Clear all documents (use with caution)"""
try:
vector_store.clear()
return jsonify({'success': True, 'message': 'Vector store cleared'}), 200
except Exception as e:
logger.error(f"Error clearing vector store: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/metrics', methods=['GET'])
def metrics():
"""Prometheus metrics endpoint"""
return generate_latest(), 200, {'Content-Type': CONTENT_TYPE_LATEST}
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5002))
logger.info(f"Starting Vector Database Service on port {port}")
app.run(host='0.0.0.0', port=port, debug=False)
services/vector-database/requirements.txt
flask==3.0.0
flask-cors==4.0.0
numpy==1.26.2
scikit-learn==1.3.2
sentence-transformers==2.2.2
prometheus-client==0.19.0
gunicorn==21.2.0
services/vector-database/Dockerfile
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
PIP_NO_CACHE_DIR=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create data directory for persistence
RUN mkdir -p /app/data && \
chown -R appuser:appuser /app
USER appuser
EXPOSE 5002
HEALTHCHECK --interval=30s --timeout=5s --start-period=60s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5002/health', timeout=3)" || exit 1
CMD ["gunicorn", "--bind", "0.0.0.0:5002", "--workers", "2", "--timeout", "300", "--access-logfile", "-", "--error-logfile", "-", "app:app"]
services/vector-database/.dockerignore
__pycache__
*.pyc
*.pyo
*.pyd
.Python
env/
venv/
.pytest_cache
.coverage
*.log
.git
tests/
*.md
data/
3. LLM Inference Service
services/llm-inference/app.py
"""
LLM Inference Service for RAG-based Chatbot
Supports multiple backends: CPU, CUDA, ROCm
"""
from flask import Flask, request, jsonify
from flask_cors import CORS
import torch
from transformers import AutoModelForCausalLM, AutoTokenizer, BitsAndBytesConfig
import logging
import os
from typing import Dict, Optional
import time
import signal
from prometheus_client import Counter, Histogram, Gauge, generate_latest, CONTENT_TYPE_LATEST
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
CORS(app)
# Configuration
MODEL_NAME = os.environ.get('MODEL_NAME', 'TinyLlama/TinyLlama-1.1B-Chat-v1.0')
MAX_LENGTH = int(os.environ.get('MAX_LENGTH', '512'))
TEMPERATURE = float(os.environ.get('TEMPERATURE', '0.7'))
USE_QUANTIZATION = os.environ.get('USE_QUANTIZATION', 'false').lower() == 'true'
# Prometheus metrics
REQUEST_COUNT = Counter('http_requests_total', 'Total HTTP requests', ['method', 'endpoint', 'status'])
REQUEST_DURATION = Histogram('http_request_duration_seconds', 'HTTP request duration', ['method', 'endpoint'])
GENERATION_DURATION = Histogram('generation_duration_seconds', 'LLM generation duration')
TOKENS_GENERATED = Counter('tokens_generated_total', 'Total tokens generated')
MODEL_LOADED = Gauge('model_loaded', 'Whether model is loaded')
class GracefulShutdown:
"""Handle graceful shutdown"""
def __init__(self):
self.shutdown_requested = False
signal.signal(signal.SIGTERM, self.handle_sigterm)
signal.signal(signal.SIGINT, self.handle_sigterm)
def handle_sigterm(self, signum, frame):
logger.info("Shutdown signal received")
self.shutdown_requested = True
def is_shutting_down(self):
return self.shutdown_requested
shutdown_handler = GracefulShutdown()
def detect_device() -> str:
"""Detect the best available device"""
if torch.cuda.is_available():
device = 'cuda'
gpu_name = torch.cuda.get_device_name(0)
logger.info(f"NVIDIA CUDA available: {gpu_name}")
elif hasattr(torch.backends, 'mps') and torch.backends.mps.is_available():
device = 'mps'
logger.info("Apple Silicon (MPS) available")
else:
device = 'cpu'
logger.info("No GPU available, using CPU")
return device
class LLMInferenceEngine:
"""Handles LLM model loading and inference"""
def __init__(self, model_name: str, device: str, use_quantization: bool):
self.model_name = model_name
self.device = device
self.use_quantization = use_quantization and device in ['cuda']
logger.info(f"Loading model: {model_name}")
logger.info(f"Device: {device}")
logger.info(f"Quantization: {self.use_quantization}")
try:
# Load tokenizer
self.tokenizer = AutoTokenizer.from_pretrained(
model_name,
trust_remote_code=True
)
# Set padding token
if self.tokenizer.pad_token is None:
self.tokenizer.pad_token = self.tokenizer.eos_token
# Load model with appropriate configuration
if self.use_quantization:
quantization_config = BitsAndBytesConfig(
load_in_4bit=True,
bnb_4bit_compute_dtype=torch.float16,
bnb_4bit_use_double_quant=True,
bnb_4bit_quant_type="nf4"
)
self.model = AutoModelForCausalLM.from_pretrained(
model_name,
quantization_config=quantization_config,
device_map="auto",
trust_remote_code=True
)
else:
self.model = AutoModelForCausalLM.from_pretrained(
model_name,
torch_dtype=torch.float16 if device != 'cpu' else torch.float32,
trust_remote_code=True,
low_cpu_mem_usage=True
)
self.model.to(device)
# Set to evaluation mode
self.model.eval()
MODEL_LOADED.set(1)
logger.info("Model loaded successfully")
except Exception as e:
logger.error(f"Failed to load model: {str(e)}", exc_info=True)
MODEL_LOADED.set(0)
raise
def generate(
self,
prompt: str,
max_length: int = 512,
temperature: float = 0.7,
top_p: float = 0.9,
top_k: int = 50
) -> str:
"""Generate text based on prompt"""
try:
start_time = time.time()
# Tokenize input
inputs = self.tokenizer(
prompt,
return_tensors="pt",
truncation=True,
max_length=2048
)
# Move to device
inputs = {k: v.to(self.device) for k, v in inputs.items()}
# Generate
with torch.no_grad():
outputs = self.model.generate(
**inputs,
max_new_tokens=max_length,
temperature=temperature,
top_p=top_p,
top_k=top_k,
do_sample=True,
pad_token_id=self.tokenizer.pad_token_id,
eos_token_id=self.tokenizer.eos_token_id
)
# Decode output
generated_text = self.tokenizer.decode(
outputs[0],
skip_special_tokens=True
)
# Remove prompt from output
if generated_text.startswith(prompt):
generated_text = generated_text[len(prompt):].strip()
# Update metrics
duration = time.time() - start_time
GENERATION_DURATION.observe(duration)
TOKENS_GENERATED.inc(len(outputs[0]))
logger.info(f"Generated {len(outputs[0])} tokens in {duration:.2f}s")
return generated_text
except Exception as e:
logger.error(f"Error during generation: {str(e)}", exc_info=True)
raise
def get_model_info(self) -> Dict:
"""Get model information"""
return {
'model_name': self.model_name,
'device': self.device,
'quantization': self.use_quantization,
'parameters': sum(p.numel() for p in self.model.parameters()),
'dtype': str(next(self.model.parameters()).dtype)
}
# Initialize inference engine
logger.info("Initializing LLM inference engine...")
device = detect_device()
inference_engine = LLMInferenceEngine(MODEL_NAME, device, USE_QUANTIZATION)
@app.before_request
def before_request():
"""Track request start time"""
if shutdown_handler.is_shutting_down():
return jsonify({'error': 'Service is shutting down'}), 503
request.start_time = time.time()
@app.after_request
def after_request(response):
"""Track request metrics"""
if hasattr(request, 'start_time'):
duration = time.time() - request.start_time
REQUEST_DURATION.labels(
method=request.method,
endpoint=request.endpoint or 'unknown'
).observe(duration)
REQUEST_COUNT.labels(
method=request.method,
endpoint=request.endpoint or 'unknown',
status=response.status_code
).inc()
return response
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint"""
model_info = inference_engine.get_model_info()
return jsonify({
'status': 'healthy',
'service': 'llm-inference',
'version': '1.0.0',
'model': model_info
}), 200
@app.route('/ready', methods=['GET'])
def readiness_check():
"""Readiness check endpoint"""
if shutdown_handler.is_shutting_down():
return jsonify({'status': 'not ready'}), 503
return jsonify({'status': 'ready'}), 200
@app.route('/generate', methods=['POST'])
def generate():
"""Generate text based on prompt"""
try:
data = request.get_json()
if not data or 'prompt' not in data:
return jsonify({'error': 'No prompt provided'}), 400
prompt = data['prompt']
max_length = data.get('max_length', MAX_LENGTH)
temperature = data.get('temperature', TEMPERATURE)
top_p = data.get('top_p', 0.9)
top_k = data.get('top_k', 50)
# Validate parameters
if max_length > 2048:
return jsonify({'error': 'max_length cannot exceed 2048'}), 400
if not 0 < temperature <= 2:
return jsonify({'error': 'temperature must be between 0 and 2'}), 400
# Generate response
logger.info(f"Generating response for prompt of length {len(prompt)}")
generated_text = inference_engine.generate(
prompt=prompt,
max_length=max_length,
temperature=temperature,
top_p=top_p,
top_k=top_k
)
return jsonify({
'success': True,
'prompt': prompt,
'generated_text': generated_text,
'model': MODEL_NAME
}), 200
except Exception as e:
logger.error(f"Error in generate: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/model-info', methods=['GET'])
def model_info():
"""Get detailed model information"""
info = inference_engine.get_model_info()
return jsonify(info), 200
@app.route('/metrics', methods=['GET'])
def metrics():
"""Prometheus metrics endpoint"""
return generate_latest(), 200, {'Content-Type': CONTENT_TYPE_LATEST}
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5003))
logger.info(f"Starting LLM Inference Service on port {port}")
app.run(host='0.0.0.0', port=port, debug=False, threaded=True)
services/llm-inference/requirements.txt
flask==3.0.0
flask-cors==4.0.0
torch==2.1.0
transformers==4.36.0
accelerate==0.25.0
sentencepiece==0.1.99
prometheus-client==0.19.0
gunicorn==21.2.0
services/llm-inference/Dockerfile.cpu
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
PIP_NO_CACHE_DIR=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create cache directory for models
RUN mkdir -p /app/.cache && \
chown -R appuser:appuser /app
USER appuser
EXPOSE 5003
HEALTHCHECK --interval=30s --timeout=10s --start-period=180s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5003/health', timeout=5)" || exit 1
CMD ["gunicorn", "--bind", "0.0.0.0:5003", "--workers", "1", "--threads", "4", "--timeout", "600", "--access-logfile", "-", "--error-logfile", "-", "app:app"]
services/llm-inference/Dockerfile.cuda
FROM nvidia/cuda:12.1.0-runtime-ubuntu22.04
# Install Python
RUN apt-get update && apt-get install -y \
python3.11 \
python3-pip \
&& rm -rf /var/lib/apt/lists/*
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
PIP_NO_CACHE_DIR=1 \
CUDA_VISIBLE_DEVICES=0
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install PyTorch with CUDA
RUN pip3 install --no-cache-dir \
torch==2.1.0 \
--index-url https://download.pytorch.org/whl/cu121
# Install other dependencies
COPY requirements.txt .
RUN pip3 install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Create cache directory
RUN mkdir -p /app/.cache && \
chown -R appuser:appuser /app
USER appuser
EXPOSE 5003
HEALTHCHECK --interval=30s --timeout=10s --start-period=180s --retries=3 \
CMD python3 -c "import requests; requests.get('http://localhost:5003/health', timeout=5)" || exit 1
CMD ["gunicorn", "--bind", "0.0.0.0:5003", "--workers", "1", "--threads", "4", "--timeout", "600", "--access-logfile", "-", "--error-logfile", "-", "app:app"]
4. API Gateway Service
services/api-gateway/app.py
"""
API Gateway Service for RAG-based LLM Chatbot
Orchestrates requests between all microservices
"""
from flask import Flask, request, jsonify
from flask_cors import CORS
import requests
import logging
import os
from typing import List, Dict, Optional
import time
import signal
from prometheus_client import Counter, Histogram, generate_latest, CONTENT_TYPE_LATEST
from pybreaker import CircuitBreaker
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
app = Flask(__name__)
CORS(app)
# Service endpoints
INGESTION_SERVICE = os.environ.get('INGESTION_SERVICE', 'http://localhost:5001')
VECTOR_SERVICE = os.environ.get('VECTOR_SERVICE', 'http://localhost:5002')
LLM_SERVICE = os.environ.get('LLM_SERVICE', 'http://localhost:5003')
# Configuration
TOP_K_DOCUMENTS = int(os.environ.get('TOP_K_DOCUMENTS', '3'))
REQUEST_TIMEOUT = int(os.environ.get('REQUEST_TIMEOUT', '30'))
# Prometheus metrics
REQUEST_COUNT = Counter('http_requests_total', 'Total HTTP requests', ['method', 'endpoint', 'status'])
REQUEST_DURATION = Histogram('http_request_duration_seconds', 'HTTP request duration', ['method', 'endpoint'])
RAG_PIPELINE_DURATION = Histogram('rag_pipeline_duration_seconds', 'RAG pipeline duration')
SERVICE_CALLS = Counter('service_calls_total', 'Total service calls', ['service', 'status'])
class GracefulShutdown:
"""Handle graceful shutdown"""
def __init__(self):
self.shutdown_requested = False
signal.signal(signal.SIGTERM, self.handle_sigterm)
signal.signal(signal.SIGINT, self.handle_sigterm)
def handle_sigterm(self, signum, frame):
logger.info("Shutdown signal received")
self.shutdown_requested = True
def is_shutting_down(self):
return self.shutdown_requested
shutdown_handler = GracefulShutdown()
class RAGOrchestrator:
"""Orchestrates the RAG pipeline across microservices"""
def __init__(self, ingestion_url: str, vector_url: str, llm_url: str, timeout: int):
self.ingestion_url = ingestion_url
self.vector_url = vector_url
self.llm_url = llm_url
self.timeout = timeout
# Circuit breakers
self.ingestion_breaker = CircuitBreaker(
fail_max=5,
timeout_duration=60,
name='ingestion-service'
)
self.vector_breaker = CircuitBreaker(
fail_max=5,
timeout_duration=60,
name='vector-service'
)
self.llm_breaker = CircuitBreaker(
fail_max=3,
timeout_duration=120,
name='llm-service'
)
logger.info("Initialized RAG orchestrator")
logger.info(f"Ingestion: {ingestion_url}")
logger.info(f"Vector: {vector_url}")
logger.info(f"LLM: {llm_url}")
def ingest_document(self, text: str, metadata: Optional[Dict] = None) -> Dict:
"""Ingest document through ingestion service"""
@self.ingestion_breaker
def _ingest():
response = requests.post(
f"{self.ingestion_url}/ingest",
json={'text': text, 'metadata': metadata or {}},
timeout=self.timeout
)
response.raise_for_status()
SERVICE_CALLS.labels(service='ingestion', status='success').inc()
return response.json()
try:
return _ingest()
except Exception as e:
SERVICE_CALLS.labels(service='ingestion', status='error').inc()
logger.error(f"Error calling ingestion service: {str(e)}")
raise
def add_to_vector_store(self, documents: List[Dict]) -> Dict:
"""Add documents to vector store"""
@self.vector_breaker
def _add():
response = requests.post(
f"{self.vector_url}/add",
json={'documents': documents},
timeout=self.timeout
)
response.raise_for_status()
SERVICE_CALLS.labels(service='vector', status='success').inc()
return response.json()
try:
return _add()
except Exception as e:
SERVICE_CALLS.labels(service='vector', status='error').inc()
logger.error(f"Error calling vector service: {str(e)}")
raise
def search_documents(self, query: str, top_k: int = 3) -> List[Dict]:
"""Search for relevant documents"""
@self.vector_breaker
def _search():
response = requests.post(
f"{self.vector_url}/search",
json={'query': query, 'top_k': top_k},
timeout=self.timeout
)
response.raise_for_status()
SERVICE_CALLS.labels(service='vector', status='success').inc()
return response.json().get('results', [])
try:
return _search()
except Exception as e:
SERVICE_CALLS.labels(service='vector', status='error').inc()
logger.error(f"Error calling vector service: {str(e)}")
return []
def generate_response(
self,
prompt: str,
max_length: int = 512,
temperature: float = 0.7
) -> str:
"""Generate response using LLM service"""
@self.llm_breaker
def _generate():
response = requests.post(
f"{self.llm_url}/generate",
json={
'prompt': prompt,
'max_length': max_length,
'temperature': temperature
},
timeout=self.timeout * 3 # LLM needs more time
)
response.raise_for_status()
SERVICE_CALLS.labels(service='llm', status='success').inc()
return response.json().get('generated_text', '')
try:
return _generate()
except Exception as e:
SERVICE_CALLS.labels(service='llm', status='error').inc()
logger.error(f"Error calling LLM service: {str(e)}")
raise
def construct_rag_prompt(self, query: str, context_docs: List[Dict]) -> str:
"""Construct prompt with retrieved context"""
context_texts = [
doc['document']['text']
for doc in context_docs
if 'document' in doc and 'text' in doc['document']
]
context = "\n\n".join(context_texts)
prompt = f"""You are a helpful assistant. Use the following context to answer the question. If you cannot answer based on the context, say so.
Context:
{context}
Question: {query}
Answer:"""
return prompt
def answer_query(
self,
query: str,
top_k: int = 3,
max_length: int = 512,
temperature: float = 0.7
) -> Dict:
"""Complete RAG pipeline"""
start_time = time.time()
try:
# Retrieve relevant documents
logger.info(f"Searching for documents: {query}")
context_docs = self.search_documents(query, top_k)
if not context_docs:
logger.warning("No relevant documents found")
return {
'answer': "I don't have enough information to answer this question.",
'sources': [],
'query': query
}
# Construct prompt
prompt = self.construct_rag_prompt(query, context_docs)
# Generate answer
logger.info("Generating answer with LLM")
answer = self.generate_response(prompt, max_length, temperature)
# Format response
sources = [
{
'text': doc['document']['text'][:200] + '...',
'similarity': doc['similarity'],
'metadata': doc['document'].get('metadata', {})
}
for doc in context_docs
]
duration = time.time() - start_time
RAG_PIPELINE_DURATION.observe(duration)
return {
'answer': answer,
'sources': sources,
'query': query,
'duration': duration
}
except Exception as e:
logger.error(f"Error in RAG pipeline: {str(e)}", exc_info=True)
raise
# Initialize orchestrator
orchestrator = RAGOrchestrator(
INGESTION_SERVICE,
VECTOR_SERVICE,
LLM_SERVICE,
REQUEST_TIMEOUT
)
@app.before_request
def before_request():
"""Track request start time"""
if shutdown_handler.is_shutting_down():
return jsonify({'error': 'Service is shutting down'}), 503
request.start_time = time.time()
@app.after_request
def after_request(response):
"""Track request metrics"""
if hasattr(request, 'start_time'):
duration = time.time() - request.start_time
REQUEST_DURATION.labels(
method=request.method,
endpoint=request.endpoint or 'unknown'
).observe(duration)
REQUEST_COUNT.labels(
method=request.method,
endpoint=request.endpoint or 'unknown',
status=response.status_code
).inc()
return response
@app.route('/health', methods=['GET'])
def health_check():
"""Health check endpoint"""
services_status = {}
for service_name, service_url in [
('ingestion', INGESTION_SERVICE),
('vector', VECTOR_SERVICE),
('llm', LLM_SERVICE)
]:
try:
response = requests.get(f"{service_url}/health", timeout=5)
services_status[service_name] = 'healthy' if response.ok else 'unhealthy'
except Exception:
services_status[service_name] = 'unreachable'
overall_status = 'healthy' if all(
status == 'healthy' for status in services_status.values()
) else 'degraded'
return jsonify({
'status': overall_status,
'service': 'api-gateway',
'version': '1.0.0',
'services': services_status
}), 200
@app.route('/ready', methods=['GET'])
def readiness_check():
"""Readiness check endpoint"""
if shutdown_handler.is_shutting_down():
return jsonify({'status': 'not ready'}), 503
return jsonify({'status': 'ready'}), 200
@app.route('/ingest', methods=['POST'])
def ingest():
"""Ingest document into the system"""
try:
data = request.get_json()
if not data or 'text' not in data:
return jsonify({'error': 'No text provided'}), 400
text = data['text']
metadata = data.get('metadata', {})
# Chunk document
logger.info("Ingesting document")
ingestion_result = orchestrator.ingest_document(text, metadata)
chunks = ingestion_result.get('chunks', [])
if not chunks:
return jsonify({'error': 'Failed to chunk document'}), 500
# Add to vector store
logger.info(f"Adding {len(chunks)} chunks to vector store")
vector_result = orchestrator.add_to_vector_store(chunks)
return jsonify({
'success': True,
'chunks_created': len(chunks),
'chunks_stored': vector_result.get('added', 0)
}), 200
except Exception as e:
logger.error(f"Error in ingest: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/query', methods=['POST'])
def query():
"""Answer query using RAG pipeline"""
try:
data = request.get_json()
if not data or 'query' not in data:
return jsonify({'error': 'No query provided'}), 400
user_query = data['query']
top_k = data.get('top_k', TOP_K_DOCUMENTS)
max_length = data.get('max_length', 512)
temperature = data.get('temperature', 0.7)
# Execute RAG pipeline
logger.info(f"Processing query: {user_query}")
result = orchestrator.answer_query(
user_query,
top_k,
max_length,
temperature
)
return jsonify({
'success': True,
**result
}), 200
except Exception as e:
logger.error(f"Error in query: {str(e)}", exc_info=True)
return jsonify({'error': str(e)}), 500
@app.route('/metrics', methods=['GET'])
def metrics():
"""Prometheus metrics endpoint"""
return generate_latest(), 200, {'Content-Type': CONTENT_TYPE_LATEST}
if __name__ == '__main__':
port = int(os.environ.get('PORT', 5000))
logger.info(f"Starting API Gateway Service on port {port}")
app.run(host='0.0.0.0', port=port, debug=False)
services/api-gateway/requirements.txt
flask==3.0.0
flask-cors==4.0.0
requests==2.31.0
pybreaker==1.0.1
prometheus-client==0.19.0
gunicorn==21.2.0
services/api-gateway/Dockerfile
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1 \
PIP_NO_CACHE_DIR=1
# Create non-root user
RUN groupadd -r appuser && useradd -r -g appuser appuser
WORKDIR /app
# Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copy application
COPY app.py .
# Set ownership
RUN chown -R appuser:appuser /app
USER appuser
EXPOSE 5000
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
CMD python -c "import requests; requests.get('http://localhost:5000/health', timeout=3)" || exit 1
CMD ["gunicorn", "--bind", "0.0.0.0:5000", "--workers", "4", "--timeout", "300", "--access-logfile", "-", "--error-logfile", "-", "app:app"]
5. Kubernetes Manifests
k8s/namespace.yaml
apiVersion: v1
kind: Namespace
metadata:
name: rag-chatbot
labels:
name: rag-chatbot
environment: production
pod-security.kubernetes.io/enforce: baseline
pod-security.kubernetes.io/audit: baseline
pod-security.kubernetes.io/warn: baseline
k8s/configmaps.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: document-ingestion-config
namespace: rag-chatbot
data:
CHUNK_SIZE: "512"
CHUNK_OVERLAP: "50"
MAX_FILE_SIZE: "10485760"
PORT: "5001"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: vector-database-config
namespace: rag-chatbot
data:
EMBEDDING_MODEL: "all-MiniLM-L6-v2"
TOP_K: "5"
PORT: "5002"
PERSIST_PATH: "/app/data/vectors.pkl"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: llm-inference-config
namespace: rag-chatbot
data:
MODEL_NAME: "TinyLlama/TinyLlama-1.1B-Chat-v1.0"
MAX_LENGTH: "512"
TEMPERATURE: "0.7"
USE_QUANTIZATION: "false"
PORT: "5003"
---
apiVersion: v1
kind: ConfigMap
metadata:
name: api-gateway-config
namespace: rag-chatbot
data:
TOP_K_DOCUMENTS: "3"
REQUEST_TIMEOUT: "30"
PORT: "5000"
INGESTION_SERVICE: "http://document-ingestion-service:80"
VECTOR_SERVICE: "http://vector-database-service:80"
LLM_SERVICE: "http://llm-inference-service:80"
k8s/document-ingestion/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: document-ingestion
namespace: rag-chatbot
labels:
app: document-ingestion
component: backend
version: v1
spec:
replicas: 2
selector:
matchLabels:
app: document-ingestion
template:
metadata:
labels:
app: document-ingestion
component: backend
version: v1
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "5001"
prometheus.io/path: "/metrics"
spec:
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
containers:
- name: document-ingestion
image: rag-chatbot/document-ingestion:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5001
name: http
protocol: TCP
envFrom:
- configMapRef:
name: document-ingestion-config
resources:
requests:
memory: "256Mi"
cpu: "250m"
limits:
memory: "512Mi"
cpu: "500m"
livenessProbe:
httpGet:
path: /health
port: 5001
initialDelaySeconds: 15
periodSeconds: 20
timeoutSeconds: 3
failureThreshold: 3
readinessProbe:
httpGet:
path: /ready
port: 5001
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 3
failureThreshold: 3
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
readOnlyRootFilesystem: false
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 15"]
k8s/document-ingestion/service.yaml
apiVersion: v1
kind: Service
metadata:
name: document-ingestion-service
namespace: rag-chatbot
labels:
app: document-ingestion
spec:
selector:
app: document-ingestion
ports:
- protocol: TCP
port: 80
targetPort: 5001
name: http
type: ClusterIP
sessionAffinity: None
k8s/vector-database/pvc.yaml
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: vector-database-pvc
namespace: rag-chatbot
spec:
accessModes:
- ReadWriteOnce
resources:
requests:
storage: 10Gi
storageClassName: standard
k8s/vector-database/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: vector-database
namespace: rag-chatbot
labels:
app: vector-database
component: backend
version: v1
spec:
replicas: 1
selector:
matchLabels:
app: vector-database
template:
metadata:
labels:
app: vector-database
component: backend
version: v1
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "5002"
prometheus.io/path: "/metrics"
spec:
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
containers:
- name: vector-database
image: rag-chatbot/vector-database:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5002
name: http
protocol: TCP
envFrom:
- configMapRef:
name: vector-database-config
volumeMounts:
- name: data
mountPath: /app/data
resources:
requests:
memory: "1Gi"
cpu: "500m"
limits:
memory: "2Gi"
cpu: "1000m"
livenessProbe:
httpGet:
path: /health
port: 5002
initialDelaySeconds: 60
periodSeconds: 30
timeoutSeconds: 5
failureThreshold: 3
readinessProbe:
httpGet:
path: /ready
port: 5002
initialDelaySeconds: 45
periodSeconds: 20
timeoutSeconds: 5
failureThreshold: 3
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 15"]
volumes:
- name: data
persistentVolumeClaim:
claimName: vector-database-pvc
k8s/vector-database/service.yaml
apiVersion: v1
kind: Service
metadata:
name: vector-database-service
namespace: rag-chatbot
labels:
app: vector-database
spec:
selector:
app: vector-database
ports:
- protocol: TCP
port: 80
targetPort: 5002
name: http
type: ClusterIP
k8s/llm-inference/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: llm-inference
namespace: rag-chatbot
labels:
app: llm-inference
component: backend
version: v1
spec:
replicas: 1
selector:
matchLabels:
app: llm-inference
template:
metadata:
labels:
app: llm-inference
component: backend
version: v1
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "5003"
prometheus.io/path: "/metrics"
spec:
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
containers:
- name: llm-inference
image: rag-chatbot/llm-inference:1.0-cpu
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5003
name: http
protocol: TCP
envFrom:
- configMapRef:
name: llm-inference-config
resources:
requests:
memory: "4Gi"
cpu: "2000m"
limits:
memory: "8Gi"
cpu: "4000m"
livenessProbe:
httpGet:
path: /health
port: 5003
initialDelaySeconds: 180
periodSeconds: 60
timeoutSeconds: 10
failureThreshold: 3
readinessProbe:
httpGet:
path: /ready
port: 5003
initialDelaySeconds: 150
periodSeconds: 30
timeoutSeconds: 10
failureThreshold: 3
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 30"]
k8s/llm-inference/service.yaml
apiVersion: v1
kind: Service
metadata:
name: llm-inference-service
namespace: rag-chatbot
labels:
app: llm-inference
spec:
selector:
app: llm-inference
ports:
- protocol: TCP
port: 80
targetPort: 5003
name: http
type: ClusterIP
k8s/api-gateway/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: api-gateway
namespace: rag-chatbot
labels:
app: api-gateway
component: frontend
version: v1
spec:
replicas: 2
selector:
matchLabels:
app: api-gateway
template:
metadata:
labels:
app: api-gateway
component: frontend
version: v1
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "5000"
prometheus.io/path: "/metrics"
spec:
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
containers:
- name: api-gateway
image: rag-chatbot/api-gateway:1.0
imagePullPolicy: IfNotPresent
ports:
- containerPort: 5000
name: http
protocol: TCP
envFrom:
- configMapRef:
name: api-gateway-config
resources:
requests:
memory: "128Mi"
cpu: "100m"
limits:
memory: "256Mi"
cpu: "200m"
livenessProbe:
httpGet:
path: /health
port: 5000
initialDelaySeconds: 10
periodSeconds: 20
timeoutSeconds: 5
failureThreshold: 3
readinessProbe:
httpGet:
path: /ready
port: 5000
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 5
failureThreshold: 3
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 15"]
k8s/api-gateway/service.yaml
apiVersion: v1
kind: Service
metadata:
name: api-gateway-service
namespace: rag-chatbot
labels:
app: api-gateway
spec:
selector:
app: api-gateway
ports:
- protocol: TCP
port: 80
targetPort: 5000
name: http
type: ClusterIP
k8s/ingress.yaml
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: rag-chatbot-ingress
namespace: rag-chatbot
annotations:
nginx.ingress.kubernetes.io/rewrite-target: /
nginx.ingress.kubernetes.io/proxy-body-size: "10m"
nginx.ingress.kubernetes.io/proxy-read-timeout: "300"
nginx.ingress.kubernetes.io/proxy-send-timeout: "300"
nginx.ingress.kubernetes.io/cors-allow-origin: "*"
nginx.ingress.kubernetes.io/enable-cors: "true"
spec:
ingressClassName: nginx
rules:
- host: rag-chatbot.local
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: api-gateway-service
port:
number: 80
k8s/hpa.yaml
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: api-gateway-hpa
namespace: rag-chatbot
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: api-gateway
minReplicas: 2
maxReplicas: 10
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 80
behavior:
scaleDown:
stabilizationWindowSeconds: 300
policies:
- type: Percent
value: 50
periodSeconds: 60
scaleUp:
stabilizationWindowSeconds: 0
policies:
- type: Percent
value: 100
periodSeconds: 30
- type: Pods
value: 2
periodSeconds: 30
selectPolicy: Max
---
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: document-ingestion-hpa
namespace: rag-chatbot
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: document-ingestion
minReplicas: 2
maxReplicas: 8
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
k8s/network-policies.yaml
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: api-gateway-policy
namespace: rag-chatbot
spec:
podSelector:
matchLabels:
app: api-gateway
policyTypes:
- Ingress
- Egress
ingress:
- from:
- namespaceSelector:
matchLabels:
name: ingress-nginx
ports:
- protocol: TCP
port: 5000
egress:
- to:
- podSelector:
matchLabels:
app: document-ingestion
ports:
- protocol: TCP
port: 5001
- to:
- podSelector:
matchLabels:
app: vector-database
ports:
- protocol: TCP
port: 5002
- to:
- podSelector:
matchLabels:
app: llm-inference
ports:
- protocol: TCP
port: 5003
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 53
- protocol: UDP
port: 53
---
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: document-ingestion-policy
namespace: rag-chatbot
spec:
podSelector:
matchLabels:
app: document-ingestion
policyTypes:
- Ingress
- Egress
ingress:
- from:
- podSelector:
matchLabels:
app: api-gateway
ports:
- protocol: TCP
port: 5001
egress:
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 53
- protocol: UDP
port: 53
---
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: vector-database-policy
namespace: rag-chatbot
spec:
podSelector:
matchLabels:
app: vector-database
policyTypes:
- Ingress
- Egress
ingress:
- from:
- podSelector:
matchLabels:
app: api-gateway
ports:
- protocol: TCP
port: 5002
egress:
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 53
- protocol: UDP
port: 53
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 443
---
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: llm-inference-policy
namespace: rag-chatbot
spec:
podSelector:
matchLabels:
app: llm-inference
policyTypes:
- Ingress
- Egress
ingress:
- from:
- podSelector:
matchLabels:
app: api-gateway
ports:
- protocol: TCP
port: 5003
egress:
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 53
- protocol: UDP
port: 53
- to:
- namespaceSelector: {}
ports:
- protocol: TCP
port: 443
6. Docker Compose for Local Development
docker-compose.yml
version: '3.8'
services:
document-ingestion:
build:
context: ./services/document-ingestion
dockerfile: Dockerfile
container_name: document-ingestion
ports:
- "5001:5001"
environment:
- CHUNK_SIZE=512
- CHUNK_OVERLAP=50
- MAX_FILE_SIZE=10485760
- PORT=5001
healthcheck:
test: ["CMD", "python", "-c", "import requests; requests.get('http://localhost:5001/health')"]
interval: 30s
timeout: 3s
retries: 3
start_period: 5s
networks:
- rag-network
vector-database:
build:
context: ./services/vector-database
dockerfile: Dockerfile
container_name: vector-database
ports:
- "5002:5002"
environment:
- EMBEDDING_MODEL=all-MiniLM-L6-v2
- TOP_K=5
- PORT=5002
- PERSIST_PATH=/app/data/vectors.pkl
volumes:
- vector-data:/app/data
healthcheck:
test: ["CMD", "python", "-c", "import requests; requests.get('http://localhost:5002/health')"]
interval: 30s
timeout: 5s
retries: 3
start_period: 60s
networks:
- rag-network
llm-inference:
build:
context: ./services/llm-inference
dockerfile: Dockerfile.cpu
container_name: llm-inference
ports:
- "5003:5003"
environment:
- MODEL_NAME=TinyLlama/TinyLlama-1.1B-Chat-v1.0
- MAX_LENGTH=512
- TEMPERATURE=0.7
- USE_QUANTIZATION=false
- PORT=5003
healthcheck:
test: ["CMD", "python", "-c", "import requests; requests.get('http://localhost:5003/health')"]
interval: 30s
timeout: 10s
retries: 3
start_period: 180s
networks:
- rag-network
deploy:
resources:
limits:
cpus: '4'
memory: 8G
reservations:
cpus: '2'
memory: 4G
api-gateway:
build:
context: ./services/api-gateway
dockerfile: Dockerfile
container_name: api-gateway
ports:
- "5000:5000"
environment:
- TOP_K_DOCUMENTS=3
- REQUEST_TIMEOUT=30
- PORT=5000
- INGESTION_SERVICE=http://document-ingestion:5001
- VECTOR_SERVICE=http://vector-database:5002
- LLM_SERVICE=http://llm-inference:5003
depends_on:
- document-ingestion
- vector-database
- llm-inference
healthcheck:
test: ["CMD", "python", "-c", "import requests; requests.get('http://localhost:5000/health')"]
interval: 30s
timeout: 5s
retries: 3
start_period: 10s
networks:
- rag-network
networks:
rag-network:
driver: bridge
volumes:
vector-data:
7. Deployment Scripts
scripts/build-images.sh
#!/bin/bash
set -e
echo "Building Docker images for RAG Chatbot..."
# Build document ingestion service
echo "Building document-ingestion..."
docker build -t rag-chatbot/document-ingestion:1.0 ./services/document-ingestion
# Build vector database service
echo "Building vector-database..."
docker build -t rag-chatbot/vector-database:1.0 ./services/vector-database
# Build LLM inference service (CPU version)
echo "Building llm-inference (CPU)..."
docker build -f ./services/llm-inference/Dockerfile.cpu -t rag-chatbot/llm-inference:1.0-cpu ./services/llm-inference
# Build API gateway service
echo "Building api-gateway..."
docker build -t rag-chatbot/api-gateway:1.0 ./services/api-gateway
echo "All images built successfully!"
docker images | grep rag-chatbot
scripts/deploy-local.sh
#!/bin/bash
set -e
echo "Deploying RAG Chatbot locally with Docker Compose..."
# Build images first
./scripts/build-images.sh
# Start services
echo "Starting services..."
docker-compose up -d
# Wait for services to be healthy
echo "Waiting for services to be healthy..."
sleep 10
# Check health
echo "Checking service health..."
curl -f http://localhost:5001/health || echo "Document Ingestion not ready"
curl -f http://localhost:5002/health || echo "Vector Database not ready"
curl -f http://localhost:5003/health || echo "LLM Inference not ready (this may take a few minutes)"
curl -f http://localhost:5000/health || echo "API Gateway not ready"
echo ""
echo "Deployment complete!"
echo "API Gateway: http://localhost:5000"
echo ""
echo "To view logs: docker-compose logs -f"
echo "To stop: docker-compose down"
scripts/deploy-k8s.sh
#!/bin/bash
set -e
echo "Deploying RAG Chatbot to Kubernetes..."
# Check if kubectl is available
if ! command -v kubectl &> /dev/null; then
echo "kubectl not found. Please install kubectl first."
exit 1
fi
# Build images
echo "Building Docker images..."
./scripts/build-images.sh
# Load images into minikube (if using minikube)
if command -v minikube &> /dev/null && minikube status &> /dev/null; then
echo "Loading images into minikube..."
minikube image load rag-chatbot/document-ingestion:1.0
minikube image load rag-chatbot/vector-database:1.0
minikube image load rag-chatbot/llm-inference:1.0-cpu
minikube image load rag-chatbot/api-gateway:1.0
fi
# Create namespace
echo "Creating namespace..."
kubectl apply -f k8s/namespace.yaml
# Apply ConfigMaps
echo "Applying ConfigMaps..."
kubectl apply -f k8s/configmaps.yaml
# Deploy services
echo "Deploying services..."
kubectl apply -f k8s/vector-database/pvc.yaml
kubectl apply -f k8s/document-ingestion/
kubectl apply -f k8s/vector-database/
kubectl apply -f k8s/llm-inference/
kubectl apply -f k8s/api-gateway/
# Apply Ingress
echo "Applying Ingress..."
kubectl apply -f k8s/ingress.yaml
# Apply HPA
echo "Applying HorizontalPodAutoscaler..."
kubectl apply -f k8s/hpa.yaml
# Apply Network Policies (optional)
echo "Applying Network Policies..."
kubectl apply -f k8s/network-policies.yaml
echo ""
echo "Deployment initiated!"
echo ""
echo "Check deployment status:"
echo " kubectl get pods -n rag-chatbot"
echo ""
echo "View logs:"
echo " kubectl logs -n rag-chatbot -l app=api-gateway"
echo ""
echo "If using minikube, run: minikube tunnel"
echo "Then add to /etc/hosts: 127.0.0.1 rag-chatbot.local"
echo "Access at: http://rag-chatbot.local"
scripts/test-system.sh
#!/bin/bash
set -e
# Determine base URL
if [ "$1" == "k8s" ]; then
BASE_URL="http://rag-chatbot.local"
else
BASE_URL="http://localhost:5000"
fi
echo "Testing RAG Chatbot System at $BASE_URL"
echo "=========================================="
# Test health endpoint
echo ""
echo "1. Testing health endpoint..."
curl -s $BASE_URL/health | jq .
# Ingest a document
echo ""
echo "2. Ingesting test document..."
curl -s -X POST $BASE_URL/ingest \
-H "Content-Type: application/json" \
-d '{
"text": "Docker is a platform for developing, shipping, and running applications in containers. Containers are lightweight, standalone, executable packages that include everything needed to run a piece of software. Kubernetes is an open-source container orchestration platform that automates deploying, scaling, and managing containerized applications. It was originally developed by Google and is now maintained by the Cloud Native Computing Foundation.",
"metadata": {"source": "docker-k8s-intro", "topic": "containers"}
}' | jq .
echo ""
echo "3. Waiting for indexing..."
sleep 5
# Query the system
echo ""
echo "4. Querying: What is Docker?"
curl -s -X POST $BASE_URL/query \
-H "Content-Type: application/json" \
-d '{
"query": "What is Docker?",
"top_k": 2,
"max_length": 150
}' | jq .
echo ""
echo "5. Querying: What is Kubernetes?"
curl -s -X POST $BASE_URL/query \
-H "Content-Type: application/json" \
-d '{
"query": "What is Kubernetes?",
"top_k": 2,
"max_length": 150
}' | jq .
echo ""
echo "=========================================="
echo "Testing complete!"
scripts/cleanup.sh
#!/bin/bash
echo "Cleaning up RAG Chatbot deployment..."
# Ask for confirmation
read -p "This will delete all resources. Continue? (y/n) " -n 1 -r
echo
if [[ ! $REPLY =~ ^[Yy]$ ]]; then
exit 1
fi
# Docker Compose cleanup
if [ -f "docker-compose.yml" ]; then
echo "Stopping Docker Compose services..."
docker-compose down -v
fi
# Kubernetes cleanup
if command -v kubectl &> /dev/null; then
echo "Deleting Kubernetes resources..."
kubectl delete namespace rag-chatbot --ignore-not-found=true
fi
# Remove Docker images
echo "Removing Docker images..."
docker rmi rag-chatbot/document-ingestion:1.0 2>/dev/null || true
docker rmi rag-chatbot/vector-database:1.0 2>/dev/null || true
docker rmi rag-chatbot/llm-inference:1.0-cpu 2>/dev/null || true
docker rmi rag-chatbot/api-gateway:1.0 2>/dev/null || true
echo "Cleanup complete!"
8. Makefile
Makefile
.PHONY: help build deploy-local deploy-k8s test clean
help:
@echo "RAG Chatbot - Makefile Commands"
@echo "================================"
@echo "build - Build all Docker images"
@echo "deploy-local - Deploy with Docker Compose"
@echo "deploy-k8s - Deploy to Kubernetes"
@echo "test-local - Test local deployment"
@echo "test-k8s - Test Kubernetes deployment"
@echo "clean - Clean up all resources"
@echo "logs-local - View Docker Compose logs"
@echo "logs-k8s - View Kubernetes logs"
build:
@chmod +x scripts/build-images.sh
@./scripts/build-images.sh
deploy-local:
@chmod +x scripts/deploy-local.sh
@./scripts/deploy-local.sh
deploy-k8s:
@chmod +x scripts/deploy-k8s.sh
@./scripts/deploy-k8s.sh
test-local:
@chmod +x scripts/test-system.sh
@./scripts/test-system.sh local
test-k8s:
@chmod +x scripts/test-system.sh
@./scripts/test-system.sh k8s
clean:
@chmod +x scripts/cleanup.sh
@./scripts/cleanup.sh
logs-local:
@docker-compose logs -f
logs-k8s:
@kubectl logs -n rag-chatbot -l app=api-gateway -f
9. README.md
README.md
# Production-Ready RAG-Based LLM Chatbot
A complete, production-ready implementation of a Retrieval-Augmented Generation (RAG) chatbot using Docker and Kubernetes.
## Architecture
The system consists of four microservices:
1. **Document Ingestion Service** - Chunks documents for processing
2. **Vector Database Service** - Stores embeddings and performs similarity search
3. **LLM Inference Service** - Generates responses using a language model
4. **API Gateway Service** - Orchestrates the RAG pipeline
## Quick Start
### Prerequisites
- Docker and Docker Compose
- Kubernetes (minikube for local development)
- kubectl
- At least 8GB RAM available
- jq (for testing scripts)
### Local Deployment with Docker Compose
```bash
# Build and deploy
make deploy-local
# Test the system
make test-local
# View logs
make logs-local
# Cleanup
make clean
Kubernetes Deployment
# Start minikube
minikube start --cpus=4 --memory=8192
# Enable ingress
minikube addons enable ingress
minikube addons enable metrics-server
# Deploy to Kubernetes
make deploy-k8s
# In a separate terminal, run:
minikube tunnel
# Add to /etc/hosts:
echo "127.0.0.1 rag-chatbot.local" | sudo tee -a /etc/hosts
# Test the system
make test-k8s
# View logs
make logs-k8s
API Usage
Ingest a Document
curl -X POST http://localhost:5000/ingest \
-H "Content-Type: application/json" \
-d '{
"text": "Your document text here...",
"metadata": {"source": "example"}
}'
Query the System
curl -X POST http://localhost:5000/query \
-H "Content-Type: application/json" \
-d '{
"query": "Your question here?",
"top_k": 3,
"max_length": 200
}'
Health Check
curl http://localhost:5000/health
Monitoring
All services expose Prometheus metrics at /metrics:
- Document Ingestion: http://localhost:5001/metrics
- Vector Database: http://localhost:5002/metrics
- LLM Inference: http://localhost:5003/metrics
- API Gateway: http://localhost:5000/metrics
Configuration
Configuration is managed through environment variables and ConfigMaps:
CHUNK_SIZE- Token size for document chunks (default: 512)EMBEDDING_MODEL- Sentence transformer model (default: all-MiniLM-L6-v2)MODEL_NAME- LLM model (default: TinyLlama/TinyLlama-1.1B-Chat-v1.0)TOP_K_DOCUMENTS- Number of documents to retrieve (default: 3)
Scaling
Horizontal Pod Autoscaling is configured for:
- API Gateway: 2-10 replicas
- Document Ingestion: 2-8 replicas
Security Features
- Non-root containers
- Read-only root filesystems where possible
- Network policies for pod-to-pod communication
- Security contexts with dropped capabilities
- Pod Security Standards enforcement
Development
Running Tests
# Install test dependencies
pip install pytest pytest-cov
# Run unit tests
pytest services/*/tests/
# Run integration tests
pytest tests/integration/
Building Individual Services
# Document Ingestion
docker build -t rag-chatbot/document-ingestion:1.0 ./services/document-ingestion
# Vector Database
docker build -t rag-chatbot/vector-database:1.0 ./services/vector-database
# LLM Inference (CPU)
docker build -f ./services/llm-inference/Dockerfile.cpu -t rag-chatbot/llm-inference:1.0-cpu ./services/llm-inference
# API Gateway
docker build -t rag-chatbot/api-gateway:1.0 ./services/api-gateway
Troubleshooting
Services not starting
# Check pod status
kubectl get pods -n rag-chatbot
# View logs
kubectl logs -n rag-chatbot <pod-name>
# Describe pod for events
kubectl describe pod -n rag-chatbot <pod-name>
LLM Inference taking too long
The LLM service downloads the model on first startup, which can take several minutes. Check logs:
kubectl logs -n rag-chatbot -l app=llm-inference -f
Out of memory errors
Increase resource limits in deployment manifests or allocate more memory to minikube:
minikube delete
minikube start --cpus=4 --memory=16384
License
MIT License
Contributing
Contributions welcome! Please open an issue or submit a pull request.
10. Make Scripts Executable```bash chmod +x scripts/*.sh
Usage Instructions
- Clone and Setup:
git clone <repository>
cd rag-chatbot
- Local Development:
make deploy-local
make test-local
- Kubernetes Deployment:
minikube start --cpus=4 --memory=8192
minikube addons enable ingress
minikube addons enable metrics-server
make deploy-k8s
# In separate terminal:
minikube tunnel
# Add to /etc/hosts:
echo "127.0.0.1 rag-chatbot.local" | sudo tee -a /etc/hosts
make test-k8s