Почему БД — нетривиальный источник для RAG

На первый взгляд задача простая: выбрать строки из таблицы и индексировать. Но реляционные данные — структурированные, нормализованные, разбитые на колонки — плохо подходят для семантического поиска в сыром виде. Рассмотрим конкретный пример:

Таблица products:
┌────┬───────────────────────┬──────────────┬───────┬───────────┬───────────┐
│ id │ name                  │ category     │ price │ in_stock  │ rating    │
├────┼───────────────────────┼──────────────┼───────┼───────────┼───────────┤
│ 42 │ Ergonomic Chair Pro   │ Furniture    │ 42900 │ true      │ 4.7       │
└────┴───────────────────────┴──────────────┴───────┴───────────┴───────────┘

Что увидит embedding-модель, если передать сырую строку:
"42 Ergonomic Chair Pro Furniture 42900 true 4.7"

Что нужно для хорошего поиска по запросу "удобное офисное кресло":
"Ergonomic Chair Pro — эргономичное офисное кресло из категории Furniture.
Цена: 42 900 ₽. В наличии. Рейтинг покупателей: 4.7/5."

Разница принципиальная: числа и булевы значения без контекста не несут семантики для языковой модели. Слово true не значит «в наличии», пока нет поля in_stock рядом с ним. Цена 42900 без единицы и подписи — просто число.

Три ключевых вопроса при работе с БД:

  1. Сериализация: как превратить кортеж значений в осмысленный текст.
  2. Гранулярность: один Document на строку, группу строк или всю таблицу.
  3. Синхронизация: как узнать, что запись изменилась, и обновить индекс.

Теория: схема важнее данных

Прежде чем сериализовать строки, нужно понять схему таблицы. Имена колонок — это уже половина контекста. Знать типы данных — значит знать, как форматировать значения. Знать внешние ключи — значит понимать, куда идти за связанными данными.

products
id INTEGER PK В metadata, не в текст
name VARCHAR(255) Главное поле → в начало текста
description TEXT Длинный текст → основа чанка
category_id INTEGER FK JOIN с categories.name → текстовое значение
price NUMERIC(10,2) Форматировать: «42 900 ₽», не «42900»
in_stock BOOLEAN Конвертировать: «в наличии» / «под заказ»
updated_at TIMESTAMPTZ Для инкрементальной синхронизации
categories
id INTEGER PK
name VARCHAR(100) «Мебель для офиса» вместо category_id=7
path LTREE Иерархия: «Офис > Мебель > Кресла» → в metadata.tags
💡 Правило разделения полей: не все колонки должны идти в page_content. Первичные ключи, технические ID, даты изменения — в metadata (для фильтрации и ссылок). Описательные поля — в текст. Числовые метрики — с единицами и контекстом: «цена 42 900 ₽», не просто «42900».

Стратегии сериализации строки в текст

Существует четыре основных подхода. Выбор зависит от типа данных и требуемого качества текста.

Template-based
Шаблоны · f-string / Jinja2
"{name} — {category}. Цена: {price} ₽. {stock_str}. Рейтинг: {rating}/5. {description}"
Плюсы: детерминировано, быстро, без LLM
Минусы: жёсткая привязка к схеме
Когда: структурированные записи (товары, контакты)
Key-Value формат
Авто-генерация из колонок
Название: Ergonomic Chair Pro Категория: Мебель для офиса Цена: 42 900 ₽ В наличии: да Рейтинг: 4.7
Плюсы: универсально, не нужен шаблон
Минусы: не самый естественный язык
Когда: много разных таблиц, прототип
NLG через LLM
Генерация естественного описания
"Ergonomic Chair Pro — эргономичное офисное кресло с поясничной поддержкой и регулируемыми подлокотниками. Цена 42 900 ₽, высокий рейтинг 4.7/5."
Плюсы: лучшее качество текста
Минусы: медленно, дорого, галлюцинации
Когда: важна читаемость в ответе LLM
Column-as-Document
Один документ = одна колонка
Описание товара #42: "Эргономичное кресло с... (полный текст описания)" → doc per description column only
Плюсы: подходит для длинных текстовых полей
Минусы: теряется контекст других полей
Когда: TEXT-колонка с полным описанием
100%
колёсико — масштаб  ·  зажать и тянуть — перемещение
PostgreSQL id · name · price · … 42 · Chair Pro · 42900 · … 43 · Standing Desk · … 44 · Monitor Arm · … updated_at ↑ инкр. синхронизация Schema col names · types · FK Schema Extractor column map · type hints Row Query SELECT + JOIN + WHERE Serializer template / kv / nlg NULL handling · format Document page_content: сериализованный текст metadata: id · table · updated_at · … → Chunking → Embedding → Vector DB Sync State last_sync_at · row hashes

