Загрузка выбирает записи по времени изменения, но у нескольких записей одинаковая метка времени. Как не пропустить записи между запусками?
Нужно хранить составной курсор, например пару время изменения плюс стабильный уникальный идентификатор, и упорядочивать записи по этой паре. Следующий запуск выбирает записи строго после последнего успешно обработанного курсора, а сам курсор обновляет только после успешной фиксации загрузки.
Полная перегрузка источника при каждом запуске проста, но становится дорогой по мере роста объёма данных. Поэтому пакетные процессы перешли к инкрементальной загрузке, когда между запусками извлекается только новая или изменившаяся часть данных.
Наивный вариант использует одну метку времени и сохраняет максимальное увиденное значение. Однако временная метка часто имеет недостаточную точность: несколько записей могут получить одинаковое значение. Это создаёт необходимость в детерминированном порядке внутри одной временной точки.
Предположим, загрузка обработала часть записей с меткой времени 12:00 и сохранила это время как позицию продолжения. Если следующий запуск ищет только записи с временем строго больше 12:00, остальные записи с той же меткой будут пропущены.
Если же следующий запуск повторно читает все записи с временем 12:00, система должна выдерживать повторную обработку. Без этого возможны дубли, неполные витрины или расхождение между источником и приёмником.
Курсор должен включать время изменения и уникальный упорядочиваемый идентификатор записи. Записи сортируются по этим двум полям, а следующий запуск выбирает те, чья пара лексикографически больше последнего сохранённого курсора: сначала сравнивается время, затем идентификатор.
Например, после обработки пары 12:00 и 105 следующий запуск должен взять записи с временем позже 12:00 либо записи с временем 12:00 и идентификатором больше 105. Поэтому одинаковые временные метки больше не приводят к потере данных.
Курсор нужно сохранять атомарно вместе с результатом загрузки либо после подтверждения успешной записи результата. Если сохранить его раньше, сбой может привести к тому, что источник уже будет считаться обработанным, хотя данные фактически не попали в приёмник.
Стабильный идентификатор должен однозначно различать записи и сохранять порядок при повторных чтениях. Если записи могут изменяться после первоначальной загрузки, одного курсора недостаточно: механизм должен снова увидеть такие изменения, например через надёжное поле версии, CDC или период повторного чтения с последующей идемпотентной записью.
Временные метки могут зависеть от часов разных узлов, иметь ограниченную точность или изменяться назад. Поэтому при отсутствии надёжного порядка лучше использовать CDC с позициями журнала. Компромиссный вариант — перекрывающееся окно чтения: он повышает устойчивость к задержкам, но требует дедупликации и увеличивает нагрузку на источник.
В интернет-магазине пакетная загрузка каждые пять минут переносила изменённые заказы по полю времени обновления. Несколько заказов обновлялись одной операцией и получали одинаковую временную метку. Процесс сохранял только время последней обработки, поэтому часть заказов периодически не попадала в аналитическую витрину.
Рассматривались три варианта. Полная загрузка устраняла риск пропуска, но была слишком дорогой. Перекрытие интервалов уменьшало риск потери, но создавало повторы и дополнительную нагрузку. Переход на составной курсор из времени и идентификатора устранил пропуски при сохранении эффективности; в приёмнике также использовалась идемпотентная обработка для безопасных повторов.
После этого граница между запусками стала однозначной, а курсор начал продвигаться только после успешной фиксации данных. Для изменений, которые могли приходить с задержкой, команда дополнительно выбрала CDC, поскольку обычное время обновления не гарантировало полного порядка событий.
Нет, если идентификаторы не отражают порядок изменений. Больший идентификатор может относиться к старой записи, а обновление записи с меньшим идентификатором произойти позже. Для инкрементальной выборки нужен признак изменения, а идентификатор используется как детерминирующий компонент внутри одинакового значения этого признака.
Такой подход действительно предотвращает пропуски на границе, но повторно читает все записи с последней временной меткой. Это допустимо только при дедупликации и идемпотентной записи в приёмник. Составной курсор позволяет читать строго следующий диапазон и уменьшает число повторов без потери полноты.
Если курсор обновляется независимо от результата, последующий запуск может перескочить через ещё не записанные данные. Поэтому запись данных и продвижение курсора должны быть согласованы, либо повторный запуск должен быть безопасным: данные загружаются через временную область, затем публикуются атомарно, а целевая запись выполняется идемпотентно по ключу события или версии.