Слой ingestion — как это работает на самом деле
Статус: current (as-built), сверено на commit
544d4e51. Целевая сервисная граница ingestion описана в../../30-services/ingestion/README.md. Этот документ — не альтернатива ей, а проверка по коду: что из целевого дизайна реально существует, как называются реальные символы и где проходит реальный путь выполнения.
Назначение
Слой ingestion отвечает за то, чтобы данные от поставщиков (цены, остатки, характеристики) регулярно попадали в систему и доезжали до нормализации. В коде это не отдельный сервис, а worker-процесс supplier-sync, который поднимается как fx-приложение и работает по тикерам — без публичного HTTP.
Участники
supplier-sync(backend/cmd/supplier-sync/main.go) — процесс-контейнер: поднимает fx-граф, запускает per-supplier тикеры, отдаёт loopback/healthz,/readyz,/metrics. Точка входа зафиксирована в../entry-points.md.- Pipelines (
backend/internal/core/ingestion/app/) —PipelineRouterвыбирает пайплайн по Capabilities поставщика:BulkSnapshotPipelineImplдля batch-поставщиков (Russvet, ETM, IEK, Systeme) или cursor-пайплайн для DKC. - Connectors (
backend/internal/core/ingestion/infra/<supplier>/и смежные пакеты) — HTTP-клиенты конкретных поставщиков, вызываемые из пайплайнов. - Normalization — преобразование сырых ответов поставщика во внутреннюю модель (
NormalizedOffer,Observation); подробный as-built разбор появится отдельным use-case файлом. - Offers — запись нормализованных данных в PostgreSQL (
supplier_offers,offer_observations) и транзакционный outbox для Kafka.
Как это работает: тик и диспетчер конвейеров
Тик (tick) — одно «срабатывание часов»: момент, когда воркер supplier-sync просыпается по таймеру и делает одну порцию работы для одного поставщика. В проде на каждого поставщика висит свой тикер со своим интервалом (напр. russvet — раз в 30 минут / 12 часов). Каждое срабатывание = один тик = «сходи к поставщику и забери данные один раз». В тестах тик дёргают вручную, без таймера.
PipelineRouter (backend/internal/core/ingestion/app/pipeline_router.go) — диспетчер. Сам данные не тянет: на каждый тик он смотрит «возможности» поставщика (Capabilities, паспорт «как этот поставщик умеет отдавать данные») и передаёт работу одному из трёх конвейеров:
| Конвейер | Когда выбирается | Что делает | Пример |
|---|---|---|---|
| cursor | IncrementalSync = revision/date | тянет только изменения с прошлой метки | DKC |
| bulk | ScheduledRefresh.UsesBulkSnapshot() + есть источник | каждый тик выгружает каталог пачкой | Russvet, ETM, IEK, Systeme |
| full | иначе (fallback) | legacy полный каталог | — |
Порядок проверки в RunTick: cursor → bulk → full. Для Russvet выбирается bulk. Диспетчер также пишет метрику «какой конвейер выбран» (PipelineRouted) — это телеметрия, на маршрутизацию она не влияет.
Отдельно: у Russvet сам bulk-конвейер разделён на два job kind — catalog (описания + характеристики) и commerce (цены + остатки); дефолтный тик обрабатывает catalog, коммерческие оси приходят отдельным commerce-тиком.
Use-case файлы слоя
russvet.md— сквозной путь Russvet: supplier-sync → БД/raw-хранилище → charnorm-worker → matcher-worker/LLM → canonical-assignment-worker → БД. Появится следующей задачей инициативы.
По мере разбора добавляются use-case файлы для остальных поставщиков и внутренних переходов слоя (нормализация, запись офферов).
Куда дальше
- Сквозная карта точек входа всей системы:
../entry-points.md. - Целевой дизайн сервисной границы ingestion:
../../30-services/ingestion/README.md. - Как читать раздел
60-flows/целиком:../README.md.