Kurs Python · Moduł 11: RAG i systemy wieloagentowe

Production RAG - Enterprise Systems

10 min czytania
W tej lekcji8

Twój prototyp RAG działa świetnie na laptopie. Potem korzysta z niego sto osób naraz, rachunek za API rośnie z godziny na godzinę, jedno pytanie czeka dziesięć sekund, a gdy coś się psuje, nikt nie wie, na którym etapie. Obóz na jedną noc to nie to samo co stała baza wyprawy.

Budowanie produkcyjnych systemów RAG wymaga uwzględnienia czterech rzeczy naraz: wydajności, skalowalności, monitoringu i bezpieczeństwa. To jak projektowanie całego ekosystemu - każdy element musi współgrać z innymi.

Architektura produkcyjnego RAG

Schemat pokazuje warstwy, przez które przechodzi każde zapytanie:

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└─────────────────────────────────────────────────────────────────────┘

Od góry: load balancer i API Gateway, warstwa serwisu RAG z trzema etapami, niżej baza wektorowa, cache Redis i API modeli, a pod wszystkim warstwa obserwowalności.

FastAPI RAG Service

Serwis piszemy w FastAPI, bo obsługuje async i sam waliduje dane. Najpierw importy i modele Pydantic opisujące żądanie i odpowiedź:

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# Modele
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 przyjmuje pytanie, liczbę dokumentów i opcjonalne filtry. Żądanie bez pola question FastAPI odrzuci samo.

Klasa serwisu łączy się z trzema usługami. Klienty są asynchroniczne, więc czekanie na jedną usługę nie blokuje innych żądań:

1class RAGService:
2    """Produkcyjny serwis RAG."""
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        """Inicjalizacja połączeń."""
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()

Połączenia tworzymy w initialize, bo wymagają await. Nazwy hostów qdrant i redis to nazwy usług z pliku Docker Compose, który zobaczysz na końcu lekcji.

Serce serwisu to metoda query. Najpierw sprawdza cache, a dopiero przy chybieniu wykonuje pełny pipeline:

1    async def query(self, request: QueryRequest) -> QueryResponse:
2        """Przetwarza zapytanie RAG."""
3        import time
4        start = time.time()
5
6        # 1. Sprawdź 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

Klucz cache liczymy przez SHA-256, a nie wbudowane hash(). To poprawka ważna w produkcji: hash() dla stringów jest losowany przy każdym starcie procesu, więc dwa workery liczyłyby różne klucze dla tego samego pytania. setex zapisuje odpowiedź na 3600 sekund, czyli godzinę.

Trzy metody pomocnicze robią embedding, wyszukiwanie i generowanie:

