Python course Β· Module 11: RAG and Multi-Agent Systems

Production RAG - Enterprise Systems

11 min read
In this lesson8

Your RAG prototype works great on your laptop. Then a hundred people use it at once, the API bill grows by the hour, one question waits ten seconds, and when something breaks nobody knows at which stage. A camp for one night is not the same as a permanent expedition base.

Building production RAG systems requires taking four things into account at once: performance, scalability, monitoring and security. It is like designing an entire ecosystem - every element must work in harmony with the others.

Production RAG Architecture

The diagram shows the layers every request passes through:

1β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
2β”‚                    Production RAG Architecture                       β”‚
3β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€
4β”‚                                                                      β”‚
5β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”              β”‚
6β”‚  β”‚  Client  │───▢│ Load Balancer│───▢│  API Gateway β”‚              β”‚
7β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜              β”‚
8β”‚                                               β”‚                      β”‚
9β”‚                                               β–Ό                      β”‚
10β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”‚
11β”‚  β”‚                     RAG Service Layer                       β”‚    β”‚
12β”‚  β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”           β”‚    β”‚
13β”‚  β”‚  β”‚  Retrieval β”‚  β”‚ Augment    β”‚  β”‚ Generation β”‚           β”‚    β”‚
14β”‚  β”‚  β”‚  Service   β”‚  β”‚ Service    β”‚  β”‚ Service    β”‚           β”‚    β”‚
15β”‚  β”‚  β””β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜           β”‚    β”‚
16β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β”‚
17β”‚           β”‚               β”‚               β”‚                        β”‚
18β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”  β”Œβ”€β”€β”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”€β”€β”                β”‚
19β”‚  β”‚ Vector Databaseβ”‚  β”‚  Cache  β”‚  β”‚  LLM APIs   β”‚                β”‚
20β”‚  β”‚ (Pinecone/     β”‚  β”‚ (Redis) β”‚  β”‚ (OpenAI/    β”‚                β”‚
21β”‚  β”‚  Qdrant)       β”‚  β”‚         β”‚  β”‚  Anthropic) β”‚                β”‚
22β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                β”‚
23β”‚                                                                      β”‚
24β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”‚
25β”‚  β”‚                    Observability Layer                         β”‚ β”‚
26β”‚  β”‚  Prometheus β”‚ Grafana β”‚ Jaeger β”‚ LangSmith β”‚ Weights & Biases β”‚ β”‚
27β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚
28β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

From the top: load balancer and API Gateway, the RAG service layer with its three stages, below that the vector database, the Redis cache and the model APIs, and under everything the observability layer.

FastAPI RAG Service

We write the service in FastAPI, because it supports async and validates data by itself. First the imports and the Pydantic models describing the request and the response:

1from fastapi import FastAPI, HTTPException, BackgroundTasks
2from pydantic import BaseModel
3from typing import Optional
4import asyncio
5import hashlib
6from contextlib import asynccontextmanager
7import uvicorn
8
9# Models
10class QueryRequest(BaseModel):
11    question: str
12    top_k: int = 5
13    filters: Optional[dict] = None
14
15class QueryResponse(BaseModel):
16    answer: str
17    sources: list[dict]
18    latency_ms: float

QueryRequest accepts a question, the number of documents and optional filters. FastAPI rejects a request without the question field by itself.

The service class connects to three services. The clients are asynchronous, so waiting for one service does not block other requests:

1class RAGService:
2    """Production RAG service."""
3
4    def __init__(self):
5        self.vector_store = None
6        self.llm = None
7        self.cache = None
8
9    async def initialize(self):
10        """Initialize connections."""
11        # Vector store
12        from qdrant_client import AsyncQdrantClient
13        self.vector_store = AsyncQdrantClient(host="qdrant", port=6333)
14
15        # Cache
16        import redis.asyncio as redis
17        self.cache = redis.Redis(host="redis", port=6379)
18
19        # LLM
20        from openai import AsyncOpenAI
21        self.llm = AsyncOpenAI()

We create the connections in initialize, because they need await. The host names qdrant and redis are the service names from the Docker Compose file you will see at the end of the lesson.

The heart of the service is the query method. It checks the cache first, and only on a miss runs the full pipeline:

1    async def query(self, request: QueryRequest) -> QueryResponse:
2        """Processes a RAG query."""
3        import time
4        start = time.time()
5
6        # 1. Check cache
7        cache_key = f"rag:{hashlib.sha256(request.question.encode()).hexdigest()}"
8        cached = await self.cache.get(cache_key)
9        if cached:
10            return QueryResponse.model_validate_json(cached)
11
12        # 2. Embedding
13        embedding = await self._get_embedding(request.question)
14
15        # 3. Retrieval
16        docs = await self._retrieve(embedding, request.top_k, request.filters)
17
18        # 4. Generation
19        answer = await self._generate(request.question, docs)
20
21        # 5. Response
22        response = QueryResponse(
23            answer=answer,
24            sources=[{"text": d.payload["text"], "score": d.score} for d in docs],
25            latency_ms=(time.time() - start) * 1000
26        )
27
28        # 6. Cache
29        await self.cache.setex(cache_key, 3600, response.model_dump_json())
30
31        return response

We compute the cache key with SHA-256, not the built-in hash(). This fix matters in production: hash() for strings is randomized at every process start, so two workers would compute different keys for the same question. setex stores the answer for 3600 seconds, one hour.

Three helper methods handle embedding, retrieval and generation:

1    async def _get_embedding(self, text: str) -> list[float]:
2        """Generates an embedding."""
3        response = await self.llm.embeddings.create(
4            model="text-embedding-3-small",
5            input=text
6        )
7        return response.data[0].embedding
8
9    async def _retrieve(self, embedding, top_k, filters):
10        """Retrieves documents."""
11        response = await self.vector_store.query_points(
12            collection_name="documents",
13            query=embedding,
14            limit=top_k,
15            query_filter=filters
16        )
17        return response.points
18
19    async def _generate(self, question: str, docs) -> str:
20        """Generates an answer."""
21        context = "\n".join([d.payload["text"] for d in docs])
22
23        response = await self.llm.chat.completions.create(
24            model="gpt-4o-mini",
25            messages=[
26                {"role": "system", "content": f"Context:\n{context}"},
27                {"role": "user", "content": question}
28            ]
29        )
30        return response.choices[0].message.content

Retrieval uses query_points, because the search method is gone from new versions of qdrant-client. The document text lives in the payload, which is why we build sources from d.payload["text"].

Finally the FastAPI app and the endpoint:

1# FastAPI app
2rag_service = RAGService()
3
4@asynccontextmanager
5async def lifespan(app: FastAPI):
6    await rag_service.initialize()
7    yield
8
9app = FastAPI(title="RAG API", lifespan=lifespan)
10
11@app.post("/query", response_model=QueryResponse)
12async def query(request: QueryRequest):
13    try:
14        return await rag_service.query(request)
15    except Exception as e:
16        raise HTTPException(status_code=500, detail=str(e))

lifespan runs initialize once, at application start. Every error becomes a 500 code, but in production do not send str(e) back to the client.

Monitoring and Observability

Without measurements you learn about outages from your users. Prometheus collects three types of metrics: a counter, a histogram and a gauge:

1from prometheus_client import Counter, Histogram, Gauge
2import logging
3import structlog
4
5# Prometheus metrics
6QUERY_COUNT = Counter(
7    "rag_queries_total",
8    "Total number of RAG queries",
9    ["status"]
10)
11
12QUERY_LATENCY = Histogram(
13    "rag_query_latency_seconds",
14    "Query latency in seconds",
15    buckets=[0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0]
16)
17
18RETRIEVAL_SCORE = Gauge(
19    "rag_retrieval_score",
20    "Average retrieval score"
21)
22
23# Structured logging
24structlog.configure(
25    processors=[
26        structlog.stdlib.add_log_level,
27        structlog.processors.TimeStamper(fmt="iso"),
28        structlog.processors.JSONRenderer()
29    ]
30)
31logger = structlog.get_logger()

A Counter only goes up, and the status label separates successes from errors. A Histogram sorts response times into buckets, and structlog writes logs as JSON.

The monitored service wraps processing in a timer:

1class MonitoredRAGService:
2    """RAG with monitoring."""
3
4    async def query(self, request: QueryRequest) -> QueryResponse:
5        with QUERY_LATENCY.time():
6            try:
7                response = await self._process_query(request)
8                QUERY_COUNT.labels(status="success").inc()
9
10                # Log
11                logger.info(
12                    "query_completed",
13                    question=request.question[:50],
14                    latency_ms=response.latency_ms,
15                    num_sources=len(response.sources)
16                )
17
18                return response
19            except Exception as e:
20                QUERY_COUNT.labels(status="error").inc()
21                logger.error("query_failed", error=str(e))
22                raise