Гранулярность: строка, группа или таблица

Строка → Document
Один Document на каждую запись таблицы. Идеально для продуктов, статей, контактов — когда каждая сущность самодостаточна.
products, employees, articles, FAQ-пары
Группа строк → Document
Несколько связанных строк объединяются в один Document. Например: все комментарии к одному тикету поддержки.
support_tickets + comments, order + order_items
Таблица → Document
Вся таблица как один Document — для небольших справочников: таблица валют, список стран, настройки системы.
currency_rates, country_codes, system_settings

Извлечение схемы: автоматическое определение колонок

Прежде чем сериализовать строки, нужно знать структуру таблицы: какие колонки есть, какого они типа, какие из них являются FK. Это позволяет строить шаблоны автоматически и форматировать значения правильно.

# db/schema.py
from dataclasses import dataclass, field
from typing import Any
import asyncpg


@dataclass
class ColumnInfo:
    name: str
    data_type: str          # "text", "integer", "boolean", "numeric", "timestamptz", …
    is_nullable: bool
    is_primary_key: bool
    foreign_key: str | None  # "categories.id" или None


@dataclass
class TableSchema:
    table_name: str
    columns: list[ColumnInfo]
    primary_key: str = "id"
    updated_at_col: str = "updated_at"   # для инкр. синхронизации

    @property
    def text_columns(self) -> list[str]:
        """Текстовые колонки — основа page_content."""
        return [c.name for c in self.columns
                if c.data_type in ("text", "varchar", "character varying")
                and not c.is_primary_key]

    @property
    def numeric_columns(self) -> list[str]:
        return [c.name for c in self.columns
                if c.data_type in ("integer", "bigint", "numeric", "real", "double precision")
                and not c.is_primary_key]


async def extract_schema(
    conn: asyncpg.Connection,
    table_name: str,
    schema: str = "public",
) -> TableSchema:
    """
    Извлекает схему таблицы из information_schema.
    Определяет PK и FK через pg_constraint.
    """
    # Колонки с типами
    col_rows = await conn.fetch("""
        SELECT
            column_name,
            data_type,
            is_nullable
        FROM information_schema.columns
        WHERE table_schema = $1 AND table_name = $2
        ORDER BY ordinal_position
    """, schema, table_name)

    # Первичные ключи
    pk_rows = await conn.fetch("""
        SELECT kcu.column_name
        FROM information_schema.table_constraints tc
        JOIN information_schema.key_column_usage kcu
          ON tc.constraint_name = kcu.constraint_name
        WHERE tc.table_schema = $1 AND tc.table_name = $2
          AND tc.constraint_type = 'PRIMARY KEY'
    """, schema, table_name)
    pk_cols = {r["column_name"] for r in pk_rows}

    # Внешние ключи
    fk_rows = await conn.fetch("""
        SELECT
            kcu.column_name,
            ccu.table_name AS foreign_table,
            ccu.column_name AS foreign_column
        FROM information_schema.table_constraints tc
        JOIN information_schema.key_column_usage kcu
          ON tc.constraint_name = kcu.constraint_name
        JOIN information_schema.constraint_column_usage ccu
          ON tc.constraint_name = ccu.constraint_name
        WHERE tc.table_schema = $1 AND tc.table_name = $2
          AND tc.constraint_type = 'FOREIGN KEY'
    """, schema, table_name)
    fk_map = {r["column_name"]: f"{r['foreign_table']}.{r['foreign_column']}"
              for r in fk_rows}

    columns = [
        ColumnInfo(
            name=r["column_name"],
            data_type=r["data_type"],
            is_nullable=(r["is_nullable"] == "YES"),
            is_primary_key=(r["column_name"] in pk_cols),
            foreign_key=fk_map.get(r["column_name"]),
        )
        for r in col_rows
    ]

    return TableSchema(table_name=table_name, columns=columns)

