Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

RAG Pipeline

🎯 Цель

После прочтения этой главы:

  • Знаете полную архитектуру RAG (Retrieval Augmented Generation)
  • Можете строить production-grade RAG pipeline
  • Понимаете стратегии Chunking и trade-off
  • Можете применять advanced RAG техники (HyDE, multi-query, re-ranking)
  • Знаете измерение и улучшение качества RAG

Что нужно изучить

  • Архитектура RAG — Naive, Advanced, Modular
  • Стратегии Chunking — fixed, semantic, sliding window, recursive
  • Стратегии Retrieval — dense, sparse, hybrid, multi-query
  • Reranking — Cross-encoder, LLM-based
  • HyDE(Hypothetical Document Embeddings)
  • Citation и source attribution
  • Context window management
  • RAG evaluation — RAGAS, custom metrics

Что такое RAG и зачем?

Проблема

LLM hallucination — может давать неправильную информацию:

  • Training data старые (до 2024 года)
  • Не знает ваши личные документы
  • Неверный ответ на точные факты

Решение — RAG

1. Пользователь задаёт вопрос: "Какова политика нашей компании?"
2. Retrieval: получение 5 похожих chunk из vector DB
3. Augment: добавление chunks в prompt
4. Generate: LLM отвечает на основе контекста
5. Cite: указание из какого chunk взято

RAG vs Fine-tuning

RAGFine-tuning
Новые знания✅ Real-time❌ Нужен Retrain
Citation✅ Точно❌ Сложно
CostPer-queryOne-time + inference
Quality on style❌ Средне✅ Хорошо
ComplexityСредняяВысокая
MaintenanceIndex updateRetrain

**Правило:**Для knowledge — RAG, для behavior/style — fine-tuning.

Архитектура RAG

Naive RAG

Query → Embed → Vector DB Search → Top-K chunks → LLM prompt → Answer

Проблемы:

  • Плохой retrieval → плохой ответ
  • Противоречия в chunks контекста
  • LLM делает hallucination вне контекста

Advanced RAG (modern)

Query
  ↓
Query Transformation:
  - Multi-query (3 ta variant)
  - HyDE (sintetik javob → embed)
  - Step-back (umumiyroq savol)
  ↓
Hybrid Retrieval:
  - Dense (semantic)
  - Sparse (BM25)
  - Metadata filter
  ↓
Reranking (Cross-encoder)
  ↓
Context Construction:
  - Deduplication
  - Sort by relevance
  - Compress (LLM summary)
  ↓
LLM Generation:
  - Structured prompt
  - Citation markers
  ↓
Post-processing:
  - Source attribution
  - Confidence score

Примеры кода

Production RAG pipeline

from dataclasses import dataclass
from openai import AsyncOpenAI
from anthropic import AsyncAnthropic
from qdrant_client import AsyncQdrantClient
from sentence_transformers import CrossEncoder

@dataclass
class RetrievedChunk:
    text: str
    source: str
    page: int
    score: float

@dataclass
class RAGAnswer:
    answer: str
    sources: list[RetrievedChunk]
    confidence: float

class RAGPipeline:
    def __init__(self):
        self.openai = AsyncOpenAI()
        self.anthropic = AsyncAnthropic()
        self.qdrant = AsyncQdrantClient(url="http://localhost:6333")
        self.reranker = CrossEncoder("BAAI/bge-reranker-base")
        self.collection = "docs"
    
    async def embed(self, text: str) -> list[float]:
        response = await self.openai.embeddings.create(
            model="text-embedding-3-small",
            input=[text],
        )
        return response.data[0].embedding
    
    async def retrieve(self, query: str, top_k: int = 20) -> list[RetrievedChunk]:
        embedding = await self.embed(query)
        results = await self.qdrant.search(
            collection_name=self.collection,
            query_vector=embedding,
            limit=top_k,
        )
        return [
            RetrievedChunk(
                text=r.payload["text"],
                source=r.payload.get("source", ""),
                page=r.payload.get("page", 0),
                score=r.score,
            )
            for r in results
        ]
    
    def rerank(self, query: str, chunks: list[RetrievedChunk], top_k: int = 5):
        pairs = [(query, c.text) for c in chunks]
        scores = self.reranker.predict(pairs)
        ranked = sorted(zip(scores, chunks), key=lambda x: -x[0])
        # Сохранить новый score
        for new_score, chunk in ranked[:top_k]:
            chunk.score = float(new_score)
        return [c for _, c in ranked[:top_k]]
    
    def build_prompt(self, query: str, chunks: list[RetrievedChunk]) -> str:
        context = "\n\n".join([
            f"[Source {i+1}: {c.source}, page {c.page}]\n{c.text}"
            for i, c in enumerate(chunks)
        ])
        
        return f"""Ты опытный assistant. Ответь на вопрос точно на основе следующего контекста.

ПРАВИЛА:
1. Отвечай ТОЛЬКО на основе данного контекста
2. Если ответа в контексте нет, ответь "В предоставленных данных ответ не найден"
3. Для каждого факта указывай ссылку в формате [Source N]
4. Отвечай на узбекском языке

КОНТЕКСТ:
{context}ВОПРОС: {query}ОТВЕТ:"""
    
    async def generate(self, prompt: str) -> tuple[str, float]:
        response = await self.anthropic.messages.create(
            model="claude-sonnet-4-6",
            max_tokens=1024,
            messages=[{"role": "user", "content": prompt}],
        )
        text = response.content[0].text
        # Confidence estimation (simple heuristic)
        confidence = 0.9 if "[Source" in text else 0.3
        return text, confidence
    
    async def query(self, query: str) -> RAGAnswer:
        # 1. Retrieve
        chunks = await self.retrieve(query, top_k=20)
        
        # 2. Rerank
        top_chunks = self.rerank(query, chunks, top_k=5)
        
        # 3. Build prompt
        prompt = self.build_prompt(query, top_chunks)
        
        # 4. Generate
        answer, confidence = await self.generate(prompt)
        
        return RAGAnswer(
            answer=answer,
            sources=top_chunks,
            confidence=confidence,
        )