QUERY_LATENCY.time() measures the duration of the with block, and an exception increments the error counter and propagates thanks to raise. _process_query stands for the logic of the service's query method.

LangChain Chains

In the production ecosystem you will often meet LangChain. Its basic pattern is a chain: a prompt connected to a model with the | operator and called with the invoke method:

1from langchain_core.prompts import ChatPromptTemplate
2from langchain_openai import ChatOpenAI
3
4prompt = ChatPromptTemplate.from_template("Answer the question briefly: {input}")
5chain = prompt | ChatOpenAI(model="gpt-4o-mini")
6
7user_query = "What is a cache in a RAG system?"
8result = chain.invoke({"input": user_query})
9print(result.content)

invoke takes a dictionary whose keys match the variables in the template, here {input}. The result is a model message, and the text is in the content field.

LangSmith Tracing

LangSmith is the LangChain team's platform for tracing calls. The @traceable decorator records the whole function, and trace creates nested stages inside it:

1from langsmith import traceable, trace
2from langsmith.wrappers import wrap_openai
3from openai import OpenAI
4
5# Every call through this client reaches LangSmith as a separate stage
6openai_client = wrap_openai(OpenAI())
7
8@traceable(name="RAG Query")
9async def traced_query(question: str) -> str:
10    """RAG query with tracing."""
11
12    # Embedding
13    with trace("embedding", metadata={"model": "text-embedding-3-small"}):
14        embedding = get_embedding(question)
15
16    # Retrieval
17    with trace("retrieval") as span:
18        docs = retrieve(embedding)
19        span.add_metadata({"num_docs": len(docs)})
20
21    # Generation
22    with trace("generation", metadata={"model": "gpt-4o-mini"}):
23        response = generate(question, docs)
24
25    return response

Note: trace is a separate function from the langsmith package, not a method of the Client object. get_embedding, retrieve and generate are functions from earlier lessons.

Evaluation runs the system on a set of questions with expected answers. The traced_query function is asynchronous, so we use aevaluate:

1from langsmith import aevaluate
2
3async def target(inputs: dict) -> dict:
4    """Adapter: LangSmith passes a dict with the question and expects a dict with the answer."""
5    return {"answer": await traced_query(inputs["question"])}
6
7def contains_expected(outputs: dict, reference_outputs: dict) -> bool:
8    """Simple evaluator: does the answer contain the expected phrase?"""
9    return reference_outputs["answer"].lower() in outputs["answer"].lower()
10
11async def evaluate_rag():
12    """Evaluates the RAG system."""
13    # Dataset with questions and expected answers, created earlier in LangSmith
14    results = await aevaluate(
15        target,
16        data="rag-eval-dataset",
17        evaluators=[contains_expected]
18    )
19    return results

An evaluator is an ordinary function that returns True or False. Metrics such as faithfulness or context precision are computed the same way, only the judge is a language model.

Rate Limiting and Security

A public API without limits is an invitation to abuse. First a simple limiter with a one-minute window:

1from fastapi import Depends, HTTPException
2from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
3import time
4
5# Rate limiter
6class RateLimiter:
7    """Simple rate limiter."""
8
9    def __init__(self, requests_per_minute: int = 60):
10        self.rpm = requests_per_minute
11        self.requests: dict[str, list[float]] = {}
12
13    async def check(self, user_id: str) -> bool:
14        now = time.time()
15        minute_ago = now - 60
16
17        if user_id not in self.requests:
18            self.requests[user_id] = []
19
20        # Remove old requests
21        self.requests[user_id] = [
22            t for t in self.requests[user_id] if t > minute_ago
23        ]
24
25        if len(self.requests[user_id]) >= self.rpm:
26            return False
27
28        self.requests[user_id].append(now)
29        return True
30
31rate_limiter = RateLimiter()
32security = HTTPBearer()

The limiter remembers the request times of every user and removes those older than 60 seconds. It keeps data in process memory, so with several service instances move it to Redis.

Now we connect authorization and the limiter to the endpoint through Depends:

