Metadata-Version: 2.5
Name: cdc-1c
Version: 0.1.20
Summary: Change data capture (CDC) from 1C:Enterprise to your data warehouse
Project-URL: Homepage, https://github.com/pavel-v-sobolev/cdc_1C
Project-URL: Repository, https://github.com/pavel-v-sobolev/cdc_1C
Project-URL: Issues, https://github.com/pavel-v-sobolev/cdc_1C/issues
Author-email: Pavel Sobolev <pavel-v-sobolev@yandex.ru>
License-Expression: MIT
License-File: LICENSE
Keywords: 1c,1c-enterprise,cdc,data-engineering,dwh,etl,postgres,python
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Intended Audience :: Information Technology
Classifier: Intended Audience :: Science/Research
Classifier: Intended Audience :: System Administrators
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Python :: 3.14
Classifier: Topic :: Database
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Requires-Python: >=3.10
Requires-Dist: croniter>=6.0.0
Requires-Dist: dbmerge>=1.0.22
Requires-Dist: requests>=2.33.0
Requires-Dist: sqlalchemy>=2.0.49
Requires-Dist: xmltodict>=1.0.4
Provides-Extra: dev
Requires-Dist: psycopg2-binary>=2.9.12; extra == 'dev'
Requires-Dist: pytest>=8; extra == 'dev'
Provides-Extra: postgres
Requires-Dist: psycopg2-binary>=2.9.12; extra == 'postgres'
Description-Content-Type: text/markdown

[![PyPI version](https://img.shields.io/pypi/v/cdc-1c.svg)](https://pypi.org/project/cdc-1c/)
[![Python versions](https://img.shields.io/pypi/pyversions/cdc-1c.svg)](https://pypi.org/project/cdc-1c/)


**cdc-1c** is a docker container and a Python library, that provides 1C system data loading to data warehouse using Change Data Capture apporach. \
It engages standard ODATA mechanism and standard 1C exchange plan mechanism to extract data from 1C system and upsert changes to the target DB.

**cdc-1c** - это докер контейнер и python-библиотека, предназначенные для получения данных из 1С, использующий подход CDC (загрузка изменений данных). \
Продукт использует стандартный интерфейс ODATA и механизм планов обмена для выгрузки изменений данных из системы 1С и обновления данных в целевой БД.

# Что нужно для работы
1) опубликовать базу 1с на web
2) настроить план обмена в конфигураторе и включить в его состав нужные объекты 1с
3) создать пользователя для доступа к odata и дать ему необходимые права (чтение и изменение к плану обмена, чтение к загружаемым объектам)
4) дать роль чтение всем пользователям к плану обмена (иначе будут ошибки при сохранении объектов)
5) создать узел обмена с использованием внешней обработки `cdc-1c.odt`
6) запустить загрузку: докер-образом `sobolevp/cdc-1c` (см. «Запуск в docker») либо python-библиотекой `cdc-1c` (см. ниже)

# Использование библиотеки python

Основной объект библиотеки это оркестратор `Replicator1C`, который читает изменения из 1С (через OData + план обмена) и
пишет их в целевую БД Postgres, подтверждая приём пакета только после успешного сохранения.

## Установка

```bash
pip install cdc-1c
```

Для записи изменений необходим **PostgreSQL** (Другие СУБД не тестировались, хотя в теории возможны).
Для записи используется библиотека dbmerge. Все необходимые схемы, таблицы и поля модуль создает сам.

## Быстрый старт

```python
from sqlalchemy import create_engine
from cdc_1c import Replicator1C

# pool_size >= full_load_workers + 2 (+1 на каждого обработчика и каждое расписание) —
# почему столько, см. «Сколько нужно соединений к БД»
engine = create_engine("postgresql+psycopg2://user:pass@localhost:5432/cdc_1c", pool_size=5)

rep = Replicator1C(
    odata_url="http://host/base/odata/standard.odata",
    odata_auth=("odata", "secret"),        # (user, password) либо None без авторизации
    exchange_name="ДляODATA",              # имя плана обмена в 1С
    queue_guid="a9bc23c5-3689-11f1-926c-0800270bc6cb",  # Ref_Key узла обмена
    engine=engine,
    db_schema="cdc_1c",                    # None → схема БД по умолчанию (public у Postgres)
    db_temp_schema="cdc_1c_tmp",           # схема промежуточных таблиц merge; None → схема данных
    request_timeout=60,                    # таймаут HTTP-запросов к 1С, сек (по умолчанию 60 на коннект, 900 на ответ)
    full_load_workers=2,                   # число фоновых потоков полной выгрузки
)

rep.run_forever(interval=60)               # цикл опроса раз в 60 секунд
```


### Как узнать guid узла обмена

`queue_guid` — это `Ref_Key` узла плана обмена, того самого, на который 1С регистрирует изменения
(`ЭтотУзел` не подходит: он описывает саму базу-источник). Если guid неизвестен, оставьте параметр
пустым (`queue_guid=""`, или просто не задавайте `CDC1C_QUEUE_GUID`) и запустите: чтение изменений
выведет в лог список узлов плана обмена и остановится.

```
ERROR cdc_1c.change_reader: queue_guid is not set. Available nodes of exchange plan ДляODATA:
    a9bc23c5-3689-11f1-926c-0800270bc6cb  CDC  Витрина
```

Guid из первой колонки и есть искомый `queue_guid`.


## Запуск из окружения

Если хочется не писать код вовсе, есть готовый entrypoint — `python -m cdc_1c` (он же команда
`cdc-1c`). Он читает те же параметры из переменных окружения:

| Переменная | Обязательна | Значение |
|---|---|---|
| `CDC1C_ODATA_URL` | да | адрес OData-интерфейса базы 1С |
| `CDC1C_EXCHANGE_NAME` | да | имя плана обмена |
| `CDC1C_QUEUE_GUID` | да | `Ref_Key` узла обмена (очереди); не знаете — не задавайте, список узлов выведется в лог |
| `CDC1C_DB_URL` | да | строка подключения SQLAlchemy к целевой БД |
| `CDC1C_ODATA_USER` / `CDC1C_ODATA_PASSWORD` | нет | без пользователя запросы идут без авторизации |
| `CDC1C_DB_SCHEMA` | нет | схема целевой БД; не задана — схема по умолчанию |
| `CDC1C_DB_TEMP_SCHEMA` | нет | схема промежуточных таблиц merge; не задана — схема данных |
| `CDC1C_FULL_LOAD_WORKERS` | нет | число фоновых потоков полной выгрузки (по умолчанию 2) |
| `CDC1C_POLL_INTERVAL` | нет | период опроса в секундах (по умолчанию 60) |
| `CDC1C_MODE` | нет | `loop` (по умолчанию) или `once` |
| `CDC1C_LOG_LEVEL` | нет | уровень логирования (по умолчанию `INFO`) |