# Usage
rag = RAGPipeline()
result = await rag.query("Какие у нас часы работы?")
print(result.answer)
for src in result.sources:
    print(f"  - {src.source} (p.{src.page}): {src.score:.3f}")

Multi-query — разделение вопроса на 3 варианта

async def multi_query_search(query: str, top_k: int = 5):
    """Один query → 3 варианта → объединённый результат."""
    
    # 1. Generate query variants
    variant_prompt = f"""Перепишите следующий вопрос 3 разными способами:

Вопрос: {query}Варианты (каждый на новой строке):
1.
2.
3."""
    
    response = await openai.chat.completions.create(
        model="gpt-4o-mini",
        messages=[{"role": "user", "content": variant_prompt}],
    )
    variants = response.choices[0].message.content.strip().split("\n")
    variants = [v.split(". ", 1)[1] for v in variants if ". " in v]
    
    # 2. Retrieve for each
    all_chunks = []
    for q in [query] + variants:
        chunks = await retrieve(q, top_k=top_k)
        all_chunks.extend(chunks)
    
    # 3. Deduplicate (по id или content hash)
    seen = set()
    unique = []
    for c in all_chunks:
        key = hash(c.text[:100])
        if key not in seen:
            seen.add(key)
            unique.append(c)
    
    return unique

HyDE — Hypothetical Document Embeddings

async def hyde_search(query: str, top_k: int = 5):
    """Не прямой search из query, а создание синтетического 'ответа' и его embed."""
    
    # 1. Создать синтетический ответ
    hypothesis = await openai.chat.completions.create(
        model="gpt-4o-mini",
        messages=[{"role": "user", "content": 
            f"Напишите полный, подробный ответ на следующий вопрос (даже если не факт):\n{query}"}],
    )
    hypothetical_answer = hypothesis.choices[0].message.content
    
    # 2. Embed гипотетический ответ
    embedding = await openai.embeddings.create(
        model="text-embedding-3-small",
        input=[hypothetical_answer],
    )
    
    # 3. Search этим embedding (ответ → ответ similarity!)
    results = await qdrant.search(
        collection_name="docs",
        query_vector=embedding.data[0].embedding,
        limit=top_k,
    )
    
    return results

Smart стратегии chunking

from langchain.text_splitter import RecursiveCharacterTextSplitter

# Strategy 1: Fixed-size (самый простой)
fixed = RecursiveCharacterTextSplitter(
    chunk_size=1000,
    chunk_overlap=200,
)

# Strategy 2: Markdown-aware
from langchain.text_splitter import MarkdownHeaderTextSplitter

md_splitter = MarkdownHeaderTextSplitter(headers_to_split_on=[
    ("#", "Header 1"), ("##", "Header 2"), ("###", "Header 3"),
])

# Strategy 3: Semantic (LangChain experimental)
from langchain_experimental.text_splitter import SemanticChunker
from langchain_openai import OpenAIEmbeddings

semantic = SemanticChunker(
    OpenAIEmbeddings(model="text-embedding-3-small"),
    breakpoint_threshold_type="percentile",
)

# Strategy 4: Sliding window (overlap)
def sliding_window_chunks(text: str, window: int = 500, stride: int = 250):
    chunks = []
    for i in range(0, len(text) - window + 1, stride):
        chunks.append(text[i:i + window])
    return chunks

Context window management

def build_context_within_budget(
    chunks: list[RetrievedChunk],
    max_tokens: int = 8000,
    encoder=tiktoken.encoding_for_model("gpt-4o"),
) -> list[RetrievedChunk]:
    """Вернуть только chunks, которые помещаются в budget."""
    included = []
    total = 0
    
    for chunk in chunks:  # already sorted by relevance
        tokens = len(encoder.encode(chunk.text))
        if total + tokens > max_tokens:
            break
        included.append(chunk)
        total += tokens
    
    return included

