Вступить в клуб →
сложнаявопросСлои DWH и моделирование витрин

Инкремент по 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: эта колонка проставляется один раз и больше не меняется, поэтому граница окна монотонна и пропусков не будет.

🔒 Проверка ответа — для участников клуба

  • Проверка ответа
  • Подсказка, если застряли
  • Разбор с объяснением, почему так
  • Прогресс по всем задачам и виртуальные собеседования
Зарегистрироваться →

Регистрация занимает минуту

Другие задачи раздела