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
| RAG | Fine-tuning | |
|---|---|---|
| Новые знания | ✅ Real-time | ❌ Нужен Retrain |
| Citation | ✅ Точно | ❌ Сложно |
| Cost | Per-query | One-time + inference |
| Quality on style | ❌ Средне | ✅ Хорошо |
| Complexity | Средняя | Высокая |
| Maintenance | Index update | Retrain |
**Правило:**Для 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 docs — docs.ragas.io
- “RAG vs Fine-tuning” — Anthropic guide
- HyDE paper — Gao et al.
- Cohere RAG guides — production patterns
🏋️ Упражнения
🟢 Easy
- Naive RAG: на 10 документах — chunking → vector DB → query.
- Citation: указание источника в ответе в формате
[Source N]. - Сравните стратегии chunking: 500 vs 1000 vs 2000 token.
🟡 Medium
- Multi-query RAG: query → 3 варианта → объединение.
- HyDE: синтетический ответ → embed → search.
- Reranking: cross-encoder с top 20 → top 5.
🔴 Hard
- Production RAG service: FastAPI + Qdrant + Celery (ingestion) + Langfuse (observability).
- RAG evaluation: создайте test set из 100 вопросов-ответов, оцените через RAGAS.
- 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.