RAG evaluation — RAGAS

# pip install ragas

from ragas import evaluate
from ragas.metrics import (
    faithfulness,
    answer_relevancy,
    context_precision,
    context_recall,
)
from datasets import Dataset

# Test set
data = {
    "question": ["Какие часы работы?", "Где адрес?"],
    "answer": ["с 8:00 до 18:00", "Ташкент, Юнусабад"],
    "contexts": [
        ["Наши часы работы с понедельника по пятницу 8:00-18:00"],
        ["Office: Ташкент, Юнусабадский район"],
    ],
    "ground_truth": ["8:00-18:00", "Ташкент, Юнусабад"],
}

dataset = Dataset.from_dict(data)
result = evaluate(
    dataset,
    metrics=[faithfulness, answer_relevancy, context_precision, context_recall],
)
print(result)
# {faithfulness: 0.95, answer_relevancy: 0.88, ...}

Интеграция с backend

Production RAG FastAPI endpoint

from fastapi import FastAPI
from contextlib import asynccontextmanager

@asynccontextmanager
async def lifespan(app):
    app.state.rag = RAGPipeline()
    yield

app = FastAPI(lifespan=lifespan)

class RAGRequest(BaseModel):
    query: str
    session_id: str = None
    top_k: int = 5
    rerank: bool = True
    multi_query: bool = False

class RAGResponse(BaseModel):
    answer: str
    sources: list[dict]
    confidence: float
    latency_ms: int

@app.post("/rag/query", response_model=RAGResponse)
async def rag_query(req: RAGRequest):
    start = time.time()
    
    result = await app.state.rag.query(req.query)
    
    # Log for monitoring
    await log_query(
        query=req.query,
        answer=result.answer,
        sources=[s.source for s in result.sources],
        confidence=result.confidence,
        session_id=req.session_id,
    )
    
    return RAGResponse(
        answer=result.answer,
        sources=[
            {"text": s.text[:200], "source": s.source, "page": s.page, "score": s.score}
            for s in result.sources
        ],
        confidence=result.confidence,
        latency_ms=int((time.time() - start) * 1000),
    )

Streaming RAG answer (SSE)

@app.post("/rag/stream")
async def rag_stream(req: RAGRequest):
    # 1. Retrieve (non-streaming)
    chunks = await app.state.rag.retrieve(req.query)
    top_chunks = app.state.rag.rerank(req.query, chunks)
    prompt = app.state.rag.build_prompt(req.query, top_chunks)
    
    async def event_stream():
        # Send sources first
        sources = [{"source": c.source, "score": c.score} for c in top_chunks]
        yield f"data: {json.dumps({'type': 'sources', 'data': sources})}\n\n"
        
        # Stream LLM response
        async with anthropic.messages.stream(
            model="claude-sonnet-4-6",
            max_tokens=1024,
            messages=[{"role": "user", "content": prompt}],
        ) as stream:
            async for text in stream.text_stream:
                yield f"data: {json.dumps({'type': 'token', 'text': text})}\n\n"
        
        yield "data: [DONE]\n\n"
    
    return StreamingResponse(event_stream(), media_type="text/event-stream")

Ресурсы

  • “Advanced RAG Techniques” — IVAN Ilin (Medium series)
  • LlamaIndex Advanced RAG cookbook
  • RAGAS docsdocs.ragas.io
  • “RAG vs Fine-tuning” — Anthropic guide
  • HyDE paper — Gao et al.
  • Cohere RAG guides — production patterns

🏋️ Упражнения

🟢 Easy

  1. Naive RAG: на 10 документах — chunking → vector DB → query.
  2. Citation: указание источника в ответе в формате [Source N].
  3. Сравните стратегии chunking: 500 vs 1000 vs 2000 token.

🟡 Medium

  1. Multi-query RAG: query → 3 варианта → объединение.
  2. HyDE: синтетический ответ → embed → search.
  3. Reranking: cross-encoder с top 20 → top 5.

🔴 Hard

  1. Production RAG service: FastAPI + Qdrant + Celery (ingestion) + Langfuse (observability).
  2. RAG evaluation: создайте test set из 100 вопросов-ответов, оцените через RAGAS.
  3. Domain-specific tuning: специальный RAG для узбекских правовых документов (chunking, prompts).

Capstone

notebooks/month-05/06_rag_pipeline.ipynb:

  • **Проект:**RAG chatbot для Конституции Узбекистана или УК
  • Ingestion 100+ документов
  • Multi-query + HyDE + reranking
  • Citation
  • Streamlit UI
  • RAGAS evaluation

✅ Чек-лист

  • Знаю архитектуру RAG
  • Могу применять стратегии chunking (fixed, semantic)
  • Hybrid retrieval (dense + sparse)
  • Reranking (cross-encoder)
  • HyDE и Multi-query
  • Citation и source attribution
  • Streaming RAG
  • RAG evaluation (RAGAS)

Переходим к AI Agents.