Сериализатор: от строки к тексту

Сериализатор — ключевой компонент пайплайна. Он принимает словарь значений строки и схему таблицы, а возвращает строку текста. Хорошая сериализация: форматирует числа с единицами, конвертирует булевы значения в слова и пропускает поля с NULL.

# db/serializer.py
from decimal import Decimal
from datetime import datetime, date
from typing import Any
from db.schema import TableSchema


def format_value(value: Any, col_name: str, data_type: str) -> str | None:
    """
    Форматирует значение колонки в читаемый текст.
    Возвращает None для NULL и пустых значений.
    """
    if value is None:
        return None

    # Булевы → слова
    if data_type == "boolean":
        return "да" if value else "нет"

    # Деньги: добавляем форматирование по имени колонки
    if data_type in ("numeric", "real", "double precision"):
        money_hints = {"price", "cost", "amount", "total", "salary", "fee"}
        if any(hint in col_name.lower() for hint in money_hints):
            return f"{Decimal(str(value)):,.0f}".replace(",", " ") + " ₽"
        percent_hints = {"rate", "percent", "discount", "tax"}
        if any(hint in col_name.lower() for hint in percent_hints):
            return f"{float(value):.1f}%"
        return str(value)

    # Даты → локаль
    if isinstance(value, datetime):
        return value.strftime("%d.%m.%Y %H:%M")
    if isinstance(value, date):
        return value.strftime("%d.%m.%Y")

    # Списки (PostgreSQL arrays)
    if isinstance(value, (list, tuple)):
        return ", ".join(str(v) for v in value if v)

    text = str(value).strip()
    return text if text else None


# Человекочитаемые имена колонок (можно расширять)
COLUMN_LABELS = {
    "name": "Название",
    "title": "Заголовок",
    "description": "Описание",
    "category": "Категория",
    "price": "Цена",
    "in_stock": "В наличии",
    "rating": "Рейтинг",
    "author": "Автор",
    "status": "Статус",
    "created_at": "Создан",
    "updated_at": "Обновлён",
    "tags": "Теги",
    "email": "Email",
    "phone": "Телефон",
}


def col_label(col_name: str) -> str:
    return COLUMN_LABELS.get(col_name, col_name.replace("_", " ").capitalize())


class KeyValueSerializer:
    """
    Универсальная сериализация: «Название: значение» для каждой колонки.
    Не требует шаблона — работает с любой таблицей.
    """

    def __init__(
        self,
        schema: TableSchema,
        include_cols: list[str] | None = None,   # None = все
        exclude_cols: list[str] | None = None,
        title_col: str | None = None,            # колонка для первой строки
    ):
        self._schema = schema
        self._col_types = {c.name: c.data_type for c in schema.columns}
        self._include = set(include_cols) if include_cols else None
        self._exclude = set(exclude_cols or []) | {schema.primary_key, schema.updated_at_col}
        self._title_col = title_col

    def serialize(self, row: dict) -> str:
        parts = []

        # Заголовок — первым
        if self._title_col and row.get(self._title_col):
            title_val = format_value(
                row[self._title_col], self._title_col,
                self._col_types.get(self._title_col, "text")
            )
            if title_val:
                parts.append(title_val)

        # Остальные колонки
        for col in self._schema.columns:
            name = col.name
            if name == self._title_col:
                continue
            if name in self._exclude:
                continue
            if self._include and name not in self._include:
                continue
            if col.is_primary_key:
                continue

            value = row.get(name)
            formatted = format_value(value, name, col.data_type)
            if formatted:
                parts.append(f"{col_label(name)}: {formatted}")

        return "\n".join(parts)