Свой код (обработчики, о них ниже) через переменные окружения не подключить — он объявляется кодом.
Для этого случая есть готовый шаблон точки входа: [config/](config/) — каталог с `runner.py`
и пакетом `handlers/` рядом. Запускается как обычный скрипт:

```bash
python config/runner.py
```

Как есть он поднимает только репликатор; обработчики и полные выгрузки по расписанию в нём
закомментированы как примеры — раскомментируйте нужное и замените на своё.

Python сам кладёт каталог скрипта в `sys.path`, поэтому `from handlers import ...` внутри `runner.py`
находит соседний пакет. Скопируйте этот каталог себе, положите туда свои обработчики — и то же самое
станет содержимым тома, монтируемого в контейнер (см. «Запуск в docker»).


## Запуск в docker

Образ — [`sobolevp/cdc-1c`](https://hub.docker.com/r/sobolevp/cdc-1c); версия образа совпадает
с версией пакета на PyPI. Внутри только python и cdc-1c: ни 1С, ни Postgres он не поднимает.

**Без своего кода.** Контейнер запускает тот же entrypoint `cdc-1c` и читает переменные из
таблицы выше. Первый запуск удобно сделать без `-d`, чтобы лог шёл прямо в терминал:

```bash
docker run --rm --network host \
  -e CDC1C_ODATA_URL="http://192.168.56.101/trade_demo/odata/standard.odata" \
  -e CDC1C_ODATA_USER=odata_user -e CDC1C_ODATA_PASSWORD=secret \
  -e CDC1C_EXCHANGE_NAME="ДляODATA" \
  -e CDC1C_QUEUE_GUID="a9bc23c5-3689-11f1-926c-0800270bc6cb" \
  -e CDC1C_DB_URL="postgresql+psycopg2://postgres:postgres@localhost:5432/cdc_1c" \
  -e CDC1C_DB_SCHEMA=cdc_1c -e CDC1C_DB_TEMP_SCHEMA=cdc_1c_tmp \
  sobolevp/cdc-1c:latest
```

Адреса и пароли здесь — из тестового контура, подставьте свои.

**`--network host` — про доступ к базе на той же машине.** У контейнера свой сетевой стек, поэтому
`localhost` внутри него означает сам контейнер, а не хост. Postgres, установленный на хосте, по
умолчанию слушает только `127.0.0.1`, то есть из контейнера недоступен ни как `localhost`, ни через
адрес docker-шлюза (`172.17.0.1`) — на том интерфейсе он ничего не слушает. `--network host` отдаёт
контейнеру сетевой стек хоста целиком: `localhost:5432` начинает означать вашу базу, а заодно
становятся видны все сети, которые видит хост (например, host-only сеть VirtualBox с 1С).
Альтернатива, если этот режим не подходит, — разрешить Postgres слушать docker-интерфейс
(`listen_addresses` + строка для `172.17.0.0/16` в `pg_hba.conf`) и подключаться на `172.17.0.1`;
доступ при этом открывается шире, поэтому для локальной работы проще первый вариант.

Когда БД и 1С стоят на других серверах, `--network host` не нужен вовсе: их адреса резолвятся
из контейнера как обычно.

Так контейнер работает постоянно: цикл опроса раз в `CDC1C_POLL_INTERVAL` секунд, недоступные
1С или БД не роняют его — ошибка пишется в лог, и попытка повторяется с нарастающей паузой.
Разово проверить доступы, не уходя в вечный цикл, можно с `-e CDC1C_MODE=once` — один цикл
read → save → notify, и выход.

Убедились, что работает — запускайте в фоне: `-d --name cdc-1c` вместо `--rm`. С `-d` docker
печатает только id контейнера и сразу возвращает управление, поэтому смотреть, что происходит,
надо в логах:

```bash
docker logs -f cdc-1c        # что делает; Ctrl-C выходит из просмотра, контейнер работает дальше
docker ps -a --filter name=cdc-1c   # если «ничего не произошло» — здесь видно, что контейнер вышел
```

**Со своим кодом.** Обработчики и расписания объявляются кодом, поэтому свой `runner.py` и пакет
`handlers/` монтируются в `/config`. Стартовый шаблон лежит в самом образе — репозиторий клонировать
не нужно:

```bash
docker run --rm sobolevp/cdc-1c:latest tar c -C /opt/cdc-1c config | tar x
# правим config/runner.py и config/handlers/, дальше:
docker run -d --name cdc-1c --env-file .env -v "$PWD/config:/config:ro" sobolevp/cdc-1c:latest
```

Если `/config/runner.py` есть — запускается он (и все параметры можно задавать прямо в нём,
литералами); если тома нет — работает env-режим выше. Файл названием отличается от `runner.py` —
укажите путь в `CDC1C_RUNNER`.

**Часовой пояс.** Расписания `FullLoadCron` считаются в локальном времени, а в контейнере это UTC:
задавайте `TZ` (например, `-e TZ=Europe/Moscow`).

**Остановка.** По `SIGTERM` циклы дорабатывают текущую итерацию и только потом выходят — так пакет
изменений не подтверждается в 1С раньше, чем записан в БД, и не остаются висеть незавершённые merge.
Дайте контейнеру на это время: в compose `stop_grace_period: 60s`, иначе docker добьёт процесс через
10 секунд по `SIGKILL`.

**docker compose.** Пример — [docker-compose.yml](docker-compose.yml) в репозитории:

```yaml
services:
  cdc-1c:
    image: sobolevp/cdc-1c:0.1.19
    restart: unless-stopped
    env_file: .env
    environment:
      TZ: Europe/Moscow
    volumes:
      - ./config:/config:ro
    stop_grace_period: 60s
```


## Режимы: `run_once` и `run_forever`

```python
rep.run_once()                 # один цикл: read → save → notify (подтверждение только после save)
rep.run_forever(interval=60)   # бесконечный цикл run_once с паузой; фоном — полные выгрузки
```

- `run_once(notify_changes=False)` — не подтверждать приём (пакет останется в очереди 1С; сделано для отладки).
- `run_forever(interval, max_iterations=0)` — `max_iterations>0` ограничивает число итераций.

## Полная (первоначальная) выгрузка

При работе `run_forever` объекты, впервые встреченные в пакете изменений, автоматически ставятся в
очередь на полную выгрузку и грузятся фоновыми потоками. Можно запустить выгрузку и вручную:

Полная выгрузка объекта реализована на стороне python, чтобы поддержать выгрузку объектов больших размеров, 
т.к. если инициировать полную выгрузку по плану обмена в 1С, то данные поступят без возможности постраничной загрузки.
(Поэтому данная функция специально убрана из модуля 1С).

Полная выгрузка спроектирована так, чтобы работать параллельно с получением изменений объекта.

```python
rep.list_objects()             # имена объектов 1С, доступных для выгрузки (Catalog_…, Document_…, …)

rep.full_load("Catalog_Номенклатура")              
rep.full_load("Document_РеализацияТоваровУслуг", batch_size=500)
```

`batch_size` — верхняя граница, а не жёсткий размер страницы: реальный размер подбирается по весу
выданных страниц, потому что одна запись 1С может тянуть за собой и одну строку, и тысячи (все
табличные части документа, весь набор движений регистратора). Как именно устроена постраничная
выгрузка и почему keyset-курсор в 1С неприменим к ссылочным ключам — см.
[DOCUMENTATION.md](DOCUMENTATION.md).


### Фильтр по периоду

Для ручной догрузки за нужный период укажите поле даты/времени и границы (включительно):

```python
from datetime import date, datetime

# весь месяц: date-граница включает последний день целиком (даже для поля дата-время)
rep.full_load("Document_РеализацияТоваровУслуг",
              date_field="Date", date_from=date(2026, 6, 1), date_to=date(2026, 6, 30))

# точная граница по времени — передайте datetime
rep.full_load("Document_РеализацияТоваровУслуг",
              date_field="Date", date_from=datetime(2026, 6, 1, 9, 0, 0))
```

`date_field` — имя поля 1С (`Date` у документов, `Period` у регистров). Границы транслируются в OData
`$filter` и объединяются с курсором пагинации.

### Пометка строк, которых в 1С больше нет

```python
rep.full_load("Catalog_Номенклатура", mark_missing=True)
```

Физическое удаление объекта (после «Удаления помеченных объектов») в обмен **не приходит вовсе**:
понятия «объект удалён» в пакете изменений нет, а у независимого регистра сведений нет даже
scoped-удаления — регистратора у него нет, и набор нечем ограничить. Такая строка иначе остаётся в
целевой таблице навсегда. Полная выгрузка — единственное место, где это видно: она читает объект
целиком и знает, чего в нём не оказалось.

Как это работает: каждая страница дописывает ключи прочитанных строк в одноразовую таблицу
`tmpkeys_<yymmddHHMMSS>_<таблица>_<hex8>` (та же схема и тот же разбор имени, что у промежуточных
таблиц merge), а после **успешного** завершения прогона строки, которых там нет, помечаются одним
`UPDATE`. Диапазоном ключей обойтись нельзя: 1С сортирует ссылочные ключи по представлению, а не по
guid.

- **Пометка, а не удаление.** `is_deleted_or_empty = true` и поднятый `merged_on`: обработчик
  замечает изменение только по `merged_on`, и физически удалённая строка не оставила бы витрине ни
  следа. Числовые ресурсы регистра гасятся в `NULL` — как и при выпадении строки из набора.
- **Гонка с изменениями.** Помечаются только строки старше старта прогона (тот же guard по
  `merged_on`, что и у самой выгрузки): строку, переписанную изменением уже во время прогона,
  снимок не трогает.
- **Выгрузка за период.** Отсутствие строки в окне ещё не значит удаления: у документа могла
  измениться дата, у независимого регистра — поле, входящее в ключ. Поэтому при заданном периоде
  каждый кандидат перед пометкой перепроверяется запросом в 1С по ключу (пачками, без фильтра по
  периоду), и найденные не помечаются.
- **Только успешный прогон.** Упавшая на середине выгрузка ничего не помечает — иначе «пропавшим»
  оказался бы весь непрочитанный хвост объекта.
- Выключено по умолчанию: прогон становится дороже (ключи + финальный анти-join), а нужно это не
  всем. У расписания флаг тот же: `FullLoadCron(..., mark_missing=True)`.

### По расписанию

Полная выгрузка — это ещё и проверка самого CDC: она читает объект из 1С и возвращает число
**реально** изменённых строк, то есть при исправном обмене отвечает нулём. Поэтому её имеет смысл
не звать руками, а повесить на расписание: ночная перегрузка свежего хвоста стоит недорого, а
расхождение показывает сразу — и тут же его выравнивает.

`FullLoadCron` — такой же вечный цикл, как `run_forever` у репликатора и обработчиков: на объект
приходится две строки, дальше он уходит своим потоком в тот же пул и останавливается тем же SIGTERM.

```python
from datetime import timedelta
from cdc_1c import FullLoadCron

# Каждую ночь в 03:00 — хвост за трое суток; timedelta считается в момент срабатывания,
# поэтому окно едет вместе с процессом.
zakazy_cron = FullLoadCron(replicator, "Document_ZakazKlienta", cron="0 3 * * *",
                           date_field="Date", date_from=timedelta(days=3))

# По воскресеньям в 02:00 — справочник целиком: границы не нужны вовсе.
nomenklatura_cron = FullLoadCron(replicator, "Catalog_Nomenklatura", cron="0 2 * * 0")

pool.submit(zakazy_cron.run_forever)
pool.submit(nomenklatura_cron.run_forever)
```

- **Имена — те, что видны в БД** (латиница), как и у обработчиков: имя таблицы и имя колонки. Их
  же показывает реестр `metadata_objects_1c` (`object_full_name_en` / `fields_en`). Оригинальные
  имена 1С (`"Document_ЗаказКлиента"`, `"ДатаОтгрузки"`) тоже принимаются.
- **Границы независимы**: нет `date_from` — с начала, нет `date_to` — до конца, нет обеих — объект
  целиком. `timedelta` в любой из них — смещение назад от текущей даты, то есть скользящее окно.
- **Расписание — обычная crontab-строка**, время локальное: в контейнере задавайте `TZ`.
- Прогон дольше периода: пропущенные срабатывания пропускаются, пачкой подряд не выстреливают.
  Если объект в этот момент уже выгружается (фоновая выгрузка репликатора или второе расписание),
  срабатывание пропускается с предупреждением в логе — 1С не получает двойную работу.
- Каждое расписание — отдельный поток, значит `pool_size` нужен на единицу больше за каждое
  (см. «Сколько нужно соединений к БД»).
- Проверить настройку, не дожидаясь трёх часов ночи: `zakazy_cron.run_once()`.

Готовый пример с расписаниями — закомментированный блок в [config/runner.py](config/runner.py):
раскомментируйте и замените объекты на свои.
Цикл изменений при этом не обязателен: `Replicator1C` нужен как исполнитель `full_load` (метаданные,
пагинация, запись, журнал), но его `run_forever` можно просто не запускать — получится процесс,
который только выгружает по расписанию. Метаданные он прочитает лениво первым же прогоном.

## Что появляется в целевой БД

- На каждый объект 1С — таблица (имя транслитерируется, длинные имена усекаются с хэшем под лимит СУБД).
- Служебные поля строк: `merged_on`/`inserted_on` (момент merge/первой вставки), `is_deleted_or_empty`
  (строку не учитывать, см. ниже), `exchange_message_no` (номер пакета обмена — диагностическое поле,
  логика загрузки на нём не построена).
- Служебные таблицы `replicator_1c_log`, `metadata_objects_1c` и `handlers_1c` (см. ниже).

`merged_on` — момент **последнего реального изменения** строки, а не последнего пакета, в котором она
приехала. 1С регулярно переписывает объекты, не меняя реквизитов; такие записи строку не трогают,
иначе инкрементальная материализация пересчитывала бы группы впустую. Поля, которые 1С меняет при
каждой записи (`exchange_message_no`, `DataVersion`), сами по себе изменением не считаются — они
записываются, только когда строку обновило что-то ещё.

### `is_deleted_or_empty` — универсальный признак «строку не учитывать»

Данные из 1С приходят инкрементально, поэтому строку нельзя просто выбросить: её исчезновение —
такое же событие, как изменение, и оно должно быть видно потребителям. Вместо удаления строка
остаётся с поднятым флагом.

Флаг сводит в одно булево поле все причины, по которым строку не следует учитывать в расчётах.
Их пять, и приходят они из разных мест:

| Причина | Откуда | Что означает |
|---|---|---|
| Пометка удаления | поле 1С `DeletionMark` | объект помечен на удаление в 1С |
| Неактивная запись | поле 1С `Active` = `false` | движение не участвует в итогах 1С |
| Строка выпала из набора | проставляется при merge | строки больше нет в наборе движений / табличной части |
| Пустой набор движений | запись сформирована при загрузке | все движения регистратора удалены |
| Опустевшая табличная часть | запись сформирована при загрузке | в табличной части не осталось строк |

Первые две — факты самой 1С, они просто переносятся во флаг.

Третья — про строки, пропавшие из набора. Набор движений регистратора и табличная часть приходят
целиком и целиком заменяют сохранённые: строки, которой в наборе больше нет, в 1С больше не
существует. Такая строка **не удаляется, а помечается**, и вместе с флагом ей поднимается
`merged_on`, а числовые ресурсы гасятся в `NULL`. Причина: инкрементальная материализация ищет
изменившиеся группы по `merged_on`, а исчезнувшая строка следа не оставляет — витрина навсегда
сохранила бы удалённое движение. Обнуление ресурсов — вторая линия обороны: `SUM` игнорирует `NULL`,
поэтому итог остаётся верным даже в запросе, забывшем фильтр по флагу. Остальные поля надгробия
сохраняются, так что видно, что это была за строка.

Последние две 1С сообщает отсутствием строк, а отсутствие в инкрементальный пакет не помещается:
чтобы «набор опустел» вообще доехало, формируется одна фиктивная запись с реальным ключом набора
(регистратор или `Ref_Key` владельца) и поднятым флагом. Она же вводит группу в пакет, без чего
пометка выпавших строк не сработала бы. Номер строки у неё — 1: как только набор снова наполнится,
первая настоящая строка перезапишет фиктивную, и та не осядет в таблице навсегда.

Флаг сбрасывается сам: если объект сняли с пометки удаления, запись снова стала активной или строка
вернулась в набор — приходит обычное изменение с `false`, и строка возвращается в расчёты вместе с
восстановленными значениями ресурсов.

Что это значит для запросов: **любой расчёт по сырым таблицам обязан учитывать флаг**, иначе в сумму
попадут удалённые, неактивные и выбывшие из наборов строки — они физически остаются в таблицах. В примерах материализации это сделано множителем
`* (NOT "is_deleted_or_empty")::int`, обнуляющим значение погашенной строки. Поля 1С `DeletionMark`
и `Active` при этом сохраняются как есть — если нужно различать причины, они рядом.

## Служебные таблицы

### Схема промежуточных таблиц (`db_temp_schema`)

Каждый merge заводит промежуточную (staging) таблицу, куда сначала складывается порция данных.
По умолчанию она создаётся в схеме данных; `db_temp_schema` уводит их в отдельную схему (её создаёт
сам `dbmerge`), и тогда таблицы с данными не перемешиваются с рабочими. Тот же параметр есть у
`HandlerLoop` — обработчик получает схему в `context.temp_schema` и передаёт её в свой `dbmerge`:

```python
handler = HandlerLoop(engine=engine, schema="cdc_1c", temp_schema="cdc_1c_tmp", handler=ZakazyKlientov())

# в обработчике
dbmerge(context.engine, table_name="ZakazyKlientov", schema=context.schema,
        temp_schema=context.temp_schema, ...)
```

Промежуточная таблица называется `tmp_<yymmddHHMMSS>_<таблица>_<hex8>` и после merge удаляется.
На Postgres она `UNLOGGED` (так быстрее), а `UNLOGGED` — это обычная постоянная таблица, поэтому
процесс, умерший посреди merge, оставляет её в базе. По имени такую понятно и опознать, и датировать:
Postgres времени создания таблиц не хранит, поэтому оно и вынесено в имя — по часам БД и первым
элементом, чтобы обычный список таблиц схемы сортировался по возрасту. Отдельная схема нужна ровно
за этим — в ней по определению нет ничего ценного, и разбирать такие остатки безопасно.

### `replicator_1c_log` — журнал загрузок

Строка на каждую загрузку объекта: пакет изменений или полная выгрузка.

| Колонка | Назначение |
|---|---|
| `id` | суррогатный ключ |
| `exchange` | имя плана обмена |
| `object` | имя объекта 1С |
| `type` | `changes` (пакет изменений) или `full` (полная выгрузка) |
| `message_no` | номер пакета обмена; `NULL` для полной выгрузки |
| `started_at` / `finished_at` | начало и конец загрузки; `finished_at IS NULL` — не завершена (упала) |
| `inserted_row_count` / `updated_row_count` / `deleted_row_count` | счётчики строк merge |
| `total_time` | суммарное время merge, сек |

Предназначено для мониторинга: незавершённые строки (`finished_at IS NULL`) — упавшие загрузки; по `type` и
`object` видно, что и когда грузилось.

### `metadata_objects_1c` — реестр объектов и состояние полной выгрузки

Синхронизируется с метаданными 1С — строка на каждый объект, встреченный в обмене. Ключ таблицы — полное
имя объекта (регистр и документ могут иметь одинаковое короткое имя).

| Колонка | Назначение |
|---|---|
| `object_full_name` | полное имя объекта 1С (ключ), например `Catalog_Номенклатура` |
| `object_full_name_en` | транслит = имя таблицы объекта в БД |
| `object_name` / `object_type` | имя и тип объекта (`Catalog` / `Document` / `AccumulationRegister` / …) |
| `fields` / `fields_en` | список полей объекта: имена 1С и их транслит (= колонки в БД) |
| `full_load_is_required` | объект ожидает полной выгрузки |
| `last_full_load_dt` | когда объект был полностью выгружен; `NULL` — ни разу |
| `last_full_load_rows_modified` | сколько строк выгрузка реально изменила |
| `last_full_load_minutes` | сколько она заняла, минуты (дробное) |
| `merged_on` | момент синхронизации записи реестра |

Новый объект оркестратор помечает `full_load_is_required=true`, фоновый воркер выгружает его целиком и
проставляет `last_full_load_dt`. Отсюда же удобно посмотреть список доступных объектов и имена их таблиц.

**`last_full_load_rows_modified` — это проверка самого CDC.** Полная выгрузка читает объект из 1С
целиком и сравнивает с тем, что уже лежит в БД. Если изменения доезжают исправно, менять ей нечего и
значение должно быть **0**. Ненулевое — значит часть изменений в обмен не попала, и стоит разобраться,
что именно: перевыгрузить объект и посмотреть, повторится ли.

```sql
SELECT object_full_name, last_full_load_dt, last_full_load_rows_modified, last_full_load_minutes
FROM cdc_1c.metadata_objects_1c
WHERE last_full_load_rows_modified > 0
ORDER BY last_full_load_rows_modified DESC;
```

### `handlers_1c` — состояние обработчиков

Строка на обработчика: подписка, отметка последнего успешного прогона, заказ пересборки. Создаётся
всегда, вместе с остальными служебными таблицами, — даже если своего кода вы не подключали: флаг
поднимать репликатор должен уметь сразу, а обработчик может появиться в другом процессе и позже.
Пустая таблица и означает «обработчиков нет». Колонки описаны
[ниже](#таблица-состояния-handlers_1c) вместе с самим механизмом.

### `writes_in_process_1c` — незавершённые записи

Строка живёт ровно столько, сколько идёт одна запись: появляется перед ней и исчезает после коммита.
Нужна обработчикам — по ней считается верхняя граница их окна (подробно
[ниже](#как-репликатор-зовёт-обработчика)). В покое таблица пуста; строки, висящие дольше нескольких
минут, означают либо очень толстую страницу полной выгрузки, либо умерший процесс.

| колонка | смысл |
|---|---|
| `id` | ключ строки: имя плана обмена + счётчик |
| `owner` | чья это запись — имя плана обмена (у обработчика `handler:<имя>`) |
| `object_name` | имя таблицы в БД, в которую идёт запись |
| `started_at` | момент старта по часам БД — то, к чему прижимается граница окна |
| `heartbeat_at` | отметка живости; не обновляется дольше 90 с — процесс считается умершим |

`heartbeat_at` отвечает на вопрос, на который `started_at` ответить не может: строка висит потому,
что merge долгий, — или потому, что процесс умер? Снаружи это неразличимо, а цена ошибки в обе
стороны высока. Признать живую запись мёртвой — перешагнуть её строки и потерять их молча; не
признать мёртвую — навсегда заморозить границу окна всем, кто читает эту таблицу (`kill -9`, OOM,
выселенный контейнер оставляют строку висеть). Поэтому владелец обновляет отметку раз в 20 секунд,
пока его записи в полёте, а читатели игнорируют строки старше 90 секунд: транзакции умершего
процесса СУБД уже откатила, держать по ним границу не за чем. Через час такие строки удаляются
совсем — владелец может не вернуться никогда.

## Логирование

Из коробки библиотека вешает вывод на логгер `cdc_1c` (INFO), если приложение не настроило логирование
само. Настроили своё — библиотека молчит и пишет через стандартный `logging`.

## Свой код по событию изменения: обработчики (handlers)

Витрина (или отправка изменений во внешнюю систему) сама не знает, что данные приехали: она либо
опрашивает БД вхолостую, либо ждёт, пока её запустят руками. Репликатор эту информацию имеет — он
и сохраняет данные, — поэтому он и сообщает, что пора считать.

Обработчик — класс, унаследованный от `Handler1C`. Готовые примеры — в
[config/](config/).

```python
from cdc_1c import Handler1C

class ZakazyKlientov(Handler1C):
    # имена ТАБЛИЦ в целевой БД, а не имена объектов 1С
    ON = ["AccumulationRegister_ZakazyKlientov", "Catalog_Nomenklatura"]
    ON_FULL_LOAD = True     # звать ли на страницах полной выгрузки (по умолчанию да)
    MIN_INTERVAL = 0        # не чаще раза в N секунд

    def setup(self, context):   # один раз за процесс: DDL вьюшек и целевых таблиц
        self.execute(context, DDL)

    def handle(self, context):  # полезная работа за окно
        ...
```

Подключается явным списком — никакого сканирования каталогов:

```python
from cdc_1c import HandlerLoop
from handlers import ZakazyKlientov, OtpravkaVOchered

handler_zakazy = HandlerLoop(engine=engine, schema="cdc_1c", handler=ZakazyKlientov())
handler_ochered = HandlerLoop(engine=engine, schema="cdc_1c",
                                handler=OtpravkaVOchered(queue="cdc"))

with ThreadPoolExecutor(max_workers=2) as pool:
    pool.submit(handler_zakazy.run_forever)
    pool.submit(handler_ochered.run_forever)
```

`run_forever` у обработчика блокирующий — ровно как у репликатора, поэтому запускаются они
одинаково, и где именно крутиться, решает точка входа. Останавливает всех `SIGTERM`; отдельный
цикл — `request_stop()`.

Репликатору их **не передают**: он о них не знает и знать не должен — общаются они через базу
(см. ниже). Поэтому запускать обработчиков можно где угодно: в одном процессе с репликатором, в
соседнем или в другом контейнере.

`HandlerLoop` — цикл **одного** обработчика: его состояние и его прогоны. По циклу на обработчика,
чтобы тяжёлая витрина не задерживала остальные.

Обратная сторона: **порядок между обработчиками не определён** — потоки независимы, и на «витрина
поверх витрины считается после базовой» полагаться нельзя. Два обработчика, пишущие в одну целевую
таблицу, тоже могут делать это одновременно.

Никакого особого каталога обработчикам не нужно: это обычные python-модули, поэтому лежат где
угодно, лишь бы импортировались, а общий код между ними подключается обычным `import`. Обновление
библиотеки этот код не трогает.

В `HandlerLoop` передаётся **экземпляр**, а не класс: так обработчик можно параметризовать конструктором
(одна логика на две схемы — два экземпляра). Имя (по умолчанию — имя класса) служит ключом состояния
в `handlers_1c`, поэтому у параметризованных экземпляров оно обязано различаться:
`ZakazyKlientov(name='ZakazyKlientov_mart2')`. Одноимённые отвергаются на старте.

Наследование не обязательно: достаточно `ON` и `handle` — годится и модуль целиком, и
функция с этими атрибутами. `Handler1C` даёт `setup()`, `since(context)`, `changed_since(context, *колонки)`,
`execute(context, sql)` и `query(context, sql)` — то, что иначе копируется из обработчика
в обработчик. В `sql` обоих последних `{schema}` подставляется целевой схемой.

В `ON` перечисляются имена **таблиц в целевой БД** (транслит), а не имена объектов 1С: обработчик
пишет SQL по таблицам, имя 1С он в глаза не видит. Имя 1С в `ON` не совпало бы ни с чем и обработчик
молча никогда бы не сработал — поэтому такой список отвергается на старте с подсказкой, как это имя
выглядит в базе.

Это не только «на что реагировать», но и «что я читаю»: по тому же списку считается верхняя граница
окна, поэтому перечислять надо **все** таблицы, из которых обработчик выбирает данные.

**Данные в обработчик не передаются** — он делает свой `SELECT`. Ему передаётся окно времени, за
которое надо отработать:

```sql
WHERE merged_on > :last_run_at
```

- `context.last_run_at` — отметка предыдущего **успешного** запуска; `NULL` — с начала времён;
- `context.boundary` — верхняя граница окна: то, что обработчик получит как `last_run_at` в следующий раз;
- `context.full_rebuild` — витрину просят собрать заново (см. ниже);
- `context.engine`, `context.schema`, `context.logger`;
- `context.objects` и `context.sources` — **только для логов**. Точности в них немного, и это
  осознанно: репликатор сообщает об изменении булевым флагом, а флаг не несёт ни имени таблицы, ни
  источника. Поэтому `objects` — это весь `ON` обработчика, а `sources` — `db_signal` (изменение),
  `startup` (первый проход процесса) или `full_rebuild` (заказ пересборки). Отличить по ним бэкфилл
  полной выгрузки от живого изменения нельзя — для этого есть `ON_FULL_LOAD`.

Обе границы — по часам БД (тем же `now()`, которым `dbmerge` штампует `merged_on`), поэтому
расхождение часов между хостами роли не играет.

В `WHERE` верхняя граница не нужна, хотя окно ею закрывается. От пропуска строк защищает не условие
выборки, а само значение `context.boundary`: оно прижато к старту незавершённого merge, чьи строки
обработчику всё равно не видны. Видимую строку правее границы обработчик посчитает раньше срока — и
посчитает ещё раз в следующем окне, потому что `last_run_at` станет `boundary`. Это лишняя работа,
а не пропуск, и взамен витрина получается свежее. Добавить `merged_on <= :boundary` имеет смысл
только там, где повтор дорог сам по себе — например, при отправке во внешнюю систему.

### Когда обработчик зовут

Только когда merge **реально что-то изменил** (вставил, обновил или пометил удалённой хоть одну
строку). 1С регистрирует изменение объекта на любую перезапись, и в пакет приезжает масса записей,
идентичных тому, что уже лежит в БД; шумные поля (`DataVersion`, номер сообщения обмена) при
сравнении не учитываются. Звать на таком пакете незачем — `SELECT` по окну всё равно вернёт пусто.

Полная выгрузка поднимает флаг **на каждую сохранённую страницу**, а не один раз в конце. Лишних
вызовов это не даёт: флаг булев, тысяча страниц поднимет его один раз, — зато витрина начинает
наполняться после первой же страницы, а не через часы. Обработчику, которому бэкфилл не нужен
(рассылка уведомлений, отправка в очередь), ставьте `ON_FULL_LOAD = False`: это значение он
объявляет в `handlers_1c`, и репликатор на страницах выгрузки его просто не трогает.

По той же причине DDL живёт в `setup()`, а не в `handle()`: `setup` вызывается один раз за процесс,
иначе `CREATE OR REPLACE VIEW` выполнялся бы на каждую страницу выгрузки.

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

### Сколько нужно соединений к БД

Потоков несколько сортов, и пулы у них раздельные: цикл изменений (главный поток),
`full_load_workers` потоков полной выгрузки, поток отметки живости незавершённых merge, по потоку
на каждого обработчика и по потоку на каждое расписание (`FullLoadCron` пишет страницы сам, а не
через пул `full_load_workers`). Занять чужие потоки они не могут, но `engine` у них общий, и
одновременно держать соединение могут все сразу:

```
pool_size >= full_load_workers + 2 + число обработчиков + число расписаний в этом процессе
```

Если обработчики вынесены в отдельный процесс, у каждого процесса свой `engine` и свой счёт. Не
хватает соединений — кто-то встаёт в ожидание, и затормозить может как обработчик, так и сама
выгрузка.

### Несколько планов обмена

Планов обмена может быть несколько — тогда на каждый заводится свой `Replicator1C`. Специально
делать для этого ничего не нужно: обработчику безразлично, сколько планов его кормит и в одном ли
они с ним процессе. Все репликаторы поднимают тот же флаг в `handlers_1c` и публикуют свои
незавершённые merge в общую `writes_in_process_1c`.

Если репликаторы живут в одном процессе, помните: `run_forever` блокирует поток, так что второй
запускается в своём. В [runner.py](config/runner.py) для этого достаточно завести второй
`Replicator1C` со своими `exchange_name` и `queue_guid` и отправить его `run_forever` в тот же пул —
больше ничего. `SIGTERM`/`SIGINT` останавливают **все** циклы сразу, в каком бы потоке они ни
крутились; отдельный цикл можно остановить программно через `rep.request_stop()`.

Пул считается по сумме: `sum(full_load_workers) + 2 × число репликаторов + число обработчиков +
число расписаний`.

### Как репликатор зовёт обработчика

Через базу, и только через неё. Обработчик при старте объявляет себя в `handlers_1c` и записывает
в `update_on` таблицы, на которые подписан. Репликатор эту колонку читает и, увидев изменение
подписанной таблицы, поднимает `update_is_required` — а дальше обработчик сам замечает флаг своим
циклом и снимает его успешным прогоном.

```
репликатор            handlers_1c              обработчик
──────────            ───────────              ──────────
сохранил объект  →  update_is_required=true  →  увидел флаг, посчитал витрину,
                                                снял флаг и записал last_run_at
```

Отсюда главное следствие: **репликатору не нужны ни объекты обработчиков, ни их импорт**. Их можно
поднять отдельным процессом или вообще в другом контейнере — связывает всех одна общая БД. Такой
процесс получается из [runner.py](config/runner.py) вычёркиванием: убираете `Replicator1C`
и его строку в пуле, остаются одни `HandlerLoop`. В самих обработчиках не меняется ничего.

Верхняя граница окна при этом тоже берётся из базы — из таблицы `writes_in_process_1c`, куда
репликатор пишет свои незавершённые merge. Обработчик обязан их видеть, чей бы процесс их ни начал:
их строки уже имеют `merged_on` в прошлом, но ещё не видны, и граница, взятая как «сейчас», их бы
перешагнула. Строки процесса, который умер, отсекаются по отметке живости — она обновляется, пока
репликатор жив, поэтому его падение обработчиков не морозит.

### Таблица состояния `handlers_1c`

| колонка | смысл |
|---|---|
| `name` | `NAME`, либо имя класса обработчика |
| `update_on` | таблицы, на которые подписан; заполняет обработчик, читает репликатор |
| `on_full_load` | звать ли на страницах полной выгрузки; тоже объявляет обработчик |
| `update_is_required` | репликатор увидел изменение подписанной таблицы |
| `enabled` | выключенный обработчик не зовут (и окно за время простоя не копят) |
| `last_run_at` | граница последнего успешного запуска |
| `last_error` | traceback последнего падения |
| `full_rebuild_is_required` | заказ на пересборку витрины |
| `rebuild_cursor` | метка последнего завершённого блока идущей пересборки; `NULL` — пересборка не идёт |
| `last_full_rebuild_dt` | когда пересборка отработала в последний раз |
| `last_full_rebuild_minutes` | сколько она заняла, минуты (дробное) |

Попросить собрать витрину заново:

```sql
UPDATE cdc_1c.handlers_1c SET full_rebuild_is_required = true WHERE name = 'ZakazyKlientov';
```

Заказ сам ставит обработчик в очередь — ждать изменений по подписанным объектам не нужно. Он
получит окно с начала времён (`context.last_run_at` = `None`) и `context.full_rebuild` = `True`.
Первый в жизни прогон — тоже пересборка, флаг для этого ставить не нужно.

После успеха требование снимается, а в `last_full_rebuild_dt` записывается время.

#### Пересборка по блокам

Пересборка большой витрины идёт десятками минут, и всё это время изменения не применяются — витрина
стоит холодной. Чтобы этого не было, пересборку можно нарезать на блоки: объявите `rebuild` —
генератор, где `yield <метка>` означает «блок закончен и закоммичен». В этих точках цикл вклинивается
и применяет накопившиеся изменения.

```python
def rebuild(self, context):
    for year in self.years(context):
        if context.rebuild_from and str(year) <= context.rebuild_from:
            continue                      # этот год уже посчитан до перезапуска процесса
        self.merge_year(context, year)
        yield str(year)
```

Чем нарезать — решает обработчик: период, диапазон ключей, склад, организация. Библиотека в блоки не
заглядывает, метка нужна ей только чтобы записать её в `handlers_1c.rebuild_cursor` и вернуть в
`context.rebuild_from` после перезапуска — генератор живёт в памяти и рестарта не переживает, а
начинать сорокаминутную пересборку заново из-за него не хочется.

Не объявили `rebuild` — прежнее поведение: один проход, весь `handle()` целиком.

**Почему это безопасно и почему это не параллельность.** Блок и инкремент выполняет один и тот же
поток, поэтому одновременно они не работают никогда — затирать друг друга им нечем. Считай
пересборка в своём потоке, случилось бы именно это: она читает данные в момент старта, а коммитит
десятки минут спустя, и её `delete_condition` снёс бы группу целиком вместе со свежим результатом
инкремента. Причём **безвозвратно**: инкремент уже сдвинул `last_run_at` за это изменение, и группу
никто не пересчитает, пока источник не изменится снова. Платим за безопасность задержкой в один
блок вместо задержки во всю пересборку.

Отсюда единственное требование к блоку: он должен читать данные **актуальные на момент своего
выполнения**, а не снимок на старте пересборки. Обработчики так и написаны (верхней границы в
`WHERE` у них нет), поэтому специально делать ничего не нужно — но «прочитаю всё один раз в начале,
разложу по блокам потом» сломает ровно это.

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

```python
merge.exec(source_condition=source_period == period,
           delete_condition=target_period == period)
```

Заодно это само чинит переезд между блоками (поправили дату документа): из старого месяца строку
удалит его блок, в новом создаст его собственный.

Отсюда требование к колонке, по которой нарезано: она должна быть **в самой витрине** и у каждой её
строки принимать одно значение. В [zakazy_klientov_grouped.py](config/handlers/zakazy_klientov_grouped.py)
это видно хорошо — витрина агрегатная, поэтому месяц берётся из даты документа (из неё же берётся
`Year` в ключе группы, значит у группы он ровно один), а не из периода движения, которых у группы
может быть несколько.

Список блоков берите из источника **и из витрины**: месяц, которого в источнике больше нет, а в
витрине он есть, иначе не очистится никогда — блок по нему просто не запустится.

Метку блока удобно делать сортируемой строкой (`YYYY-MM`) — точка возобновления сравнивается именно
как строка, и «всё, что меньше метки» должно быть уже сделано.

`last_run_at` пересборка не обнуляет: если она упадёт, прежняя граница останется на месте, а
незаконченная пересборка продолжится с `rebuild_cursor`.

Тот же заказ репликатор ставит сам, когда в таблице объекта появляется **новая колонка** (новый реквизит
в 1С). Инкремент её не увидел бы никогда: окно строится по `merged_on`, а `merged_on` двигается
только у строк, у которых изменились значения — добавление колонки не меняет ни одного значения, и
все уже лежащие строки остались бы левее окна навсегда.

### Что нужно помнить

- Доставка **at-least-once**: упали после работы обработчика, но до записи `last_run_at` — окно
  повторится. Для витрины это безразлично (merge идемпотентен), для отправки во внешнюю систему
  потребитель должен быть идемпотентным.
- Упавший обработчик границу не двигает и остаётся «грязным» — повтор произойдёт сам, без нового
  изменения.
- Удаления в окно попадают: из целевых таблиц ничего не исчезает физически (у справочников и
  документов пометка приезжает из 1С, у регистров выпавшие из набора строки помечаются
  `is_deleted_or_empty` с обнулением ресурсов), а пометка — это `UPDATE`, который двигает
  `merged_on`. Поэтому `SELECT` обработчика **не должен** отфильтровывать `is_deleted_or_empty`:
  для витрины это строка с нулём, для очереди — событие удаления.

## Дальнейшая материализация и сборка денормализованных таблиц

1С хранит данные в нормализованном виде: чтобы дотянуться, например, из регистра заказов до кода
товара, нужен `JOIN` со справочником номенклатуры по guid. Для задач DWH обычно нужна менее строгая
нормализация, поэтому в [config/handlers/](config/handlers/) лежат два разобранных
примера инкрементальной витрины:

- [zakazy_klientov.py](config/handlers/zakazy_klientov.py) — ключ таблицы фактов сохраняется
  (строки регистра как есть, плюс артикул из справочника);
- [zakazy_klientov_grouped.py](config/handlers/zakazy_klientov_grouped.py) — ключ меняется
  (`GROUP BY` по номеру документа, году и артикулу), пересчёт остаётся инкрементальным, а ключ
  группы составной — видно, как это отражается на фильтрах (`tuple_(a, b).in_(...)`).

Оба построены на одном правиле: **вьюшка отдаёт `merged_on` каждого участника `JOIN` отдельной
колонкой** (`merged_on`, `Nomenklatura_merged_on`, …), и в инкремент попадает всё, что стало свежее
хотя бы по одному источнику.

Так надо потому, что объекты 1С приезжают в обмене независимо и в произвольном порядке: справочник
номенклатуры может доехать (или измениться) позже регистра. Если ориентироваться только на
`merged_on` регистра, такая строка уже не попадёт в инкремент и витрина навсегда останется с `NULL`
или старым артикулом. Отдельные колонки заодно не смешивают несвязанные «часы» в одну.

Собирать эти отметки в один `WHERE ... OR ...` на объёме нельзя, и дело не в отсутствии индексов:
ветки `OR` живут на **разных** таблицах соединения, поэтому ни по одной нельзя отфильтровать до
`JOIN` — условие проверяется на готовой паре и вырождается в `Join Filter` поверх полного соединения.
`BitmapOr` из индексных сканов Postgres строит только когда все ветки `OR` на одной таблице, а
переписать `OR` в `UNION` через границу джойна он не умеет (тип соединения ни при чём: `INNER`
вместо `LEFT` плана не меняет).

Поэтому «какие группы изменились» спрашивается через `UNION ALL` по той же вьюшке — по ветке на
отметку:

```sql
SELECT "Number", "Year" FROM "..._rows_view" WHERE "merged_on" > :since
UNION ALL
SELECT "Number", "Year" FROM "..._rows_view" WHERE "ZakazKlienta_merged_on" > :since
UNION ALL
SELECT "Number", "Year" FROM "..._rows_view" WHERE "Nomenklatura_merged_on" > :since
```

В каждой ветке остаётся предикат ровно по одной базовой таблице, планировщик опускает его внутрь
вьюшки, в скан этой таблицы, и стоимость начинает зависеть от размера окна, а не от размера таблиц.
Переписывать `JOIN`-ы руками не нужно — логика соединений остаётся в одном месте. На синтетике
(регистр 2 млн строк, окно ~0.5%): `OR` — 620 мс с полным сканом регистра, `UNION ALL` — 98 мс без
единого `Seq Scan`.

Двух вещей это требует: `random_page_cost` под SSD (`1.1` вместо дефолтных `4` — иначе планировщик
всё равно предпочтёт полный скан, и выигрыш будет вдвое скромнее) и индекса по колонке соединения
со справочником: репликатор создаёт только индексы `merged_on`, а что с чем соединяется, знает
витрина.

Пересчитывать при этом нужно **группу целиком** (тот же ключ, по которому идёт удаление), а не
отдельные «свежие» строки — иначе `delete_condition` снесёт соседние строки группы, которые в
инкремент не попали.

Хранить границу обработки в самой витрине (по отметке на источник) больше не нужно — она приходит
готовой в `context.last_run_at`, одна на все источники сразу.

