Инкремент по updated_at теряет строки из-за поздно приходящих данных
Витрина грузится инкрементом: where updated_at > (select max(updated_at) from dm.orders_mart). Источник — реплика с плавающим лагом, часть транзакций доезжает до реплики с задержкой в несколько минут после того, как их updated_at уже проставлен.
При сверке выяснилось: каждый день безвозвратно теряется около 0.1% строк — их updated_at меньше уже сохранённого максимума, поэтому в следующий инкремент они не попадают. Полностью перезаливать витрину нельзя: 1.5 млрд строк, окно загрузки не позволяет.
Какое решение закроет потери?
- Запускать тот же DAG чаще, каждые 15 минут вместо часа: окно загрузки станет короче лага реплики, и опоздавшие строки попадут уже в следующий запуск.
- Грузить с перекрытием окна — брать данные за последние N часов от границы, а не строго после максимума, и сделать загрузку идемпотентной по ключу, чтобы перекрытие не давало дублей.
- Заменить строгое сравнение на нестрогое: updated_at >= (select max(updated_at) from dm.orders_mart) — граничные строки перестанут пропадать, а возникшие дубли уберёт DISTINCT в витрине.
- Перейти на фильтр по created_at вместо updated_at: эта колонка проставляется один раз и больше не меняется, поэтому граница окна монотонна и пропусков не будет.