class TemplateSerializer:
    """
    Сериализация через f-string шаблон.
    Для конкретных таблиц, где важен порядок и формат.
    """

    def __init__(self, template: str, schema: TableSchema):
        self._template = template
        self._col_types = {c.name: c.data_type for c in schema.columns}

    def serialize(self, row: dict) -> str:
        # Форматируем все значения заранее
        formatted = {}
        for key, value in row.items():
            fv = format_value(value, key, self._col_types.get(key, "text"))
            formatted[key] = fv or ""
        try:
            return self._template.format(**formatted).strip()
        except KeyError as e:
            raise ValueError(f"Template references missing column: {e}") from e

DB Loader: полный пайплайн

Теперь соединяем всё вместе: подключение к БД, выполнение запроса с JOIN-ами, сериализация и формирование Document-объектов.

# db/loader.py
import os
import asyncpg
from pathlib import Path
from core.document import Document
from db.schema import TableSchema, extract_schema
from db.serializer import KeyValueSerializer, TemplateSerializer


class DBLoader:
    """
    Загружает строки таблицы как Document-объекты.
    Поддерживает: JOIN-ы, инкрементальную синхронизацию,
    кастомный SQL-запрос.
    """

    def __init__(
        self,
        dsn: str | None = None,
        table: str | None = None,
        query: str | None = None,           # кастомный SQL (с JOIN-ами)
        serializer=None,                    # KeyValueSerializer или TemplateSerializer
        metadata_cols: list[str] | None = None,  # колонки → только в metadata
        batch_size: int = 500,
    ):
        self._dsn = dsn or os.environ["DATABASE_URL"]
        self._table = table
        self._query = query
        self._serializer = serializer
        self._metadata_cols = set(metadata_cols or [])
        self._batch_size = batch_size

    async def load(self, modified_after: str = "") -> list[Document]:
        """
        Загружает записи.
        modified_after — ISO timestamp для инкрементальной загрузки.
        """
        conn = await asyncpg.connect(self._dsn)
        try:
            # Если схема не задана — извлекаем автоматически
            if self._table and not self._serializer:
                schema = await extract_schema(conn, self._table)
                self._serializer = KeyValueSerializer(
                    schema=schema,
                    title_col=self._detect_title_col(schema),
                )

            sql, params = self._build_query(modified_after)
            documents = []
            offset = 0

            while True:
                paginated = f"{sql} LIMIT {self._batch_size} OFFSET {offset}"
                rows = await conn.fetch(paginated, *params)
                if not rows:
                    break

                for row in rows:
                    doc = self._row_to_document(dict(row))
                    if doc:
                        documents.append(doc)

                offset += self._batch_size
                if len(rows) < self._batch_size:
                    break   # последняя страница

            return documents
        finally:
            await conn.close()

    def _build_query(self, modified_after: str) -> tuple[str, list]:
        params = []
        if self._query:
            sql = self._query
            if modified_after and "updated_at" in sql.lower():
                # Пробуем добавить фильтр (если в запросе есть WHERE — AND, иначе WHERE)
                clause = f"updated_at > ${len(params) + 1}"
                params.append(modified_after)
                if "where" in sql.lower():
                    sql += f" AND {clause}"
                else:
                    sql += f" WHERE {clause}"
        else:
            sql = f"SELECT * FROM {self._table}"
            if modified_after:
                sql += f" WHERE updated_at > $1"
                params.append(modified_after)
            sql += " ORDER BY id"
        return sql, params

    def _row_to_document(self, row: dict) -> Document | None:
        if not self._serializer:
            raise RuntimeError("No serializer configured")

        text = self._serializer.serialize(row)
        if not text.strip():
            return None

        # Метаданные: технические поля + явно указанные
        metadata = {
            "source":    f"db://{self._table or 'query'}/{row.get('id', '')}",
            "table":     self._table or "custom_query",
            "record_id": str(row.get("id", "")),
            "updated_at": str(row.get("updated_at", "")),
        }
        for col in self._metadata_cols:
            if col in row:
                metadata[col] = str(row[col]) if row[col] is not None else ""

        return Document(page_content=text, metadata=metadata)

    @staticmethod
    def _detect_title_col(schema: TableSchema) -> str | None:
        """Эвристика: ищем колонку name/title/subject."""
        for hint in ("name", "title", "subject", "full_name", "headline"):
            for col in schema.columns:
                if col.name.lower() == hint:
                    return col.name
        return None

