Что такое indexing pipeline и зачем он нужен
RAG состоит из двух фаз: offline (indexing) и online (retrieval + generation). Indexing pipeline — это оффлайн-часть: он превращает сырые документы в вектора, готовые для поиска. Эта работа делается заранее, один раз для всего корпуса, и затем повторяется при обновлении данных.
Важно понять: качество retrieval целиком зависит от качества индексирования. Если на этапе парсинга потерялся текст или метаданные — их не восстановить в момент запроса. Если чанки нарезаны неудачно — embedding не поможет. Pipeline — фундамент.
Архитектура: как данные текут через pipeline
Каждый этап преобразует данные определённого типа в данные следующего. Понимание этих типов помогает правильно спроектировать структуры данных и интерфейсы между компонентами.
Ключевой принцип дизайна: каждый этап — чистая функция без побочных эффектов. Parser принимает путь к файлу, возвращает Document. Splitter принимает Document, возвращает список Chunk. Embedder принимает список Chunk, возвращает список EmbeddedChunk. Storage принимает список EmbeddedChunk, пишет в БД. Это делает каждый компонент тестируемым и заменяемым независимо от остальных.
Модели данных: что течёт между этапами
Прежде чем писать код pipeline, нужно определить структуры данных.
Строгая типизация — не бюрократия: она ловит ошибки на этапе разработки,
а не в 3 часа ночи в проде, когда половина документов проиндексирована с пустым source.
# ── Типы данных pipeline ─────────────────────────────────────────────
@dataclass
class Document:
text: str # извлечённый текст
source: str # путь к файлу или URL
title: str # заголовок документа
content_hash: str # MD5 текста — для дедупликации
doc_id: str # hash(source) — стабильный ID
metadata: dict # author, created_at, tags, …
@dataclass
class Chunk:
text: str # текст чанка
doc_id: str # ID родительского документа
chunk_id: str # hash(doc_id + chunk_index)
chunk_index: int # порядковый номер в документе
metadata: dict # унаследовано от Document
@dataclass
class EmbeddedChunk:
chunk_id: str # из Chunk
vector: np.ndarray # shape (dim,), L2-нормализован
text: str # оригинальный текст (для payload)
metadata: dict # из Chunk
Этап 1: Парсинг документов
Парсер — точка входа данных. Его задача: извлечь чистый текст. «Чистый» означает: без артефактов форматирования, колонтитулов, нумерации страниц, HTML-тегов и прочего технического мусора, который ухудшит качество embedding.
doc.core_properties.title, .author.
Сохраняй заголовки разделов (стили Heading1/Heading2) как структурные метаданные.<script>, <style>, <nav>, <footer>.
Извлекай <title>, <meta description>.
trafilatura — готовая библиотека для извлечения основного текста страниц.## заголовки помогают семантическому chunking.
Альтернатива: markdownify для очистки в plain text.import hashlib
from dataclasses import dataclass, field
from datetime import datetime
from pathlib import Path
@dataclass
class Document:
text: str
source: str
title: str = ""
doc_id: str = ""
content_hash: str = ""
metadata: dict = field(default_factory=dict)
def __post_init__(self):
if not self.doc_id:
self.doc_id = hashlib.md5(self.source.encode()).hexdigest()
if not self.content_hash:
self.content_hash = hashlib.md5(self.text.encode()).hexdigest()
def parse_pdf(path: str) -> Document:
import pdfplumber
texts = []
title = Path(path).stem
with pdfplumber.open(path) as pdf:
# Пытаемся вытащить title из метаданных PDF
if pdf.metadata.get("Title"):
title = pdf.metadata["Title"]
for page in pdf.pages:
page_text = page.extract_text(x_tolerance=2, y_tolerance=2)
if page_text:
texts.append(page_text)
return Document(
text="\n\n".join(texts),
source=path,
title=title,
metadata={"pages": len(pdf.pages), "type": "pdf"},
)
def parse_docx(path: str) -> Document:
from docx import Document as DocxDocument
doc = DocxDocument(path)
props = doc.core_properties
paragraphs = [p.text for p in doc.paragraphs if p.text.strip()]
return Document(
text="\n\n".join(paragraphs),
source=path,
title=props.title or Path(path).stem,
metadata={
"author": props.author or "",
"created": str(props.created or ""),
"type": "docx",
},
)
def parse_html(path_or_url: str) -> Document:
from bs4 import BeautifulSoup
if path_or_url.startswith("http"):
import httpx
html = httpx.get(path_or_url, timeout=10).text
source = path_or_url
else:
html = Path(path_or_url).read_text(encoding="utf-8")
source = path_or_url
soup = BeautifulSoup(html, "html.parser")
# Удаляем технический мусор
for tag in soup(["script", "style", "nav", "header", "footer", "aside"]):
tag.decompose()
title = soup.title.string.strip() if soup.title else Path(source).stem
text = soup.get_text(separator="\n", strip=True)
# Убираем строки из одних пробелов/пустышек
text = "\n".join(line for line in text.splitlines() if len(line.strip()) > 15)
return Document(text=text, source=source, title=title, metadata={"type": "html"})
def parse_markdown(path: str) -> Document:
import re
text = Path(path).read_text(encoding="utf-8")
# Ищем заголовок h1 в начале файла
match = re.match(r'^#\s+(.+)', text.strip())
title = match.group(1) if match else Path(path).stem
return Document(text=text, source=path, title=title, metadata={"type": "markdown"})
# ── Маршрутизатор по расширению ─────────────────────────────────────
PARSERS = {
".pdf": parse_pdf,
".docx": parse_docx,
".html": parse_html,
".htm": parse_html,
".md": parse_markdown,
".txt": lambda p: Document(
text=Path(p).read_text(encoding="utf-8"),
source=p,
title=Path(p).stem,
metadata={"type": "txt"},
),
}
def parse_document(path: str) -> Document | None:
ext = Path(path).suffix.lower()
parser = PARSERS.get(ext)
if not parser:
print(f"[SKIP] Неизвестный формат: {path}")
return None
try:
doc = parser(path)
if len(doc.text.strip()) < 50:
print(f"[SKIP] Слишком мало текста: {path}")
return None
return doc
except Exception as e:
print(f"[ERROR] Парсинг {path}: {e}")
return None # не падаем на одном плохом файле
Этап 2: Chunking с наследованием метаданных
Chunking в pipeline имеет одно принципиальное требование сверх базового сплиттинга:
каждый чанк должен нести метаданные родительского документа.
В момент retrieval у нас есть только чанк — документ уже недоступен.
Если в чанке нет source или title, LLM не сможет сослаться на источник.
import hashlib
from dataclasses import dataclass, field
@dataclass
class Chunk:
text: str
doc_id: str
chunk_id: str = ""
chunk_index: int = 0
metadata: dict = field(default_factory=dict)
def __post_init__(self):
if not self.chunk_id:
raw = f"{self.doc_id}:{self.chunk_index}"
self.chunk_id = hashlib.md5(raw.encode()).hexdigest()
def split_document(
doc: Document,
chunk_size: int = 512,
chunk_overlap: int = 64,
) -> list[Chunk]:
"""
Рекурсивный сплиттер: сначала по \n\n, потом \n, потом по предложениям.
Каждый чанк наследует метаданные Document.
"""
separators = ["\n\n", "\n", ". ", " ", ""]
raw_chunks = _recursive_split(doc.text, chunk_size, chunk_overlap, separators)
chunks = []
for i, text in enumerate(raw_chunks):
if not text.strip():
continue
chunks.append(Chunk(
text=text.strip(),
doc_id=doc.doc_id,
chunk_index=i,
metadata={
**doc.metadata, # всё из Document
"source": doc.source,
"title": doc.title,
"doc_id": doc.doc_id,
"chunk_index": i,
"total_chunks": 0, # заполним после
"content_hash": doc.content_hash,
},
))
# Обновляем total_chunks теперь, когда знаем итоговое количество
for ch in chunks:
ch.metadata["total_chunks"] = len(chunks)
return chunks
def _recursive_split(
text: str,
chunk_size: int,
chunk_overlap: int,
separators: list[str],
) -> list[str]:
"""Рекурсивный сплит: пробует разделители по очереди."""
if len(text) <= chunk_size:
return [text]
sep = next((s for s in separators if s in text), "")
splits = text.split(sep) if sep else [text[i:i+chunk_size]
for i in range(0, len(text), chunk_size)]
chunks, current = [], ""
for split in splits:
candidate = (current + sep + split).strip() if current else split
if len(candidate) <= chunk_size:
current = candidate
else:
if current:
chunks.append(current)
# Если сплит сам по себе больше chunk_size — рекурсивно дробим дальше
if len(split) > chunk_size:
chunks.extend(_recursive_split(split, chunk_size, chunk_overlap, separators[1:]))
current = ""
else:
current = split
if current:
chunks.append(current)
# Добавляем overlap: берём хвост предыдущего чанка
if chunk_overlap > 0 and len(chunks) > 1:
overlapped = [chunks[0]]
for i in range(1, len(chunks)):
prev_tail = chunks[i-1][-chunk_overlap:]
overlapped.append(prev_tail + " " + chunks[i])
return overlapped
return chunks
Этап 3: Батчевая векторизация
Наивная реализация — embedding по одному чанку — в 20–50 раз медленнее батчевой. При корпусе в 10 000 чанков разница: 10 минут против 20 секунд.
import numpy as np
from dataclasses import dataclass
from sentence_transformers import SentenceTransformer
from tqdm import tqdm
@dataclass
class EmbeddedChunk:
chunk_id: str
vector: np.ndarray # shape (dim,), L2-нормализован
text: str
metadata: dict
def embed_chunks(
chunks: list[Chunk],
model_name: str = "BAAI/bge-small-en-v1.5",
batch_size: int = 64,
show_progress: bool = True,
) -> list[EmbeddedChunk]:
"""
Батчевая векторизация списка чанков.
Возвращает EmbeddedChunk с L2-нормализованными векторами.
"""
model = SentenceTransformer(model_name)
texts = [ch.text for ch in chunks]
# encode() принимает список и обрабатывает батчами внутри
# normalize_embeddings=True → косинусное сходство = dot product
vectors = model.encode(
texts,
batch_size=batch_size,
normalize_embeddings=True,
show_progress_bar=show_progress,
convert_to_numpy=True,
)
return [
EmbeddedChunk(
chunk_id=ch.chunk_id,
vector=vectors[i],
text=ch.text,
metadata=ch.metadata,
)
for i, ch in enumerate(chunks)
]
BAAI/bge-small-en-v1.5 (33M параметров, dim=384) — хорошее соотношение
скорость/качество для английского. BAAI/bge-m3 (570M) — мультиязычная,
поддерживает русский. text-embedding-3-small от OpenAI — через API,
без GPU, dim=1536. Для Russian: intfloat/multilingual-e5-small или
sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2.
Этап 4: Запись в векторную БД
Хранение — последний этап, но самый важный для производительности в проде. Два ключевых решения: как делать upsert (не insert!) и какой размер батча использовать.
Upsert вместо insert — обязательно. При повторном индексировании (например, документ обновился) insert создаст дубликаты, upsert перезапишет по ID. Все production векторные БД поддерживают upsert.
import uuid
from qdrant_client import QdrantClient
from qdrant_client.models import (
Distance, VectorParams, PointStruct, OptimizersConfigDiff,
)
def create_collection_if_not_exists(
client: QdrantClient,
collection_name: str,
dim: int,
):
"""Создаёт коллекцию, если её нет. Idempotent."""
existing = [c.name for c in client.get_collections().collections]
if collection_name not in existing:
client.create_collection(
collection_name=collection_name,
vectors_config=VectorParams(size=dim, distance=Distance.COSINE),
# Отключаем оптимизатор на время индексирования → быстрее запись
optimizers_config=OptimizersConfigDiff(indexing_threshold=0),
)
def store_chunks_qdrant(
client: QdrantClient,
collection_name: str,
embedded_chunks: list[EmbeddedChunk],
batch_size: int = 100,
) -> int:
"""
Upsert батчами. Возвращает количество записанных точек.
"""
total = 0
for i in range(0, len(embedded_chunks), batch_size):
batch = embedded_chunks[i : i + batch_size]
points = [
PointStruct(
id=str(uuid.uuid5(uuid.NAMESPACE_DNS, ec.chunk_id)),
vector=ec.vector.tolist(),
payload={
"text": ec.text,
"chunk_id": ec.chunk_id,
**ec.metadata,
},
)
for ec in batch
]
client.upsert(collection_name=collection_name, points=points)
total += len(points)
return total
# ── Альтернатива: Chroma ────────────────────────────────────────────
def store_chunks_chroma(
collection, # chromadb.Collection
embedded_chunks: list[EmbeddedChunk],
batch_size: int = 100,
) -> int:
total = 0
for i in range(0, len(embedded_chunks), batch_size):
batch = embedded_chunks[i : i + batch_size]
collection.upsert(
ids=[ec.chunk_id for ec in batch],
embeddings=[ec.vector.tolist() for ec in batch],
documents=[ec.text for ec in batch],
metadatas=[ec.metadata for ec in batch],
)
total += len(batch)
return total
indexing_threshold=0
перед загрузкой, а после — верни обратно (indexing_threshold=20000).
Скорость записи 10k точек: с индексированием ~4 мин, без ~40 сек.
Полный pipeline: сборка всех этапов
import logging
from pathlib import Path
from dataclasses import dataclass, field
from qdrant_client import QdrantClient
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger(__name__)
@dataclass
class IndexingConfig:
collection_name: str = "documents"
model_name: str = "BAAI/bge-small-en-v1.5"
chunk_size: int = 512
chunk_overlap: int = 64
embed_batch: int = 64
store_batch: int = 100
class IndexingPipeline:
def __init__(self, config: IndexingConfig, qdrant_url: str = "http://localhost:6333"):
self.cfg = config
self.client = QdrantClient(url=qdrant_url)
# Определяем размерность до создания коллекции
from sentence_transformers import SentenceTransformer
self.model = SentenceTransformer(config.model_name)
dim = self.model.get_sentence_embedding_dimension()
create_collection_if_not_exists(self.client, config.collection_name, dim)
# ─────────────────────────────────────────────────────────
def index_files(self, paths: list[str]) -> dict:
"""
Индексирует список файлов. Возвращает статистику.
Не падает на одном плохом файле — логирует ошибку и идёт дальше.
"""
stats = {"parsed": 0, "skipped": 0, "chunks": 0, "stored": 0, "errors": 0}
for path in paths:
try:
result = self._index_one(path)
if result is None:
stats["skipped"] += 1
else:
stats["parsed"] += 1
stats["chunks"] += result["chunks"]
stats["stored"] += result["stored"]
except Exception as e:
log.error(f"Ошибка при индексировании {path}: {e}")
stats["errors"] += 1
log.info(f"Готово: {stats}")
return stats
def index_directory(self, directory: str, glob: str = "**/*") -> dict:
"""Рекурсивно индексирует директорию."""
paths = [
str(p) for p in Path(directory).glob(glob)
if p.suffix.lower() in PARSERS and p.is_file()
]
log.info(f"Найдено файлов: {len(paths)}")
return self.index_files(paths)
# ─────────────────────────────────────────────────────────
def _index_one(self, path: str) -> dict | None:
"""Полный цикл для одного документа. Возвращает None, если пропущен."""
# Парсинг
doc = parse_document(path)
if doc is None:
return None
# Chunking
chunks = split_document(doc, self.cfg.chunk_size, self.cfg.chunk_overlap)
if not chunks:
log.warning(f"Нет чанков: {path}")
return None
# Embedding
embedded = embed_chunks(
chunks,
model_name=self.cfg.model_name,
batch_size=self.cfg.embed_batch,
show_progress=False,
)
# Storage
stored = store_chunks_qdrant(
self.client,
self.cfg.collection_name,
embedded,
batch_size=self.cfg.store_batch,
)
log.info(f"✓ {Path(path).name}: {len(chunks)} чанков, {stored} точек")
return {"chunks": len(chunks), "stored": stored}
# ── Использование ───────────────────────────────────────────────────
if __name__ == "__main__":
config = IndexingConfig(collection_name="my_docs", model_name="BAAI/bge-small-en-v1.5")
pipeline = IndexingPipeline(config, qdrant_url="http://localhost:6333")
stats = pipeline.index_directory("./docs", glob="**/*.{pdf,docx,md,html}")
print(stats)
Инкрементальный индексинг: обновления без перестройки
Полная перестройка индекса при каждом обновлении документа — дорого и долго. Инкрементальный подход: проверяем хеш документа, пересчитываем только изменившиеся.
Три сценария обновления:
doc_id → индексируем целиком.content_hash отличается от хранящегося → удаляем старые чанки по doc_id, индексируем заново.content_hash совпадает → пропускаем, экономим время и деньги.doc_id из БД.from qdrant_client.models import Filter, FieldCondition, MatchValue
def get_indexed_hashes(
client: QdrantClient,
collection_name: str,
) -> dict[str, str]:
"""
Возвращает {doc_id: content_hash} для всех документов в БД.
Используется для быстрой проверки «изменился ли документ».
"""
hashes = {}
offset = None
while True:
response, next_offset = client.scroll(
collection_name=collection_name,
scroll_filter=None,
limit=1000,
offset=offset,
with_payload=["doc_id", "content_hash"],
with_vectors=False,
)
for point in response:
doc_id = point.payload.get("doc_id")
content_hash = point.payload.get("content_hash")
if doc_id and content_hash:
hashes[doc_id] = content_hash
if next_offset is None:
break
offset = next_offset
return hashes
def delete_doc_chunks(
client: QdrantClient,
collection_name: str,
doc_id: str,
) -> None:
"""Удаляет все чанки, принадлежащие документу doc_id."""
from qdrant_client.models import FilterSelector
client.delete(
collection_name=collection_name,
points_selector=FilterSelector(
filter=Filter(
must=[FieldCondition(key="doc_id", match=MatchValue(value=doc_id))]
)
),
)
class IncrementalIndexingPipeline(IndexingPipeline):
def sync_directory(self, directory: str, glob: str = "**/*") -> dict:
"""
Синхронизирует директорию с индексом:
- новые файлы → индексирует
- изменившиеся → удаляет старые чанки, индексирует заново
- удалённые → удаляет из БД
- неизменившиеся → пропускает
"""
import hashlib
all_paths = {
str(p): hashlib.md5(p.stem.encode()).hexdigest()
for p in Path(directory).glob(glob)
if p.suffix.lower() in PARSERS and p.is_file()
}
# Текущее состояние БД
indexed = get_indexed_hashes(self.client, self.cfg.collection_name)
stats = {"new": 0, "updated": 0, "unchanged": 0, "deleted": 0}
# Обработка файлов на диске
for path, _ in all_paths.items():
doc = parse_document(path)
if doc is None:
continue
if doc.doc_id not in indexed:
# Новый документ
self._index_one(path)
stats["new"] += 1
elif indexed[doc.doc_id] != doc.content_hash:
# Документ изменился
delete_doc_chunks(self.client, self.cfg.collection_name, doc.doc_id)
self._index_one(path)
stats["updated"] += 1
else:
# Без изменений
stats["unchanged"] += 1
# Удаляем из БД файлы, которых больше нет на диске
current_doc_ids = set()
for path in all_paths:
doc = parse_document(path)
if doc:
current_doc_ids.add(doc.doc_id)
for doc_id in indexed:
if doc_id not in current_doc_ids:
delete_doc_chunks(self.client, self.cfg.collection_name, doc_id)
stats["deleted"] += 1
log.info(f"Sync: {stats}")
return stats
Мониторинг и логирование
Pipeline без мониторинга — это чёрный ящик. При корпусе в тысячи документов важно знать: сколько проиндексировано, сколько пропущено, где ошибки и какова скорость.
def pipeline_health_check(
client: QdrantClient,
collection_name: str,
) -> dict:
"""
Возвращает статистику по коллекции: количество точек,
распределение по источникам, средний размер чанков.
"""
info = client.get_collection(collection_name)
total_points = info.points_count
# Семплируем 200 точек для агрегатов
sample, _ = client.scroll(
collection_name=collection_name,
limit=200,
with_payload=["source", "text"],
with_vectors=False,
)
sources: dict[str, int] = {}
text_lengths = []
for point in sample:
src = point.payload.get("source", "unknown")
# Берём только имя файла для читаемости
src_short = src.split("/")[-1] if "/" in src else src
sources[src_short] = sources.get(src_short, 0) + 1
text_lengths.append(len(point.payload.get("text", "")))
avg_len = sum(text_lengths) / len(text_lengths) if text_lengths else 0
return {
"total_points": total_points,
"unique_sources_in_sample": len(sources),
"avg_chunk_chars": round(avg_len),
"top_sources": dict(sorted(sources.items(), key=lambda x: -x[1])[:5]),
}
# ── Пример вывода ───────────────────────────────────────────────────
# {
# "total_points": 8432,
# "unique_sources_in_sample": 23,
# "avg_chunk_chars": 487,
# "top_sources": {"api_reference.md": 41, "guide.pdf": 38, ...}
# }
Шпаргалка
1. Parse — извлечь текст + метаданные из PDF/DOCX/HTML/MD. Не падать на плохих файлах.
2. Chunk — нарезать рекурсивным сплиттером (chunk_size=512, overlap=64). Каждый чанк наследует метаданные.
3. Embed — батч 64–128 чанков, normalize_embeddings=True. Скорость: в 20–50× быстрее поодиночке.
4. Store — upsert (не insert!) по chunk_id. Батч 100 точек.
Модели данных:
•
Document → content_hash, doc_id, metadata•
Chunk → chunk_id = hash(doc_id + index), inherits metadata•
EmbeddedChunk → L2-нормализованный vector + text + metadataИнкрементальный индексинг:
• Хранить content_hash в payload → сравнивать при обновлении
• Изменился → delete by doc_id, затем reindex
• Удалён → delete by doc_id
Производительность:
• Отключить HNSW на время загрузки:
indexing_threshold=0• Embed батчами: 64–128 чанков за раз
• Store батчами: 100 точек за upsert
• Для >100k документов → asyncio.gather для параллельного парсинга
Практика
-
Запустите pipeline на своём корпусе. Возьмите папку с 20–50 PDF/MD-файлами.
Проиндексируйте через
IndexingPipeline, затем вызовитеpipeline_health_check(). Сколько чанков получилось? Какой средний размер? Есть ли файлы с 0 чанков? -
Реализуйте dry-run режим. Добавьте в
IndexingPipelineфлагdry_run=True, при котором pipeline делает все шаги кроме записи в БД и выводит статистику: сколько файлов, чанков, ожидаемый объём в БД. Полезно перед индексированием большого корпуса. -
Добавьте поддержку нового формата. Реализуйте парсер для
.epubчерез библиотекуebooklibи зарегистрируйте его в словареPARSERS. Проверьте, что метаданные (title, author) корректно извлекаются.