1    async def _get_embedding(self, text: str) -> list[float]:
2        """Generuje 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        """Pobiera dokumenty."""
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        """Generuje odpowiedź."""
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"Kontekst:\n{context}"},
27                {"role": "user", "content": question}
28            ]
29        )
30        return response.choices[0].message.content

Wyszukiwanie używa query_points, bo metoda search zniknęła z nowych wersji qdrant-client. Tekst dokumentu leży w payload, dlatego źródła budujemy z d.payload["text"].

Na końcu aplikacja FastAPI i 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 uruchamia initialize raz, przy starcie aplikacji. Każdy błąd staje się kodem 500, ale w produkcji nie odsyłaj klientowi str(e).

Monitoring i Observability

Bez pomiarów o awarii dowiesz się od użytkowników. Prometheus zbiera trzy typy metryk: licznik, histogram i wskaźnik:

1from prometheus_client import Counter, Histogram, Gauge
2import logging
3import structlog
4
5# Metryki Prometheus
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()

Counter tylko rośnie, a etykieta status rozdziela sukcesy od błędów. Histogram układa czasy odpowiedzi w przedziały, a structlog wypisuje logi jako JSON.

Serwis z monitoringiem owija przetwarzanie pomiarem czasu:

1class MonitoredRAGService:
2    """RAG z monitoringiem."""
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() mierzy czas bloku with, a wyjątek zwiększa licznik błędów i leci dalej dzięki raise. _process_query oznacza tu logikę z metody query serwisu.

Łańcuchy LangChain

W ekosystemie produkcyjnym często spotkasz LangChain. Jego podstawowy wzorzec to łańcuch: prompt połączony z modelem operatorem | i wywołany metodą invoke:

1from langchain_core.prompts import ChatPromptTemplate
2from langchain_openai import ChatOpenAI
3
4prompt = ChatPromptTemplate.from_template("Odpowiedz krótko na pytanie: {input}")
5chain = prompt | ChatOpenAI(model="gpt-4o-mini")
6
7user_query = "Czym jest cache w systemie RAG?"
8result = chain.invoke({"input": user_query})
9print(result.content)

invoke przyjmuje słownik, którego klucze pasują do zmiennych w szablonie, tu {input}. Wynik to wiadomość modelu, a tekst leży w polu content.

LangSmith Tracing

LangSmith to platforma twórców LangChain do śledzenia wywołań. Dekorator @traceable zapisuje całą funkcję, a trace tworzy w niej zagnieżdżone etapy:

1from langsmith import traceable, trace
2from langsmith.wrappers import wrap_openai
3from openai import OpenAI
4
5# Każde wywołanie przez ten klient trafia do LangSmith jako osobny etap
6openai_client = wrap_openai(OpenAI())
7
8@traceable(name="RAG Query")
9async def traced_query(question: str) -> str:
10    """Zapytanie RAG z tracingiem."""
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

Uwaga: trace to osobna funkcja z pakietu langsmith, a nie metoda obiektu Client. get_embedding, retrieve i generate to funkcje z wcześniejszych lekcji.

Ewaluacja uruchamia system na zbiorze pytań z oczekiwanymi odpowiedziami. Funkcja traced_query jest asynchroniczna, więc używamy aevaluate:

1from langsmith import aevaluate
2
3async def target(inputs: dict) -> dict:
4    """Adapter: LangSmith podaje słownik z pytaniem i oczekuje słownika z odpowiedzią."""
5    return {"answer": await traced_query(inputs["question"])}
6
7def contains_expected(outputs: dict, reference_outputs: dict) -> bool:
8    """Prosty evaluator: czy odpowiedź zawiera oczekiwaną frazę?"""
9    return reference_outputs["answer"].lower() in outputs["answer"].lower()
10
11async def evaluate_rag():
12    """Ewaluacja systemu RAG."""
13    # Dataset z pytaniami i oczekiwanymi odpowiedziami, utworzony wcześniej w LangSmith
14    results = await aevaluate(
15        target,
16        data="rag-eval-dataset",
17        evaluators=[contains_expected]
18    )
19    return results

Evaluator to zwykła funkcja zwracająca True albo False. Metryki takie jak faithfulness czy context precision liczy się podobnie, tylko sędzią jest model językowy.

Rate Limiting i Bezpieczeństwo

Publiczne API bez limitów to zaproszenie do nadużyć. Najpierw prosty limiter z oknem jednej minuty:

1from fastapi import Depends, HTTPException
2from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
3import time
4
5# Rate limiter
6class RateLimiter:
7    """Prosty 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        # Usuń stare requesty
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()

Limiter pamięta czasy żądań każdego użytkownika i usuwa te starsze niż 60 sekund. Trzyma dane w pamięci procesu, więc przy kilku instancjach serwisu przenieś go do Redisa.

Teraz podłączamy autoryzację i limiter do endpointu przez Depends:

1async def verify_token(credentials: HTTPAuthorizationCredentials = Depends(security)):
2    """Weryfikacja tokenu."""
3    # W produkcji: 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 wyciąga token z nagłówka Authorization, a verify_token zwraca 401 dla złego tokenu. Przekroczenie limitu kończy się kodem 429. Ten endpoint zastępuje wcześniejszy /query, więc w aplikacji zostaw tylko jeden, bo FastAPI użyje pierwszej pasującej trasy.