Работа с JOIN-ами: денормализация для RAG

Нормализованная БД хранит category_id = 7 вместо category = "Мебель для офиса". Для embedding это плохо: нужно либо делать JOIN в SQL-запросе, либо подгружать связанные данные на уровне приложения.

# Вариант 1: денормализующий SQL с JOIN
PRODUCTS_QUERY = """
    SELECT
        p.id,
        p.name,
        p.description,
        p.price,
        p.in_stock,
        p.rating,
        p.updated_at,
        c.name        AS category,
        c.path::text  AS category_path,
        array_agg(DISTINCT t.name) FILTER (WHERE t.name IS NOT NULL)
                      AS tags
    FROM products p
    LEFT JOIN categories c ON c.id = p.category_id
    LEFT JOIN product_tags pt ON pt.product_id = p.id
    LEFT JOIN tags t ON t.id = pt.tag_id
    GROUP BY p.id, c.name, c.path
    ORDER BY p.updated_at DESC
"""

PRODUCTS_TEMPLATE = """\
{name} — {category}
Цена: {price}. В наличии: {in_stock}.
Рейтинг: {rating}/5.
Теги: {tags}

{description}"""

# Создаём загрузчик с кастомным запросом и шаблоном
from db.schema import TableSchema, ColumnInfo
from db.serializer import TemplateSerializer

# Минимальная схема для TemplateSerializer (только типы нужных колонок)
schema_stub = TableSchema(
    table_name="products",
    columns=[
        ColumnInfo("id",          "integer", False, True,  None),
        ColumnInfo("name",        "text",    False, False, None),
        ColumnInfo("description", "text",    True,  False, None),
        ColumnInfo("price",       "numeric", False, False, None),
        ColumnInfo("in_stock",    "boolean", False, False, None),
        ColumnInfo("rating",      "numeric", True,  False, None),
        ColumnInfo("category",    "text",    True,  False, None),
        ColumnInfo("tags",        "ARRAY",   True,  False, None),
        ColumnInfo("updated_at",  "timestamptz", False, False, None),
    ],
)

loader = DBLoader(
    table="products",
    query=PRODUCTS_QUERY,
    serializer=TemplateSerializer(PRODUCTS_TEMPLATE, schema_stub),
    metadata_cols=["category", "in_stock", "rating"],
)

# Результат для одного товара:
# "Ergonomic Chair Pro — Мебель для офиса
#  Цена: 42 900 ₽. В наличии: да.
#  Рейтинг: 4.7/5.
#  Теги: кресло, офис, эргономика
#
#  Эргономичное кресло с поясничной поддержкой..."

Инкрементальная синхронизация

Таблица продуктов меняется постоянно: добавляются новые товары, меняются цены, товары снимаются с продажи. Перезагружать всё каждый раз — дорого. Нужна синхронизация только изменённых записей.

# db/sync.py
import asyncio
import asyncpg
import hashlib
import json
from datetime import datetime, timezone
from pathlib import Path
from core.document import Document
from db.loader import DBLoader


