Загрузка выбирает записи по времени изменения, но у нескольких записей одинаковая метка времени. Как не про...

Загрузка выбирает записи по времени изменения, но у нескольких записей одинаковая метка времени. Как не пропустить записи между запусками?

Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

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

Исторический контекст

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

Наивный вариант использует одну метку времени и сохраняет максимальное увиденное значение. Однако временная метка часто имеет недостаточную точность: несколько записей могут получить одинаковое значение. Это создаёт необходимость в детерминированном порядке внутри одной временной точки.

Постановка проблемы

Предположим, загрузка обработала часть записей с меткой времени 12:00 и сохранила это время как позицию продолжения. Если следующий запуск ищет только записи с временем строго больше 12:00, остальные записи с той же меткой будут пропущены.

Если же следующий запуск повторно читает все записи с временем 12:00, система должна выдерживать повторную обработку. Без этого возможны дубли, неполные витрины или расхождение между источником и приёмником.

Подробное решение

Курсор должен включать время изменения и уникальный упорядочиваемый идентификатор записи. Записи сортируются по этим двум полям, а следующий запуск выбирает те, чья пара лексикографически больше последнего сохранённого курсора: сначала сравнивается время, затем идентификатор.

Например, после обработки пары 12:00 и 105 следующий запуск должен взять записи с временем позже 12:00 либо записи с временем 12:00 и идентификатором больше 105. Поэтому одинаковые временные метки больше не приводят к потере данных.

Курсор нужно сохранять атомарно вместе с результатом загрузки либо после подтверждения успешной записи результата. Если сохранить его раньше, сбой может привести к тому, что источник уже будет считаться обработанным, хотя данные фактически не попали в приёмник.

Стабильный идентификатор должен однозначно различать записи и сохранять порядок при повторных чтениях. Если записи могут изменяться после первоначальной загрузки, одного курсора недостаточно: механизм должен снова увидеть такие изменения, например через надёжное поле версии, CDC или период повторного чтения с последующей идемпотентной записью.

Временные метки могут зависеть от часов разных узлов, иметь ограниченную точность или изменяться назад. Поэтому при отсутствии надёжного порядка лучше использовать CDC с позициями журнала. Компромиссный вариант — перекрывающееся окно чтения: он повышает устойчивость к задержкам, но требует дедупликации и увеличивает нагрузку на источник.

Ситуация из практики

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

Рассматривались три варианта. Полная загрузка устраняла риск пропуска, но была слишком дорогой. Перекрытие интервалов уменьшало риск потери, но создавало повторы и дополнительную нагрузку. Переход на составной курсор из времени и идентификатора устранил пропуски при сохранении эффективности; в приёмнике также использовалась идемпотентная обработка для безопасных повторов.

После этого граница между запусками стала однозначной, а курсор начал продвигаться только после успешной фиксации данных. Для изменений, которые могли приходить с задержкой, команда дополнительно выбрала CDC, поскольку обычное время обновления не гарантировало полного порядка событий.

Что кандидаты часто упускают

  1. Достаточно ли использовать только максимальный идентификатор записи?

Нет, если идентификаторы не отражают порядок изменений. Больший идентификатор может относиться к старой записи, а обновление записи с меньшим идентификатором произойти позже. Для инкрементальной выборки нужен признак изменения, а идентификатор используется как детерминирующий компонент внутри одинакового значения этого признака.

  1. Почему нельзя просто использовать условие больше или равно для одной временной метки?

Такой подход действительно предотвращает пропуски на границе, но повторно читает все записи с последней временной меткой. Это допустимо только при дедупликации и идемпотентной записи в приёмник. Составной курсор позволяет читать строго следующий диапазон и уменьшает число повторов без потери полноты.

  1. Что произойдёт, если загрузка упадёт после записи части данных?

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