Docker Deployment

Na koniec pakujemy serwis w kontener. Dockerfile opisuje obraz aplikacji:

1# Dockerfile
2FROM python:3.11-slim
3
4WORKDIR /app
5
6# Instalacja zależności
7COPY requirements.txt .
8RUN pip install --no-cache-dir -r requirements.txt
9
10# Kopiowanie kodu
11COPY . .
12
13# Uruchomienie
14CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]

Kopiujemy najpierw requirements.txt, a potem kod, więc zmiana kodu nie wymusza ponownej instalacji zależności.

Docker Compose uruchamia całą bazę naraz: API, Qdrant i 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"

Zapis ${OPENAI_API_KEY} Compose podmienia wartością ze zmiennej środowiskowej hosta, więc klucz nie trafia do repozytorium. Klucz version jest w aktualnym Compose przestarzały i ignorowany.

Dalsza część pliku dodaje warstwę obserwowalności:

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 zbiera metryki, a Grafana rysuje z nich wykresy. Wolumen qdrant_data sprawia, że wektory przetrwają restart kontenera.

Best Practices

Wszystkie lekcje z tego modułu zbiera lista dobrych praktyk:

1"""
2Production RAG Best Practices:
3
41. RETRIEVAL
5   - Używaj hybrid search (vector + keyword)
6   - Implementuj reranking
7   - Filtruj po metadanych
8
92. CHUNKING
10   - Eksperymentuj z rozmiarem chunks
11   - Używaj overlap
12   - Rozważ semantic chunking
13
143. CACHING
15   - Cache embeddings
16   - Cache częstych zapytań
17   - Invalidacja przy aktualizacji dokumentów
18
194. MONITORING
20   - Śledź latency na każdym etapie
21   - Monitoruj jakość retrieval
22   - Alerting na anomalie
23
245. SECURITY
25   - Rate limiting
26   - Input sanitization
27   - PII detection i filtering
28
296. SCALABILITY
30   - Horizontal scaling API
31   - Sharding vector database
32   - Async processing
33"""

Polecam zaczynać od cache i monitoringu, bo zwracają się najszybciej: cache redukuje koszty API i opóźnienia, a metryki pokazują, gdzie szukać problemu. Zwróć uwagę, że na liście nie ma zahardkodowanych promptów: w produkcji prompty trzyma się w konfiguracji albo w systemie wersjonowania.

Gratulacje! Poznałeś zaawansowane systemy RAG i Multi-Agent. W następnej lekcji poznasz LangGraph - framework do budowania złożonych przepływów agentów!

Zapamiętaj: produkcyjny RAG to stała baza wyprawy, w której każdy namiot ma swoje zadanie, a strażnik liczy każdego, kto wchodzi.

Widzisz błąd w tej lekcji?

Sprawdź się

Odpowiedz na pytania z tej lekcji. Wybierz odpowiedź, a od razu zobaczysz, czy jest poprawna.

  1. 1. Co jest kluczowe w produkcyjnym systemie RAG?

  2. 2. Dlaczego cache jest ważny w produkcyjnym RAG?

To 2 z 3 pytań do tej lekcji. Pozostałe rozwiążesz w grze.

Zadania praktyczne w grze

  • Edytor kodu

    Zaimplementuj produkcyjny endpoint RAG

  • Układanie w pionie

    Ułóż warstwy architektury produkcyjnego RAG:

  • Edytor kodu

    Zaimplementuj kompletny Simple RAG

  • Klikanie w kolejności

    Ułóż wywołanie łańcucha LangChain:

  • Edytor kodu

    Zaimplementuj Semantic Search Engine

  • Edytor kodu

    Zaimplementuj Multi-Agent Content Pipeline

Przydatne artykuły