class DBSyncManager:
    """
    Управляет инкрементальной синхронизацией таблицы с индексом.
    Стратегия: updated_at + content_hash для обнаружения реальных изменений.
    """

    def __init__(self, loader: DBLoader, state_file: Path):
        self._loader = loader
        self._state_file = state_file
        self._state: dict = self._load_state()

    def _load_state(self) -> dict:
        if self._state_file.exists():
            return json.loads(self._state_file.read_text())
        return {"last_sync_at": "", "doc_hashes": {}}

    def _save_state(self) -> None:
        self._state_file.write_text(
            json.dumps(self._state, ensure_ascii=False, indent=2)
        )

    @staticmethod
    def _hash(text: str) -> str:
        return hashlib.sha256(text.encode()).hexdigest()[:16]

    async def sync(self) -> tuple[list[Document], list[str]]:
        """
        Синхронизирует изменения.
        Возвращает (new_or_updated, deleted_ids).
        """
        last_sync = self._state.get("last_sync_at", "")

        # 1. Загружаем изменённые записи
        docs = await self._loader.load(modified_after=last_sync)

        # 2. Разделяем: реально изменённые (по хэшу) vs просто обновлённые
        changed_docs = []
        for doc in docs:
            record_id = doc.metadata.get("record_id", "")
            new_hash = self._hash(doc.page_content)
            old_hash = self._state["doc_hashes"].get(record_id, "")
            if new_hash != old_hash:
                changed_docs.append(doc)
                self._state["doc_hashes"][record_id] = new_hash

        # 3. Определяем удалённые записи (если БД поддерживает soft-delete)
        # Для hard-delete нужна отдельная таблица аудита или CDC (Change Data Capture)
        deleted_ids = await self._find_deleted()

        # 4. Обновляем состояние
        self._state["last_sync_at"] = datetime.now(timezone.utc).isoformat()
        self._save_state()

        print(
            f"Sync: {len(changed_docs)} changed, {len(deleted_ids)} deleted "
            f"(since {last_sync or 'beginning'})"
        )
        return changed_docs, deleted_ids

    async def _find_deleted(self) -> list[str]:
        """
        Определяем удалённые записи через soft-delete флаг.
        Если таблица использует is_deleted/deleted_at — проверяем их.
        """
        # Для таблиц с deleted_at:
        # SELECT id FROM products WHERE deleted_at > last_sync_at
        # Для hard-delete — только через CDC (pglogical, Debezium)
        return []


# Пример: ночная синхронизация в cron
async def nightly_sync():
    loader = DBLoader(
        table="products",
        query=PRODUCTS_QUERY,
        serializer=TemplateSerializer(PRODUCTS_TEMPLATE, schema_stub),
        metadata_cols=["category", "in_stock"],
    )
    manager = DBSyncManager(loader, Path(".sync_state/products.json"))
    changed, deleted = await manager.sync()

    # changed → upsert в vector store
    # deleted → delete из vector store

asyncio.run(nightly_sync())
💡 CDC (Change Data Capture) — продвинутый подход для high-frequency обновлений. Вместо polling по updated_at используется поток изменений из WAL (Write-Ahead Log) PostgreSQL через pglogical или Debezium. Даёт near-realtime синхронизацию, но требует отдельной инфраструктуры. Для большинства RAG-задач достаточно ночного polling.

LLM-обогащение: генерация описаний

Иногда данные в БД — это только числа и коды. Товар без описания: SKU-4892, 42900, 7, 4.7. Шаблон не поможет — нет текста. Решение: генерировать описание с помощью LLM для записей, у которых нет текстового поля.

import anthropic
import asyncio
from db.serializer import KeyValueSerializer, col_label, format_value


class LLMEnricher:
    """
    Генерирует естественный текст для записей без описания.
    Использует claude-haiku — быстро и дёшево.
    """

    def __init__(self, model: str = "claude-haiku-4-5-20251001"):
        self._client = anthropic.AsyncAnthropic()
        self._model = model

    async def generate_description(
        self,
        row: dict,
        table_context: str,    # "товар из каталога интернет-магазина"
        lang: str = "ru",
    ) -> str:
        """Генерирует связное описание записи из её полей."""
        # Формируем кратко все поля для промпта
        fields_text = "\n".join(
            f"{col_label(k)}: {v}"
            for k, v in row.items()
            if v is not None and str(v).strip()
        )
        prompt = (
            f"Вот данные о {table_context}:\n\n{fields_text}\n\n"
            f"Напиши краткое (2-4 предложения) описание на русском языке "
            f"для поисковой системы. Не придумывай характеристики, "
            f"которых нет в данных. Используй только предоставленную информацию."
        )
        response = await self._client.messages.create(
            model=self._model,
            max_tokens=300,
            messages=[{"role": "user", "content": prompt}],
        )
        return response.content[0].text.strip()

    async def enrich_batch(
        self,
        rows: list[dict],
        table_context: str,
        concurrency: int = 5,
    ) -> list[str]:
        """Параллельная генерация для батча записей."""
        semaphore = asyncio.Semaphore(concurrency)

        async def enrich_one(row: dict) -> str:
            async with semaphore:
                return await self.generate_description(row, table_context)

        return await asyncio.gather(*[enrich_one(r) for r in rows])