1async def verify_token(credentials: HTTPAuthorizationCredentials = Depends(security)):
2    """Token verification."""
3    # In production: JWT verification
4    if credentials.credentials == "invalid":
5        raise HTTPException(status_code=401, detail="Invalid token")
6    return credentials.credentials
7
8async def rate_limit(user_id: str = Depends(verify_token)):
9    """Rate limiting middleware."""
10    if not await rate_limiter.check(user_id):
11        raise HTTPException(status_code=429, detail="Too many requests")
12    return user_id
13
14@app.post("/query")
15async def query(request: QueryRequest, user: str = Depends(rate_limit)):
16    return await rag_service.query(request)

HTTPBearer extracts the token from the Authorization header, and verify_token returns 401 for a bad token. Exceeding the limit ends with a 429 code. This endpoint replaces the earlier /query, so keep only one in the app, because FastAPI uses the first matching route.

Docker Deployment

Finally we package the service in a container. The Dockerfile describes the application image:

1# Dockerfile
2FROM python:3.11-slim
3
4WORKDIR /app
5
6# Install dependencies
7COPY requirements.txt .
8RUN pip install --no-cache-dir -r requirements.txt
9
10# Copy code
11COPY . .
12
13# Run
14CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]

We copy requirements.txt first and the code afterwards, so a code change does not force reinstalling dependencies.

Docker Compose starts the whole base at once: the API, Qdrant and Redis:

1# docker-compose.yml
2version: '3.8'
3
4services:
5  rag-api:
6    build: .
7    ports:
8      - "8000:8000"
9    environment:
10      - OPENAI_API_KEY=${OPENAI_API_KEY}
11      - QDRANT_HOST=qdrant
12      - REDIS_HOST=redis
13    depends_on:
14      - qdrant
15      - redis
16
17  qdrant:
18    image: qdrant/qdrant:latest
19    ports:
20      - "6333:6333"
21    volumes:
22      - qdrant_data:/qdrant/storage
23
24  redis:
25    image: redis:alpine
26    ports:
27      - "6379:6379"

Compose replaces ${OPENAI_API_KEY} with the value of the host environment variable, so the key never lands in the repository. The version key is obsolete and ignored in current Compose.

The rest of the file adds the observability layer:

1  prometheus:
2    image: prom/prometheus
3    ports:
4      - "9090:9090"
5    volumes:
6      - ./prometheus.yml:/etc/prometheus/prometheus.yml
7
8  grafana:
9    image: grafana/grafana
10    ports:
11      - "3000:3000"
12
13volumes:
14  qdrant_data:

Prometheus collects metrics, and Grafana draws charts from them. The qdrant_data volume makes the vectors survive a container restart.

Best Practices

A list of good practices sums up all the lessons of this module:

1"""
2Production RAG Best Practices:
3
41. RETRIEVAL
5   - Use hybrid search (vector + keyword)
6   - Implement reranking
7   - Filter by metadata
8
92. CHUNKING
10   - Experiment with chunk sizes
11   - Use overlap
12   - Consider semantic chunking
13
143. CACHING
15   - Cache embeddings
16   - Cache frequent queries
17   - Invalidate on document updates
18
194. MONITORING
20   - Track latency at each stage
21   - Monitor retrieval quality
22   - Alert on anomalies
23
245. SECURITY
25   - Rate limiting
26   - Input sanitization
27   - PII detection and filtering
28
296. SCALABILITY
30   - Horizontal scaling API
31   - Sharding vector database
32   - Async processing
33"""

I recommend starting with caching and monitoring, because they pay off the fastest: a cache reduces API costs and latency, and metrics show where to look for a problem. Notice that hardcoded prompts are not on the list: in production, prompts live in configuration or in a versioning system.

Congratulations! You have learned advanced RAG and Multi-Agent systems. In the next lesson you will learn LangGraph - a framework for building complex agent workflows!

Remember: production RAG is a permanent expedition base where every tent has its job and a guard counts everyone who walks in.

Spotted a mistake in this lesson?

Check yourself

Answer the questions from this lesson. Pick an answer to see right away whether it is correct.

  1. 1. What is crucial in a production RAG system?

  2. 2. Why is cache important in production RAG?

These are 2 of 3 questions for this lesson. Solve the rest in the game.

Hands-on tasks in the game

  • Code editor

    Implement a production RAG endpoint

  • Vertical ordering

    Arrange the layers of production RAG architecture:

  • Code editor

    Implement a complete Simple RAG

  • Click in order

    Arrange the LangChain chain invocation:

  • Code editor

    Implement a Semantic Search Engine

  • Code editor

    Implement a Multi-Agent Content Pipeline

Useful articles