# Использование: обогащаем записи без описания
async def enrich_products_without_description():
    enricher = LLMEnricher()
    conn = await asyncpg.connect(os.environ["DATABASE_URL"])
    rows = await conn.fetch(
        "SELECT * FROM products WHERE description IS NULL OR description = '' LIMIT 100"
    )
    rows_dicts = [dict(r) for r in rows]
    descriptions = await enricher.enrich_batch(rows_dicts, "товар из каталога")
    # Сохраняем сгенерированные описания обратно в БД
    for row, desc in zip(rows_dicts, descriptions):
        await conn.execute(
            "UPDATE products SET description = $1 WHERE id = $2",
            desc, row["id"],
        )
    await conn.close()

PostgreSQL + pgvector: гибридный поиск внутри БД

Для PostgreSQL есть особый сценарий: хранить векторы прямо в БД через расширение pgvector. Это позволяет объединить семантический поиск с SQL-фильтрами в одном запросе — без отдельного vector store.

-- Установка расширения (один раз)
CREATE EXTENSION IF NOT EXISTS vector;

-- Добавляем колонку для embedding (1536 dimensions для text-embedding-3-small)
ALTER TABLE products ADD COLUMN embedding vector(1536);

-- Индекс для быстрого ANN-поиска (IVFFlat)
CREATE INDEX ON products USING ivfflat (embedding vector_cosine_ops)
  WITH (lists = 100);  -- lists ≈ sqrt(N строк)
# Индексирование: сериализуем → эмбеддим → сохраняем в products.embedding
import anthropic
import asyncpg
import numpy as np


async def index_products(dsn: str, loader: "DBLoader") -> None:
    """Заполняет колонку embedding для всех продуктов."""
    client = anthropic.AsyncAnthropic()
    conn = await asyncpg.connect(dsn)

    docs = await loader.load()
    texts = [d.page_content for d in docs]
    ids   = [d.metadata["record_id"] for d in docs]

    # Батчами по 100 (лимит API)
    batch_size = 100
    for i in range(0, len(texts), batch_size):
        batch_texts = texts[i:i + batch_size]
        batch_ids   = ids[i:i + batch_size]

        response = await client.beta.messages.batches.create(...)
        # Или через embeddings API (voyage-3):
        emb_response = await client.embeddings.create(
            model="voyage-3",
            input=batch_texts,
        )
        vectors = [e.embedding for e in emb_response.data]

        # Массовое обновление через unnest
        await conn.executemany(
            "UPDATE products SET embedding = $1 WHERE id = $2",
            [(v, int(rid)) for v, rid in zip(vectors, batch_ids)],
        )
        print(f"  Indexed {i + len(batch_texts)}/{len(texts)}")

    await conn.close()


async def semantic_search_with_filters(
    query: str,
    dsn: str,
    category: str | None = None,
    min_rating: float | None = None,
    in_stock: bool | None = None,
    top_k: int = 10,
) -> list[dict]:
    """
    Гибридный поиск: pgvector + SQL-фильтры.
    Сначала фильтруем по metadata, затем ищем ближайших по вектору.
    """
    client = anthropic.AsyncAnthropic()
    conn = await asyncpg.connect(dsn)

    # Эмбеддим запрос
    emb = await client.embeddings.create(model="voyage-3", input=[query])
    query_vec = emb.data[0].embedding

    # WHERE-условия из фильтров
    filters, params = [], [query_vec]
    if category:
        params.append(category)
        filters.append(f"c.name = ${len(params)}")
    if min_rating is not None:
        params.append(min_rating)
        filters.append(f"p.rating >= ${len(params)}")
    if in_stock is not None:
        params.append(in_stock)
        filters.append(f"p.in_stock = ${len(params)}")

    where_clause = "WHERE " + " AND ".join(filters) if filters else ""
    params.append(top_k)

    sql = f"""
        SELECT
            p.id, p.name, p.price, p.rating, p.in_stock,
            c.name AS category,
            1 - (p.embedding <=> $1::vector) AS similarity
        FROM products p
        LEFT JOIN categories c ON c.id = p.category_id
        {where_clause}
        ORDER BY p.embedding <=> $1::vector
        LIMIT ${len(params)}
    """

    rows = await conn.fetch(sql, *params)
    await conn.close()
    return [dict(r) for r in rows]

Типичные ошибки

Индексировать числа без контекста
Строка "42 true 4.7 42900" в индексе совершенно бесполезна — эти числа не несут семантики. Запрос «недорогое кресло в наличии» никогда не найдёт эту запись.
→ Всегда форматируйте числа с единицами и подписями: «Цена: 42 900 ₽», «В наличии: да».
Включать технические колонки в page_content
Первичные ключи, UUID, хэши, созданные_at — это шум в тексте, который не помогает поиску, но занимает токены в контексте при генерации ответа.
→ PK, FK и технические поля — только в metadata. В page_content только то, по чему пользователь будет искать.
Не обрабатывать NULL
Шаблон f"{description}" при description=None даст строку "None" в тексте. В сериализованном тексте сотни «None» мешают поиску и выглядят как мусор.
→ Функция format_value() возвращает None для пустых значений — пропускайте такие поля целиком.
Один огромный Document на всю таблицу
Сериализовать всю таблицу products (10 000 строк) в один текст — это несколько мегабайт, которые не влезут ни в один контекст и не смогут быть полезно встроены в один вектор.
→ Один Document = одна строка (или группа связанных строк). Таблица-как-документ — только для маленьких справочников (< 50 строк).
Не синхронизировать удаления
При инкрементальной синхронизации легко учесть INSERT и UPDATE, но забыть про DELETE. Удалённый товар остаётся в индексе — пользователь получает ответ с несуществующей позицией.
→ Используйте soft-delete (колонка deleted_at) или CDC. При hard-delete — храните список ID в vector store и периодически сверяйте.

Шпаргалка

📋 SQL → text пайплайн: ключевые решения
  • Сериализация — KeyValueSerializer для любых таблиц, TemplateSerializer когда важен порядок и формат.
  • Числа и булевы — всегда с единицами и словами: «42 900 ₽», «в наличии», «4.7/5».
  • NULL — пропускать поле целиком, не включать «None» в текст.
  • Технические поля (PK, FK, timestamps) — только в metadata, не в page_content.
  • JOIN-ы — делать в SQL: category_id → category_name в одном запросе.
  • Гранулярность — строка → Document (продукты, статьи); группа строк (тикет + комментарии); маленький справочник → один Document.
  • Инкрементальная синхронизация — фильтр по updated_at + content_hash для пропуска неизменённых.
  • Удаления — soft-delete или CDC; без этого в индексе будут «призраки».
  • pgvector — хранить эмбеддинги в PostgreSQL + SQL-фильтры в одном запросе.

Практическое задание

  1. Сериализатор для реальной таблицы. Возьмите любую таблицу из вашей БД (или создайте тестовую: products, employees, articles). Реализуйте KeyValueSerializer с правильными COLUMN_LABELS для ваших колонок. Загрузите 10 строк и выведите сериализованный текст. Проверьте: нет ли в тексте «None», необработанных чисел, технических ID.
  2. JOIN + шаблонная сериализация. Добавьте в таблицу внешний ключ (например, category_id →categories). Напишите SQL-запрос с JOIN, который подтягивает название категории. Создайте TemplateSerializer с шаблоном, включающим название категории. Убедитесь, что в тексте появляется «Категория: Мебель для офиса» вместо «category_id: 7».
  3. Инкрементальная синхронизация. Реализуйте DBSyncManager для вашей таблицы. Запустите первую синхронизацию (full load), измените 2–3 записи в БД, запустите вторую синхронизацию и убедитесь, что загружаются только изменённые строки. Выведите количество изменённых записей на каждом запуске.