Time Travel: запросы к историческим снапшотам и откат таблицы
Time Travel в Apache Iceberg: VERSION AS OF и TIMESTAMP AS OF в Spark SQL, snapshot-id и as-of-timestamp в PySpark DataFrameReader, механика rollback_to_snapshot и set_current_snapshot под капотом metadata.json, горизонт хранения истории и стоимость retention.
Зачем нужен Time Travel: от ручного disaster recovery к запросу за секунды¶
Представьте классическую Hive-таблицу на HDFS или S3, которую обслуживает ежедневный batch ETL. Инженер деплоит новую версию пайплайна, в коде закралась ошибка в условии JOIN - она задваивает часть строк. Пайплайн отрабатывает штатно, без единого исключения, и перезаписывает партицию вчерашнего дня через INSERT OVERWRITE. Ошибку замечают не сразу - часто только когда аналитик или бизнес-пользователь смотрит на дашборд и видит, что выручка за вчера каким-то образом превысила выручку за весь предыдущий месяц. К этому моменту старая, корректная версия партиции физически уже не существует: она была удалена в момент перезаписи, а единственный способ её вернуть - поднять резервную копию (если она настроена и достаточно свежая) или заново прогнать ETL от исходных данных в источнике (если источник всё ещё хранит нужный диапазон и ETL детерминирован). Оба пути занимают часы, а не секунды, и оба требуют участия дежурного инженера, который должен сначала понять, что случилось, потом найти нужный бэкап или пересчитать нужный диапазон, и только потом восстановить таблицу.
Apache Iceberg меняет эту картину радикально, и происходит это не благодаря отдельной системе бэкапов, а благодаря архитектуре, с которой вы уже знакомы из второго урока модуля. Snapshot Model устроена так, что каждая операция записи - INSERT, UPDATE, DELETE, MERGE INTO - не изменяет существующие файлы данных, а создаёт новый снапшот, который ссылается на новый набор файлов через новое дерево манифестов. Старые файлы данных и старые манифесты, описывающие предыдущее состояние таблицы, никуда не исчезают в момент коммита - они продолжают физически существовать в object storage до тех пор, пока отдельная процедура обслуживания (expire_snapshots, тема следующего урока) не примет явное решение их удалить. Из этого простого факта прямо следует Time Travel - возможность прочитать таблицу такой, какой она была на момент любого из сохранённых снапшотов, без копирования данных, без восстановления из бэкапа и без повторного запуска ETL. Технически это превращается в задачу «найти нужную запись в дереве метаданных и направить движок читать по ней» - операцию, которая по стоимости не зависит от объёма данных в таблице и завершается за секунды независимо от того, 10 гигабайт там или 10 петабайт.
Этот урок продолжает прямую линию модуля: шестой и седьмой уроки разобрали Copy-on-Write и Merge-on-Read как механизмы физической записи, восьмой урок показал, как MERGE INTO использует обе стратегии для production-паттерна CDC-репликации. Каждая такая операция оставляет за собой снапшот. Time Travel - это инструмент, который превращает этот побочный продукт записи (историю снапшотов) в первоклассную пользовательскую возможность: читать прошлое, сравнивать версии и - при необходимости - откатывать таблицу назад.
Зачем Time Travel бизнесу¶
С точки зрения бизнес-пользователей и аналитиков Time Travel решает три различных, но смежных класса задач.
Аудит изменений данных. Финансовый аналитик должен иметь возможность объяснить, почему отчёт, сгенерированный сегодня, отличается от отчёта, сгенерированного неделю назад, даже если оба запроса синтаксически идентичны и обращаются к «той же» таблице. Без Time Travel единственный способ ответить на этот вопрос - смотреть в журналы ETL-пайплайна и пытаться реконструировать происходившее по логам. С Time Travel ответ получается прямым SQL-запросом: достаточно прочитать таблицу TIMESTAMP AS OF нужной даты и сравнить результат с текущим состоянием.
Воспроизводимость ML-экспериментов. Модель машинного обучения, обученная на каком-то срезе данных, должна допускать повторное воспроизведение того же эксперимента - в том числе спустя месяцы, когда исходная таблица уже претерпела десятки новых загрузок и пересчётов. Если обучающий датасет был материализован как точечный снимок (CSV-экспорт, отдельная Parquet-копия), это решает проблему, но дублирует данные и создаёт собственную проблему устаревания дубликата. Time Travel позволяет вместо этого зафиксировать не сами данные, а ссылку на конкретный снапшот исходной таблицы - и в любой момент в будущем восстановить ровно тот же тренировочный набор, прочитав ту же таблицу с тем же snapshot-id, без отдельного хранения копии.
Сравнение отчётов «тогда и сейчас». Продуктовая аналитика часто сводится к вопросу «что изменилось с прошлой недели» - в данных, в метриках, в распределении значений по сегментам. Time Travel превращает этот вопрос в JOIN между двумя версиями одной и той же таблицы, прочитанными на разные моменты времени - конкретный пример такого запроса показан позже в этом уроке, в разделе про синтаксис Spark SQL.
Зачем Time Travel инженерам: disaster recovery без отдельной системы бэкапов¶
Со стороны инженерии эксплуатации Time Travel - это прежде всего ответ на человеческий фактор и баги в production-коде. Ошибочный MERGE INTO с перепутанным условием ON, UPDATE без полноценного WHERE, случайный DROP TABLE по неверно скопированному имени, баг в логике дедупликации CDC-потока, который пропустил WHERE rn = 1 из восьмого урока модуля и задвоил половину заказов - все эти инциденты в Hive-парадигме означают, что данные физически испорчены, и единственный путь назад - восстановление из бэкапа или пересчёт. В парадигме Iceberg каждая из этих катастроф - это просто очередной снапшот в цепочке истории, который можно «обойти», вернувшись к снапшоту-предшественнику.
Важно подчеркнуть формулировку точнее: Time Travel и операция отката (rollback, разобрана в отдельном разделе ниже) не «отменяют» ошибочную операцию в смысле undo - они дают доступ к состоянию, которое существовало до неё, и позволяют либо прочитать это состояние для расследования, либо официально сделать его текущим состоянием таблицы снова. Сама ошибочная операция и её снапшот при этом не удаляются из истории - они остаются доступными для дальнейшего разбора (например, чтобы понять, какие именно строки были испорчены, для последующего уведомления заинтересованных сторон).
Диаграмма выше специально устроена симметрично: первые три шага в обеих ветках идентичны - баг есть баг, и сам момент его внесения и обнаружения Time Travel никак не ускоряет и не предотвращает. Разница начинается на четвёртом шаге. В Hive-ветке четвёртый шаг требует внешнего ресурса (резервная копия) или повторного дорогого вычисления (пересчёт ETL), и оба варианта измеряются часами и требуют ручной оркестрации: найти нужный бэкап, поднять его, подтвердить целостность, переключить таблицу. В Iceberg-ветке четвёртый шаг - это обычный SELECT к системной таблице history, возвращающий список снапшотов с их made_current_at и snapshot_id; пятый шаг - вызов одной хранимой процедуры, которая переключает указатель current-snapshot-id в metadata.json, не трогая ни единого байта данных. Именно это различие - «искать внешний ресурс и пересчитывать» против «прочитать существующие метаданные и переключить указатель» - и есть архитектурная суть того, почему Time Travel в Iceberg практически бесплатен по производительности: к моменту, когда он нужен, все необходимые для него данные уже лежат на месте как побочный продукт нормальной работы snapshot-модели.
Напоминание: почему это «бесплатная» фича¶
Стоит явно подчеркнуть связь с материалом второго урока модуля, а не повторять его с нуля. Снапшот в Iceberg - это неизменяемый снимок состояния таблицы, состоящий из ссылки на manifest list, которая в свою очередь ссылается на manifest files, а они - на конкретные data files. Когда MERGE INTO или UPDATE переписывает файл (Copy-on-Write) или добавляет delete-файл (Merge-on-Read), предыдущий снапшот продолжает указывать на старый manifest list, который продолжает указывать на старые, физически не удалённые файлы. Time Travel не создаёт никакой новой инфраструктуры для хранения истории - он просто читает уже существующий граф снапшотов и манифестов, который в любом случае строится при каждой записи. Цена этой фичи - не в её использовании, а в решении не удалять старые файлы дольше, чем нужно для технических нужд; этому вопросу посвящён отдельный раздел этого урока про горизонт хранения, и его прямое продолжение - десятый урок про expire_snapshots.
Способы адресации исторических данных: Timestamp, Snapshot ID и именованные ссылки¶
Чтобы прочитать таблицу «в прошлом», нужно как-то указать движку, какую именно историческую точку имеется в виду. Iceberg поддерживает три принципиально разных способа адресации, и важно понимать не просто синтаксис каждого, а механику поиска нужного снапшота и trade-off'ы, которые определяют, какой способ использовать в конкретной ситуации.
Адресация по Snapshot ID¶
Каждый снапшот в Iceberg имеет уникальный 64-битный идентификатор snapshot_id, присваиваемый в момент коммита. Это число можно увидеть в системной таблице <table>.snapshots, в выводе процедур обслуживания, либо в логах пайплайна, который сделал коммит. Адресация по snapshot_id - самый строгий и однозначный способ: запрос либо находит снапшот с точно таким идентификатором и читает его, либо завершается ошибкой, если такого snapshot_id не существует в истории таблицы (например, потому что он уже был удалён expire_snapshots, либо потому что идентификатор просто неверен). Никакой эвристики, никакого «ближайшего совпадения» - либо точное совпадение, либо явная ошибка.
Адресация по Timestamp¶
Адресация по времени интуитивно понятнее людям: «покажи таблицу такой, какой она была вчера в 15:00» - естественная формулировка для расследования инцидента, когда известно примерное время происшествия, но не известен точный snapshot_id. Механика поиска здесь принципиально важна и часто понимается неверно, поэтому стоит сформулировать её предельно точно.
Iceberg хранит snapshot-log - хронологический список записей (timestamp_ms, snapshot_id), фиксирующих, какой снапшот был текущим начиная с какого момента (та же структура, что лежит в основе системной таблицы history, уже встречавшейся во втором уроке). Когда запрос указывает TIMESTAMP AS OF '2025-06-15 15:00:00', движок находит последний снапшот, который стал текущим не позже указанного момента времени - то есть идёт по snapshot-log назад во времени и берёт первую запись, чей timestamp_ms меньше или равен запрошенному. Это не поиск «ближайшего» снапшота в обе стороны - снапшот, созданный через минуту после запрошенного времени, никогда не будет выбран, даже если он отделён от запрошенного момента всего на несколько секунд, а предыдущий валидный снапшот отделён часами. Это поведение логически обязательно: иначе путешествие в «прошлое» иногда могло бы тайно возвращать данные из будущего относительно запрошенной точки, что разрушило бы саму идею детерминированного исторического снимка.
Диаграмма показывает ровно этот нюанс: несмотря на то, что Snapshot D (16:10) ближе к запрошенным 15:00 по абсолютной разнице во времени (1 час 10 минут), чем Snapshot C (14:50, разница чуть больше часа в противоположную сторону), движок всегда выбирает только тех кандидатов, что не позже запрошенного момента, поэтому однозначно выбирается Snapshot C. Если запрошенный timestamp окажется раньше, чем committed_at самого первого снапшота таблицы (то есть если вообще не существует снапшота, удовлетворяющего условию «не позже»), Spark вернёт явную ошибку вида Cannot find a snapshot older than ... - в этой ситуации Time Travel честно сообщает о невозможности выполнить запрос, а не возвращает пустой результат или произвольный снапшот.
Второй важный нюанс касается часовых зон. Время committed_at/made_current_at, хранящееся в метаданных, - это абсолютный момент в UTC (миллисекунды с начала эпохи). Строковый литерал TIMESTAMP AS OF '2025-06-15 15:00:00' интерпретируется Spark в соответствии с текущей сессионной настройкой spark.sql.session.timeZone - если эта настройка не совпадает с тем часовым поясом, в котором инженер мысленно представляет себе «15:00», запрос будет искать снапшот на момент, отличающийся от ожидаемого ровно на разницу часовых зон. Это не теоретический риск - в разделе с производственным кейсом этого урока разобран именно такой инцидент, стоивший команде потерянного дня воспроизведения ML-эксперимента.
Адресация по именованным ссылкам: tag и branch¶
Третий способ адресации - использование именованных ссылок (refs), представленных во втором уроке модуля при разборе поля refs в metadata.json: tag (замороженная неизменяемая метка на конкретный снапшот, например audit-2024-q2) и branch (изменяемый указатель, который продолжает двигаться вперёд с каждым новым коммитом в неё, аналог main в Git). Time Travel умеет адресоваться непосредственно по имени такой ссылки, что сочетает удобство (понятное человеку имя вместо числа) со строгостью (имя ссылается ровно на один конкретный снапшот в каждый момент - для тега навсегда, для ветки до следующего коммита в неё).
-- предполагается, что тег уже создан заранее (второй урок модуля):
-- CALL lakehouse.system.create_tag('lakehouse.analytics.orders', 'audit-2024-q2', <snapshot_id>)
SELECT *
FROM lakehouse.analytics.orders VERSION AS OF 'audit-2024-q2';
Здесь важен тонкий нюанс, отличающий адресацию по tag от адресации по branch. Тег - это замороженная точка: VERSION AS OF 'audit-2024-q2' всегда вернёт ровно тот же снапшот, сколько бы коммитов ни произошло в основной истории таблицы после создания тега, - в этом смысле обращение по тегу эквивалентно обращению по конкретному snapshot_id, просто под более удобным человеку именем. Ветка же продолжает двигаться: VERSION AS OF 'experimental-etl' сегодня и VERSION AS OF 'experimental-etl' через час после нескольких коммитов в эту ветку могут вернуть два разных снапшота - то есть обращение по имени branch - это не «путешествие в фиксированную точку прошлого», а «чтение текущей головы параллельной линии разработки». Эта разница не делает branch-адресацию разновидностью Time Travel в строгом смысле слова, но она использует тот же синтаксис VERSION AS OF, поэтому важно держать в голове, какую именно ссылку - замороженную или движущуюся - использует конкретный запрос.
Сравнение трёх способов: что выбрать и когда¶
| Способ адресации | Удобство для человека | Надёжность в автоматизации (CI/CD, пайплайны) | Типичный сценарий использования |
|---|---|---|---|
snapshot_id |
Низкое - нужно знать или предварительно найти число | Максимальная - точное совпадение или явная ошибка, никакой неоднозначности | Пайплайны, фиксирующие воспроизводимость; rollback после инцидента, когда snapshot_id уже найден через history |
timestamp |
Высокое - естественный язык расследования («вчера в 15:00») | Низкая - зависит от часового пояса сессии, от точной семантики «не позже», и от того, что относительный смысл даты со временем меняется | Ad hoc расследование инцидента аналитиком или дежурным инженером, когда snapshot_id ещё не известен |
Именованная ссылка (tag/branch) |
Высокое - осмысленное имя вместо числа или даты | Высокая, но требует заранее созданной ссылки - сама по себе не появляется без явного CREATE TAG/CREATE BRANCH |
Зафиксированные contractually значимые точки («конец квартала», «before-migration snapshot»), параллельные ветки разработки ETL-логики |
Практический паттерн, вытекающий из этой таблицы и развиваемый дальше в разделе про динамическую параметризацию: человек начинает расследование с timestamp («что-то случилось вчера около обеда»), находит через системную таблицу history точный snapshot_id, соответствующий моменту до инцидента, и именно этот snapshot_id - а не исходный приблизительный timestamp - фиксируется во всех последующих шагах: в команде rollback, в параметрах автоматизированного пайплайна, в тикете расследования. Timestamp - инструмент первого контакта с проблемой; snapshot_id - инструмент её точной и воспроизводимой фиксации.
Синтаксис Time Travel в Spark SQL¶
Spark поддерживает две взаимозаменяемые формы синтаксиса для Time Travel внутри блока FROM: «классическую» форму Iceberg VERSION AS OF / TIMESTAMP AS OF и ANSI SQL-форму FOR SYSTEM_VERSION AS OF / FOR SYSTEM_TIME AS OF, появившуюся в Spark начиная с версии 3.3 как часть общей для всех источников данных v2 спецификации временных запросов (AS OF syntax, описанной в SQL:2011). Обе формы транслируется в один и тот же физический план и поддерживаются Iceberg идентично - выбор между ними обычно определяется стилем команды или совместимостью с другими движками, которые могут понимать только одну из форм.
VERSION AS OF: адресация по снапшоту или именованной ссылке¶
-- Адресация по конкретному snapshot_id (число, без кавычек)
SELECT *
FROM lakehouse.analytics.orders VERSION AS OF 5847203991234567890;
-- Адресация по имени тега или ветки (строка, в кавычках) -
-- Spark определяет, что это ссылка, а не число, по типу литерала
SELECT *
FROM lakehouse.analytics.orders VERSION AS OF 'audit-2024-q2';
-- Полностью эквивалентная ANSI-форма того же запроса по snapshot_id
SELECT *
FROM lakehouse.analytics.orders
FOR SYSTEM_VERSION AS OF 5847203991234567890;
Обратите внимание на то, как Spark различает «снапшот по числу» и «ссылку по имени» в одной и той же синтаксической позиции: если литерал после VERSION AS OF - целое число, он интерпретируется как snapshot_id; если это строка в кавычках, Iceberg ищет среди refs таблицы тег или ветку с таким именем. Это удобно, но создаёт небольшую ловушку для невнимательного копирования: если переменная со snapshot_id случайно подставится в запрос как строка (например, через f-string без преобразования типа), число в кавычках '5847203991234567890' будет интерпретировано не как snapshot_id, а как попытка найти ссылку с таким текстовым именем - и запрос завершится ошибкой Cannot find matching snapshot ID or reference name, а не тихо вернёт неверные данные. Ошибка в данном случае - это хорошо: лучше явный отказ, чем скрытая путаница между числом и именем.
TIMESTAMP AS OF: адресация по времени¶
-- Классическая форма Iceberg
SELECT *
FROM lakehouse.analytics.orders
TIMESTAMP AS OF '2025-06-15 15:00:00';
-- Эквивалентная ANSI-форма
SELECT *
FROM lakehouse.analytics.orders
FOR SYSTEM_TIME AS OF '2025-06-15 15:00:00';
-- Литерал TIMESTAMP вместо строки - так же допустимо
SELECT *
FROM lakehouse.analytics.orders
TIMESTAMP AS OF TIMESTAMP '2025-06-15 15:00:00';
Все три формы строкового литерала, литерала TIMESTAMP и ANSI-синтаксиса абсолютно идентичны по результату - они различаются только нотацией. Механика поиска снапшота, разобранная в предыдущем разделе («последний снапшот не позже указанного момента»), применяется одинаково независимо от выбранной нотации.
Time Travel внутри произвольного запроса, а не только в простом SELECT¶
Конструкция VERSION AS OF/TIMESTAMP AS OF синтаксически привязана не к запросу целиком, а к конкретной ссылке на таблицу в блоке FROM - то есть указанную таблицу можно time-travel'ить внутри любого более сложного запроса: JOIN, подзапроса, CREATE TABLE AS SELECT. Более того, если один и тот же физический table reference встречается в запросе несколько раз (например, через self-join), каждое упоминание может быть привязано к собственной исторической точке - это открывает практический паттерн «сравнить две версии одной таблицы одним запросом», без необходимости делать два отдельных запроса и сравнивать результаты на стороне клиента.
# Сравнение состояния витрины заказов "сейчас" и "неделю назад" -
# одним SQL-запросом, без двух отдельных round-trip'ов к Spark
diff_df = spark.sql("""
SELECT
cur.order_id,
cur.status AS status_now,
old.status AS status_week_ago,
cur.amount AS amount_now,
old.amount AS amount_week_ago
FROM lakehouse.analytics.orders cur
JOIN lakehouse.analytics.orders TIMESTAMP AS OF '2025-06-08 00:00:00' AS old
ON cur.order_id = old.order_id
WHERE cur.status != old.status
OR cur.amount != old.amount
""")
diff_df.show(20, truncate=False)
Запрос выше делает ровно то, что было анонсировано в начале урока как один из бизнес-кейсов Time Travel: сравнение «тогда и сейчас» средствами одного JOIN, где левая часть читает таблицу в её текущем состоянии (без какого-либо AS OF, то есть последний снапшот), а правая - в состоянии на конкретный исторический момент. Условие WHERE отфильтровывает только реально изменившиеся заказы, превращая запрос в готовый отчёт по дрифту данных за неделю. Это невозможно реализовать в классическом Hive-подходе без отдельного хранения недельного снапшота заранее - в Iceberg это работает «бесплатно» для любой пары моментов времени, для которых ещё не истёк expire_snapshots.
Гипотетическая адресация по будущему: чего Time Travel не делает¶
Стоит явно закрыть вопрос, который иногда возникает у студентов по аналогии с названием фичи: Time Travel не позволяет «увидеть будущее» в смысле прогноза или незакоммиченных изменений - он работает только с уже зафиксированными снапшотами, то есть строго с прошлым относительно момента выполнения запроса. Если указать snapshot_id, который ещё не существует на момент запроса (потому что соответствующий коммит ещё не произошёл), запрос немедленно завершится ошибкой Cannot find matching snapshot ID, а не будет ждать его появления. Time Travel - это история, а не подписка на будущие изменения; для отслеживания новых изменений в реальном времени используется Spark Structured Streaming над Iceberg-таблицей (тема, пересекающаяся с CDC-материалом восьмого урока, но не Time Travel как таковой).
Time Travel в PySpark DataFrame API¶
Помимо SQL, Iceberg предоставляет Time Travel через стандартный DataFrameReader, передавая параметры как опции источника данных iceberg. Это нужно, когда исторический срез - не конечная цель запроса, а промежуточный DataFrame, который дальше обрабатывается средствами PySpark API (трансформации, запись в другую таблицу, передача в ML-пайплайн).
# Адресация по snapshot_id - аналог VERSION AS OF с числом
df_by_snapshot = (
spark.read
.format("iceberg")
.option("snapshot-id", 5847203991234567890)
.load("lakehouse.analytics.orders")
)
# Адресация по имени тега или ветки - аналог VERSION AS OF со строкой
df_by_tag = (
spark.read
.format("iceberg")
.option("tag", "audit-2024-q2")
.load("lakehouse.analytics.orders")
)
df_by_branch = (
spark.read
.format("iceberg")
.option("branch", "experiment-new-dedup-logic")
.load("lakehouse.analytics.orders")
)
Каждая опция отвечает за свой способ адресации, разобранный в предыдущих разделах, и они взаимоисключающие - указание сразу двух противоречащих друг другу опций (например, snapshot-id и tag одновременно) приводит к ошибке валидации на стороне Iceberg до начала чтения данных.
Критическая деталь: as-of-timestamp ожидает миллисекунды, а не строку¶
Адресация по времени в DataFrameReader API называется as-of-timestamp, и здесь возникает несовместимость форматов, о которую легко споткнуться при переносе кода из SQL в PySpark API. SQL-форма TIMESTAMP AS OF '2025-06-15 15:00:00' принимает человекочитаемую строку. PySpark-опция as-of-timestamp, напротив, ожидает целое число миллисекунд, прошедших с начала эпохи Unix (1970-01-01 00:00:00 UTC), переданное как строка или число - человекочитаемая дата здесь не распознаётся вообще и приведёт либо к ошибке парсинга, либо (что хуже) к тихой попытке интерпретировать строку как число с непредсказуемым результатом.
from datetime import datetime, timezone
# НЕПРАВИЛЬНО: человекочитаемая строка не распознаётся as-of-timestamp
# df_wrong = (
# spark.read.format("iceberg")
# .option("as-of-timestamp", "2025-06-15 15:00:00") # ошибка парсинга
# .load("lakehouse.analytics.orders")
# )
# ПРАВИЛЬНО: явное преобразование человекочитаемой даты в epoch-миллисекунды,
# с явным указанием часового пояса - та же ловушка с timezone, что и в SQL,
# здесь даже более коварна, потому что неверный результат не всегда очевиден
moment = datetime(2025, 6, 15, 15, 0, 0, tzinfo=timezone.utc)
as_of_ms = int(moment.timestamp() * 1000)
df_by_timestamp = (
spark.read
.format("iceberg")
.option("as-of-timestamp", as_of_ms)
.load("lakehouse.analytics.orders")
)
print(f"as_of_ms = {as_of_ms}")
# as_of_ms = 1750000800000
Эта несовместимость форматов между SQL и DataFrameReader API - не баг, а следствие того, что SQL-парсер Spark самостоятельно разбирает строковый литерал TIMESTAMP средствами своего стандартного типа данных (с учётом сессионного часового пояса), тогда как опции DataFrameReader передаются Iceberg как набор строковых пар «ключ-значение» без какой-либо синтаксической интерпретации на стороне Spark SQL - Iceberg получает голую строку и обязан задать для неё однозначный, не зависящий от локали и часового пояса формат, и таким форматом исторически стали epoch-миллисекунды UTC. Практический вывод: при переносе пайплайна с SQL-формы на PySpark API (или наоборот) недостаточно скопировать строку с датой - необходимо явно произвести конвертацию формата, и желательно - вынести её в отдельную, покрытую тестами функцию, а не повторять datetime.timestamp() * 1000 в каждом месте, где это нужно.
Инкрементальное чтение: соседняя, но другая возможность¶
Помимо точечного Time Travel (прочитать состояние на один момент), Iceberg поддерживает инкрементальное чтение - запрос всех строк, добавленных в диапазоне между двумя снапшотами, без повторного чтения данных, не изменившихся за этот период. Это упоминалось во втором уроке модуля при разборе parent-snapshot-id lineage и использует тот же граф родительских связей, но решает принципиально иную задачу: не «как выглядела таблица в момент X», а «что изменилось между X и Y».
incremental_df = (
spark.read
.format("iceberg")
.option("start-snapshot-id", 5847203991111111111) # не включается в результат
.option("end-snapshot-id", 5847203991234567890) # включается в результат
.load("lakehouse.analytics.orders")
)
Важно не путать инкрементальное чтение с Time Travel: start-snapshot-id/end-snapshot-id работают только с операциями типа append (то есть с накопленными вставками между двумя точками) и не предназначены для воспроизведения полного состояния таблицы на момент времени - для этого нужен именно snapshot-id/as-of-timestamp, разобранные выше. Главная сфера применения инкрементального чтения - построение собственных систем потоковой репликации или аудита изменений на основе чистого batch API, без подписки на Structured Streaming.
Динамическая параметризация: автоматизация исторических выгрузок¶
Раздел про сравнение способов адресации заканчивался практическим выводом: timestamp хорош для первого контакта с проблемой, а snapshot_id - для точной и воспроизводимой фиксации найденного решения. Этот раздел показывает, как это выглядит в виде кода, а не только как принцип.
Функция перевода «человеческого момента» в snapshot_id¶
Раз PySpark API не умеет сам резолвить timestamp в snapshot_id (в отличие от SQL, где это происходит неявно внутри TIMESTAMP AS OF), для воспроизводимых пайплайнов полезно сделать этот перевод явным шагом, который выполняется один раз и логируется, а не повторяется неявно при каждом запуске.
def resolve_snapshot_as_of(spark, table_fqn: str, as_of_ts: "datetime") -> int:
"""Находит snapshot_id, который был текущим не позже указанного момента.
Тот же алгоритм поиска, что использует TIMESTAMP AS OF внутри Spark SQL,
но выполненный явно через системную таблицу history - чтобы зафиксировать
результат как обычное число и положить его в лог запуска пайплайна,
а не оставлять резолвинг "спрятанным" внутри SQL-литерала.
"""
as_of_ms = int(as_of_ts.timestamp() * 1000)
row = spark.sql(f"""
SELECT snapshot_id
FROM {table_fqn}.history
WHERE made_current_at <= timestamp_millis({as_of_ms})
AND is_current_ancestor = true
ORDER BY made_current_at DESC
LIMIT 1
""").collect()
if not row:
raise ValueError(
f"Нет снапшота {table_fqn}, ставшего текущим не позже {as_of_ts}"
)
return row[0]["snapshot_id"]
Условие is_current_ancestor = true в запросе выше - не случайная деталь, а прямое следствие того, что history отражает весь snapshot-log, включая записи, которые перестали быть частью текущей линии после rollback (механика разобрана в следующем разделе урока). Без этого фильтра функция рискует вернуть snapshot_id, который технически был «текущим» в какой-то момент прошлого, но затем оказался «обойдён» откатом и больше не является предком сегодняшнего текущего состояния - что обычно не то поведение, которое нужно от функции, ищущей точку для воспроизводимого чтения.
Зачем фиксировать именно snapshot_id, а не timestamp, в манифесте запуска пайплайна¶
Практика, которую стоит закрепить как стандарт для любого ML-пайплайна или регламентного отчёта, претендующего на воспроизводимость: момент, когда был построен тренировочный датасет или отчёт, должен фиксироваться в виде snapshot_id исходной таблицы (или нескольких таблиц, если их несколько), а не в виде timestamp запуска. Это можно сделать в виде простого JSON-манифеста запуска, сохраняемого рядом с артефактом модели или отчёта.
import json
from datetime import datetime, timezone
def build_run_manifest(spark, table_fqn: str) -> dict:
"""Фиксирует snapshot_id таблицы-источника в момент запуска пайплайна -
манифест, по которому через месяцы можно точно воспроизвести вход."""
current_snapshot = spark.sql(f"""
SELECT snapshot_id, committed_at
FROM {table_fqn}.snapshots
ORDER BY committed_at DESC
LIMIT 1
""").collect()[0]
return {
"source_table": table_fqn,
"source_snapshot_id": current_snapshot["snapshot_id"],
"source_committed_at": str(current_snapshot["committed_at"]),
"run_started_at": datetime.now(timezone.utc).isoformat(),
}
manifest = build_run_manifest(spark, "lakehouse.analytics.orders")
print(json.dumps(manifest, indent=2, default=str))
# Сохраняем манифест рядом с артефактом модели/отчёта - через произвольное
# время воспроизведение входа сводится к одному чтению по snapshot_id:
replay_df = (
spark.read.format("iceberg")
.option("snapshot-id", manifest["source_snapshot_id"])
.load(manifest["source_table"])
)
Разница между двумя подходами становится ощутимой ровно тогда, когда воспроизведение действительно требуется - то есть не сразу, а спустя недели или месяцы, когда детали запуска уже забыты, исходная таблица успела накопить множество новых коммитов, а возможно - и поменять часовой пояс сессии у инженера, который выполняет повторный прогон. Манифест со snapshot_id отвечает на вопрос «откуда брать данные» однозначно и без интерпретации; манифест с timestamp требует повторного резолвинга через history, который зависит от того, что таблица «history» ещё хранит нужный диапазон (то есть что expire_snapshots ещё не вычистил нужный снапшот) - тема, к которой урок вернётся в разделе про горизонт хранения.
На практике пайплайн редко зависит от одной таблицы - типичный ML-пайплайн собирает признаки из нескольких источников (orders, customers, currency_rates из практического блока этого урока), и манифест должен фиксировать snapshot_id для каждого из них независимо, а не один общий timestamp для всех сразу. Причина та же, что и в случае с production-кейсом этого урока: разные таблицы коммитятся в разное время и с разной частотой, и единственный способ гарантировать, что повторный прогон через полгода прочитает ровно тот же набор строк из каждой из них - явно перечислить snapshot_id для каждой таблицы-источника по отдельности, а не понадеяться, что один timestamp одинаково корректно разрешится во всех таблицах сразу.
def build_multi_table_manifest(spark, table_fqns: list[str]) -> dict:
"""Расширение build_run_manifest на несколько источников: каждая
таблица фиксируется своим собственным snapshot_id независимо от других,
поскольку частота коммитов у них, как правило, разная."""
manifest = {
"run_started_at": datetime.now(timezone.utc).isoformat(),
"sources": {},
}
for table_fqn in table_fqns:
snap = spark.sql(f"""
SELECT snapshot_id, committed_at
FROM {table_fqn}.snapshots
ORDER BY committed_at DESC
LIMIT 1
""").collect()[0]
manifest["sources"][table_fqn] = {
"snapshot_id": snap["snapshot_id"],
"committed_at": str(snap["committed_at"]),
}
return manifest
manifest = build_multi_table_manifest(spark, [
"lakehouse.analytics.orders",
"lakehouse.analytics.customers",
"lakehouse.reference.currency_rates",
])
Архитектурный откат таблицы: Time Travel (чтение) vs Rollback (запись)¶
Всё, что было разобрано до этого момента - VERSION AS OF, TIMESTAMP AS OF, опции DataFrameReader - представляет собой read-only операции. Они не изменяют состояние таблицы ни на бит: запрос с Time Travel возвращает DataFrame, как и любой обычный SELECT, не оставляя после себя никакого нового снапшота и не трогая поле current-snapshot-id в metadata.json. Сколько угодно пользователей могут одновременно читать таблицу на разные исторические моменты, и ни один из этих запросов не повлияет на то, что увидят все остальные читатели, обращающиеся к текущей версии таблицы без AS OF.
Rollback - принципиально другая операция. Она не читает данные, а изменяет то, какой снапшот является текущим для всех будущих читателей таблицы без явного Time Travel - то есть rollback - это write-операция, требующая коммита через тот же механизм atomic compare-and-swap указателя в catalog, что был разобран во втором уроке модуля при описании Write-and-Commit Path. Это различие стоит закрепить максимально явно, потому что путаница между «прочитать прошлое» и «вернуть таблицу в прошлое для всех» - источник реальных production-инцидентов: инженер, желающий лишь посмотреть на старые данные для расследования, по ошибке выполняет rollback, и обнаруживает, что только что без предупреждения откатил taблицу для всех потребителей, включая dashboard'ы и downstream-пайплайны, которые в этот момент читают «текущую» версию.
Механика Rollback под капотом: что меняется в metadata.json, а что нет¶
Здесь важно исправить интуицию, которая иногда складывается из упрощённых описаний rollback как «создания нового снапшота, дублирующего старый». Это не совсем точно, и неточность важна для правильного понимания цены операции. Технически rollback_to_snapshot выполняет следующее: catalog атомарно заменяет указатель на metadata.json новой версией файла метаданных (тот же commit-механизм Stage 3 из второго урока), в которой поле current-snapshot-id установлено равным snapshot_id целевого (более раннего) снапшота, и в массив snapshot-log добавлена новая запись (новый timestamp, тот же snapshot_id цели). При этом массив snapshots (список самих объектов-снапшотов с их манифест-листами) не получает новой записи - целевой снапшот уже существовал там до отката и просто становится «текущим» снова. Ни один манифест, ни один data file и ни один delete file при этой операции не читается, не пишется и не перезаписывается - rollback касается исключительно одного маленького JSON-файла метаданных.
Диаграмма показывает три состояния: «до», условие, которое проверяется внутри процедуры, и «после». Ключевая деталь, которую стоит прочитать внимательно: массив snapshots идентичен до и после - в нём как до, так и после ровно 4 объекта (S1-S4), ни один не добавлен и не удалён. Единственные изменения - это новое значение поля current-snapshot-id (которое теперь снова 2) и одна новая строка в snapshot-log, фиксирующая факт и время этого переключения. Снапшоты 3 и 4 - включая тот самый испорченный MERGE INTO, который и стал причиной отката - остаются полностью на месте в истории таблицы, доступные для дальнейшего расследования через Time Travel, до тех пор, пока их явно не удалит expire_snapshots. Это прямое продолжение материала про Snapshot Lineage из второго урока: откат в Iceberg концептуально гораздо ближе к git checkout <старый коммит> (который не удаляет более новые коммиты из истории репозитория), чем к git reset --hard (который физически переписывает указатель ветки, отрезая последующие коммиты от графа достижимости, и при отсутствии других ссылок на них рискует сделать их недостижимыми для сборщика мусора).
Именно эта механика - «передвинуть один указатель, не трогая остальные структуры» - и есть причина, по которой rollback занимает constant time независимо от объёма данных в таблице: стоимость операции определяется размером одного JSON-файла метаданных (обычно десятки-сотни килобайт даже для таблиц с тысячами снапшотов), а не размером данных, на которые этот файл ссылается.
Команды управления: rollback_to_snapshot, rollback_to_timestamp и set_current_snapshot¶
Spark SQL предоставляет доступ к этим операциям через хранимые процедуры (CALL), зарегистрированные Iceberg-расширениями Spark - тот же механизм, что использовался во втором уроке для CREATE TAG/CREATE BRANCH.
-- Откат к конкретному snapshot_id. Требует, чтобы целевой снапшот
-- был ПРЕДКОМ текущего - то есть находился в цепочке parent-snapshot-id
-- между корнем истории и текущим снапшотом.
CALL lakehouse.system.rollback_to_snapshot(
'analytics.orders', 5847203991111111111
);
-- Откат к снапшоту, который был текущим не позже указанного момента -
-- внутри процедуры используется тот же алгоритм поиска "последний <= T",
-- что и в TIMESTAMP AS OF, после чего применяется то же ограничение
-- "целевой снапшот должен быть предком текущего".
CALL lakehouse.system.rollback_to_timestamp(
'analytics.orders', TIMESTAMP '2025-06-15 14:50:00'
);
Обе процедуры возвращают результирующий набор с двумя колонками, отражающими произошедшее переключение указателя:
+----------------------+----------------------+
|previous_snapshot_id |current_snapshot_id |
+----------------------+----------------------+
|5847203991234567890 |5847203991111111111 |
+----------------------+----------------------+
Если запрошенный целевой снапшот не является предком текущего - например, инженер по ошибке указывает snapshot_id из параллельной ветки, либо снапшот, который сам уже был «обойдён» более ранним откатом и больше не входит в линию предков - обе процедуры завершаются ошибкой валидации (ValidationException: Cannot roll back, snapshot is not an ancestor of the current state), а не выполняют операцию «как получится». Это намеренное ограничение, защищающее от случайного «отката в случайную сторону»: rollback предназначен исключительно для возврата строго назад по уже пройденному пути, а не для перемещения в произвольную точку графа снапшотов.
Для случаев, когда нужно именно произвольное перемещение - например, «redo» после слишком резкого отката, или перенос снапшота из экспериментальной ветки в main (паттерн, близкий к git cherry-pick) - Iceberg предоставляет третью процедуру без ограничения на предков:
-- set_current_snapshot МОЖЕТ переключить указатель на снапшот,
-- который не является предком текущего - в отличие от rollback_to_snapshot
CALL lakehouse.system.set_current_snapshot(
'analytics.orders', snapshot_id => 5847203991234567890
);
-- Альтернативная форма: переключение по имени ссылки, а не по числу
CALL lakehouse.system.set_current_snapshot(
'analytics.orders', ref => 'audit-2024-q2'
);
| Процедура | Ограничение на целевой снапшот | Типичный сценарий |
|---|---|---|
rollback_to_snapshot |
Только предок текущего снапшота | Немедленный откат после обнаруженного инцидента, путь уже пройден |
rollback_to_timestamp |
Только предок (после резолвинга timestamp → snapshot_id внутри) | То же, но когда snapshot_id ещё не найден, известен лишь примерный момент |
set_current_snapshot |
Любой существующий снапшот таблицы, включая неандцестора | "Redo" после чрезмерного отката, перенос снапшота из другой ветки в текущую линию |
Эта таблица закрывает раздел практическим правилом: если цель - «вернуться строго назад по тому пути, который таблица уже прошла», используется rollback_to_snapshot/rollback_to_timestamp с их встроенной защитой от случайного перемещения не туда; если цель - что-то более нестандартное (перенос снапшота, восстановление после слишком глубокого отката), нужен более низкоуровневый и менее защищённый set_current_snapshot, использование которого имеет смысл сопровождать дополнительным ручным подтверждением, что выбранный snapshot_id - действительно тот, что нужен.
Ограничения Time Travel и цена хранения истории¶
Всё описанное выше создаёт впечатление, что Time Travel - это безусловно бесплатная фича без какого-либо trade-off'а. Это не совсем так: бесплатно само использование Time Travel (выполнение запроса с AS OF не стоит дороже обычного SELECT к манифестам нужного снапшота), но не бесплатно хранение того диапазона истории, в пределах которого Time Travel остаётся возможным. Этот раздел разбирает именно эту вторую сторону.
Горизонт хранения: почему нельзя путешествовать в прошлое бесконечно¶
История снапшотов не растёт бесконечно сама по себе только потому, что Iceberg «забывает» удалять старые версии - она растёт, потому что никто не запускает процедуру, которая отвечает за чистку. Этой процедурой является expire_snapshots (детально - тема десятого урока), и до её запуска снапшоты продолжают накапливаться с каждым новым коммитом, без каких-либо автоматических ограничений по умолчанию, кроме двух table-level свойств, определяющих, что именно считается «безопасным к удалению» в момент следующего запуска expire_snapshots:
| Свойство | Назначение | Типичное значение по умолчанию |
|---|---|---|
history.expire.max-snapshot-age-ms |
Снапшоты старше этого возраста считаются кандидатами на удаление при следующем expire_snapshots |
5 дней (432 000 000 мс) |
history.expire.min-snapshots-to-keep |
Минимальное количество последних снапшотов, которые сохраняются независимо от возраста - защита от случая, когда все коммиты происходят реже, чем раз в max-snapshot-age-ms |
1 |
Сами по себе эти свойства не выполняют удаление - они лишь задают политику, которую применит expire_snapshots при явном вызове. Иными словами, простое наличие этих настроек в TBLPROPERTIES ничего не освобождает: пока процедура не запущена (вручную или по расписанию), снапшоты старше max-snapshot-age-ms продолжают физически существовать и продолжают быть доступны через Time Travel, даже если согласно политике они уже считаются «просроченными». Это важная деталь: горизонт хранения - это не жёсткий технический предел, как может показаться из названия, а декларация политики, которая становится действующей только в момент исполнения maintenance-процедуры.
Влияние на дисковое пространство: где здесь реальные деньги¶
Стоимость хранения истории - это не абстракция, а прямое следствие Write Amplification, разобранного в шестом уроке модуля. Если таблица использует Copy-on-Write и обслуживает частые точечные UPDATE/MERGE, то каждый такой коммит полностью переписывает все затронутые файлы - и старая версия каждого переписанного файла продолжает занимать место в object storage до тех пор, пока ссылающийся на неё снапшот не будет удалён expire_snapshots. Чем дольше горизонт хранения (выше max-snapshot-age-ms) и чем чаще происходят CoW-операции, тем больше копий одних и тех же логических данных одновременно лежит на диске.
Диаграмма переносит конкретные числа из материала шестого и седьмого уроков в контекст принятия решения о retention: при одинаковой частоте коммитов и одинаковом горизонте хранения CoW накапливает кратно больший объём «мёртвого», но пока не удалённого веса в object storage, потому что каждый коммит создаёт полные копии затронутых файлов, а не компактные дельты. Это не аргумент в пользу того, что MoR всегда лучше для retention - у MoR своя цена накопления, разобранная в седьмом уроке (деградация чтения от роста числа delete-файлов между компакциями) - но это прямое объяснение, почему вопрос «сколько дней хранить историю» нельзя решать одним универсальным числом без учёта профиля нагрузки конкретной таблицы и выбранного режима записи.
Практический ориентир, вытекающий из этого расчёта: для таблиц с высокочастотными точечными изменениями на CoW длинный горизонт хранения (недели) может оказаться неожиданно дорогим именно из-за множителя Write Amplification - в то время как для таблиц с редкими bulk-загрузками (раз в день, append-only) тот же горизонт хранения практически не создаёт лишней нагрузки на storage, потому что между коммитами не происходит переписывания существующих файлов.
Опережающий анонс: как expire_snapshots завершает этот разговор¶
Всё, что было сказано про retention в этом разделе - это описание политики, а не исполнения. Исполнение происходит в десятом уроке модуля, который разбирает процедуру expire_snapshots - она физически удаляет снапшоты старше max-snapshot-age-ms (с учётом min-snapshots-to-keep) из массива snapshots, а также запускает сборку мусора для файлов данных и манифестов, на которые больше не ссылается ни один оставшийся снапшот. Важное следствие, которое стоит зафиксировать здесь заранее: после успешного выполнения expire_snapshots Time Travel к удалённому снапшоту становится невозможным - запрос VERSION AS OF с таким snapshot_id вернёт ошибку Cannot find matching snapshot ID, ровно как если бы такого снапшота никогда не существовало. Это необратимая операция в том смысле, что после физического удаления файлов данных откатить её средствами самого Iceberg уже нельзя - единственный путь назад в такой ситуации - это внешний бэкап object storage, то есть тот самый медленный путь, который Time Travel вообще существует, чтобы сделать не нужным в первую очередь. Это прямая причина, по которой выбор горизонта хранения - не техническая деталь конфигурации, а осознанное архитектурное решение, к которому урок вернётся в чек-листе и в мостике к десятому уроку.
Практический демо-блок: ломаем витрину и чиним за 5 секунд¶
Полная конфигурация SparkSession (JDBC Catalog на self-hosted PostgreSQL, S3FileIO на MinIO) приведена в первом уроке модуля и здесь не повторяется - предполагается, что spark уже инициализирован, а каталог lakehouse подключён. Сценарий этого демо-блока намеренно собран из материала всех предыдущих уроков модуля: таблица заказов lakehouse.analytics.orders (та же, что использовалась в шестом-восьмом уроках) служит источником, а новая Gold-витрина дневной выручки - целью катастрофы и последующего восстановления.
Кейс 0: подготовка - витрина дневной выручки поверх таблицы заказов¶
spark.sql("""
CREATE TABLE IF NOT EXISTS lakehouse.gold.daily_revenue_report (
report_date DATE,
total_orders BIGINT,
total_revenue DECIMAL(18, 2)
)
USING iceberg
""")
# Базовая загрузка витрины из таблицы заказов - один агрегат на дату,
# без учёта отменённых заказов
spark.sql("""
INSERT OVERWRITE lakehouse.gold.daily_revenue_report
SELECT
CAST(updated_at AS DATE) AS report_date,
COUNT(*) AS total_orders,
SUM(amount) AS total_revenue
FROM lakehouse.analytics.orders
WHERE status != 'cancelled'
GROUP BY CAST(updated_at AS DATE)
""")
spark.sql("""
SELECT * FROM lakehouse.gold.daily_revenue_report ORDER BY report_date
""").show(truncate=False)
# +-----------+------------+-------------+
# |report_date|total_orders|total_revenue|
# +-----------+------------+-------------+
# |2025-06-14 |842 |184320.50 |
# |2025-06-15 |915 |201450.75 |
# +-----------+------------+-------------+
Это чистая, верная версия витрины - именно к ней нужно вернуться, когда чуть позже она будет испорчена. Фиксируем snapshot_id этого «чистого» состояния через системную таблицу snapshots:
spark.sql("""
SELECT snapshot_id, committed_at, operation
FROM lakehouse.gold.daily_revenue_report.snapshots
ORDER BY committed_at
""").show(truncate=False)
# +--------------------+-----------------------+---------+
# |snapshot_id |committed_at |operation|
# +--------------------+-----------------------+---------+
# |9001 |2025-06-15 09:00:00.000|overwrite|
# +--------------------+-----------------------+---------+
snapshot_id = 9001 - единственный на данный момент снапшот, соответствующий корректной витрине. Запоминаем его как точку отсчёта для всего дальнейшего сценария (в реальном расследовании этот номер находится точно так же - запросом к snapshots, а не «известен заранее»).
Кейс 1: симуляция катастрофы - инженер ломает витрину «безопасным» изменением¶
Финансовая команда попросила учитывать курс валюты при пересчёте выручки. Инженер добавляет JOIN со справочником курсов lakehouse.ref.currency_rates, не заметив, что справочник хранит две строки для рубля - старую и недавно скорректированную, обе с currency_code = 'RUB', различающиеся лишь флагом is_current (старая запись не была проставлена как неактуальная при добавлении новой). Условие JOIN не фильтрует по is_current, и в результате каждый заказ в рублях сопоставляется не с одной, а с двумя строками справочника - типичный fan-out, концептуально близкий к проблеме «множественного совпадения» из восьмого урока (там она ловилась cardinality check внутри MERGE INTO, но только для условий WHEN MATCHED, которые реально что-то обновляют по конкретному совпадению ключа - агрегирующий JOIN внутри подзапроса этой защитой не покрывается, потому что cardinality check Iceberg проверяет совпадения между target и source самого MERGE INTO, а не любые внутренние JOIN внутри источника).
spark.sql("""
MERGE INTO lakehouse.gold.daily_revenue_report t
USING (
SELECT
CAST(o.updated_at AS DATE) AS report_date,
COUNT(*) AS total_orders,
SUM(o.amount * r.rate) AS total_revenue
FROM lakehouse.analytics.orders o
JOIN lakehouse.ref.currency_rates r
ON o.currency = r.currency_code -- БАГ: пропущено "AND r.is_current = true"
WHERE o.status != 'cancelled'
GROUP BY CAST(o.updated_at AS DATE)
) s
ON t.report_date = s.report_date
WHEN MATCHED THEN UPDATE SET
t.total_orders = s.total_orders,
t.total_revenue = s.total_revenue
""")
spark.sql("""
SELECT * FROM lakehouse.gold.daily_revenue_report ORDER BY report_date
""").show(truncate=False)
# +-----------+------------+-------------+
# |report_date|total_orders|total_revenue|
# +-----------+------------+-------------+
# |2025-06-14 |1684 |368641.00 |
# |2025-06-15 |1830 |402901.50 |
# +-----------+------------+-------------+
Оба показателя - total_orders и total_revenue - выросли ровно вдвое по сравнению с Кейсом 0, потому что каждый заказ в рублях посчитан дважды из-за дублирующейся строки справочника. В 09:15, спустя пятнадцать минут после коммита, на утреннем совещании финансовый аналитик замечает, что выручка за 15 июня внезапно превысила сумму, физически невозможную при известном среднем чеке и количестве клиентов - инцидент обнаружен.
Кейс 2: исторический аудит через системные таблицы¶
Первый шаг расследования - не угадывать, а посмотреть на фактическую историю снапшотов витрины и найти точку, предшествующую инциденту.
spark.sql("""
SELECT
snapshot_id, parent_id, committed_at, operation,
summary['total-records'] AS total_records
FROM lakehouse.gold.daily_revenue_report.snapshots
ORDER BY committed_at
""").show(truncate=False)
# +--------------------+--------+-----------------------+---------+------------+
# |snapshot_id |parent_id|committed_at |operation|total_records|
# +--------------------+--------+-----------------------+---------+------------+
# |9001 |NULL |2025-06-15 09:00:00.000|overwrite|2 |
# |9002 |9001 |2025-06-15 09:10:00.000|overwrite|2 |
# +--------------------+--------+-----------------------+---------+------------+
Два снапшота: 9001 - чистое исходное состояние из Кейса 0, 9002 - тот самый испорченный MERGE INTO из Кейса 1, чей parent_id прямо указывает на 9001 как на предшественника. Расследование подтверждает гипотезу: нужный «чистый» снапшот - 9001, и он находится в прямой линии предков текущего 9002, то есть подходит для rollback_to_snapshot без каких-либо дополнительных условий.
Прежде чем откатывать таблицу для всех, разумно сначала прочитать этот исторический снапшот и убедиться, что данные внутри него действительно корректны - то есть применить Time Travel как read-only инструмент проверки гипотезы перед write-операцией.
-- SQL: читаем витрину такой, какой она была до инцидента
SELECT * FROM lakehouse.gold.daily_revenue_report VERSION AS OF 9001
ORDER BY report_date;
-- +-----------+------------+-------------+
-- |report_date|total_orders|total_revenue|
-- +-----------+------------+-------------+
-- |2025-06-14 |842 |184320.50 |
-- |2025-06-15 |915 |201450.75 |
-- +-----------+------------+-------------+
# PySpark: то же самое чтение через DataFrameReader API
clean_df = (
spark.read.format("iceberg")
.option("snapshot-id", 9001)
.load("lakehouse.gold.daily_revenue_report")
)
clean_df.orderBy("report_date").show(truncate=False)
# Идентичный результат - 842/184320.50 и 915/201450.75
Оба запроса подтверждают: данные «до инцидента» физически живы и доступны, несмотря на то что текущая (без AS OF) версия таблицы всё ещё показывает испорченные цифры из Кейса 1. Чтобы превратить это наблюдение в наглядный отчёт для финансовой команды, удобно одним запросом сравнить текущее (испорченное) состояние с чистым историческим - тот же паттерн self-join с разными AS OF на разных сторонах, что был показан в разделе про синтаксис Spark SQL:
spark.sql("""
SELECT
cur.report_date,
cur.total_revenue AS revenue_now_corrupted,
old.total_revenue AS revenue_before_incident,
cur.total_revenue - old.total_revenue AS diff
FROM lakehouse.gold.daily_revenue_report cur
JOIN lakehouse.gold.daily_revenue_report VERSION AS OF 9001 AS old
ON cur.report_date = old.report_date
ORDER BY cur.report_date
""").show(truncate=False)
# +-----------+---------------------+-----------------------+---------+
# |report_date|revenue_now_corrupted|revenue_before_incident|diff |
# +-----------+---------------------+-----------------------+---------+
# |2025-06-14 |368641.00 |184320.50 |184320.50|
# |2025-06-15 |402901.50 |201450.75 |201450.75|
# +-----------+---------------------+-----------------------+---------+
Колонка diff показывает точную величину искажения по каждой дате - в данном случае она равна самой величине корректной выручки, что численно подтверждает гипотезу «каждая сумма задвоена», ещё до того, как кто-либо полез разбираться в код MERGE INTO построчно.
Кейс 3: rollback - возвращаем витрину в исходное состояние¶
Гипотеза подтверждена, корректный снапшот найден и проверен через Time Travel. Переходим к официальному восстановлению - откату таблицы, который повлияет на всех читателей, включая dashboard финансовой команды.
CALL lakehouse.system.rollback_to_snapshot('gold.daily_revenue_report', 9001);
-- +----------------------+----------------------+
-- |previous_snapshot_id |current_snapshot_id |
-- +----------------------+----------------------+
-- |9002 |9001 |
-- +----------------------+----------------------+
Проверяем, что для обычного, ничем не помеченного запроса (без AS OF) таблица вновь выглядит так же, как до инцидента:
spark.sql("""
SELECT * FROM lakehouse.gold.daily_revenue_report ORDER BY report_date
""").show(truncate=False)
# +-----------+------------+-------------+
# |report_date|total_orders|total_revenue|
# +-----------+------------+-------------+
# |2025-06-14 |842 |184320.50 |
# |2025-06-15 |915 |201450.75 |
# +-----------+------------+-------------+
От момента обнаружения аномалии на совещании до восстановления корректной витрины прошло заведомо меньше времени, чем заняло бы поднятие резервной копии - сама операция rollback_to_snapshot выполнилась за доли секунды, поскольку, как разобрано в теоретической части урока, она лишь переключает один указатель в metadata.json.
Чтобы убедиться, что под капотом действительно ничего, кроме указателя, не изменилось, инспектируем history до и после отката:
spark.sql("""
SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor
FROM lakehouse.gold.daily_revenue_report.history
ORDER BY made_current_at
""").show(truncate=False)
# +-----------------------+-----------+--------+--------------------+
# |made_current_at |snapshot_id|parent_id|is_current_ancestor|
# +-----------------------+-----------+--------+--------------------+
# |2025-06-15 09:00:00.000|9001 |NULL |true |
# |2025-06-15 09:10:00.000|9002 |9001 |false |
# |2025-06-15 09:18:30.000|9001 |NULL |true |
# +-----------------------+-----------+--------+--------------------+
Три строки рассказывают всю историю переключений указателя current-snapshot-id, и это именно history (snapshot-log), а не snapshots - количество объектов в snapshots всё ещё равно двум (9001 и 9002), как было до отката, что можно проверить отдельным запросом к snapshots и убедиться, что новых строк там не появилось. Третья строка history - это и есть запись, добавленная вызовом rollback_to_snapshot: тот же snapshot_id = 9001, но новое значение made_current_at. Колонка is_current_ancestor для 9002 теперь показывает false - снапшот всё ещё существует в истории и доступен через Time Travel (VERSION AS OF 9002 продолжает работать и вернёт испорченные цифры, что полезно для дальнейшего разбора инцидента), но он больше не входит в линию предков текущего состояния таблицы.
Чтобы наглядно убедиться, что снапшот 9002 действительно не удалён - то есть что rollback это git checkout, а не git reset --hard - можно «передвинуть» указатель обратно на него через set_current_snapshot, а затем откатиться снова, доказывая, что вся операция полностью обратима в обе стороны без единого физического изменения файлов:
-- "Redo" - вручную возвращаем указатель на испорченный снапшот,
-- чтобы доказать, что он физически никуда не пропал
CALL lakehouse.system.set_current_snapshot(
'gold.daily_revenue_report', snapshot_id => 9002
);
-- Подтверждение: испорченные цифры снова видны как ТЕКУЩЕЕ состояние
SELECT * FROM lakehouse.gold.daily_revenue_report ORDER BY report_date;
-- 2025-06-14 | 1684 | 368641.00 <- та самая испорченная версия, цела
-- Возвращаем всё в порядок ещё раз
CALL lakehouse.system.rollback_to_snapshot('gold.daily_revenue_report', 9001);
Этот эксперимент с «redo» в реальном production-инциденте обычно не нужен (после подтверждённого отката незачем снова показывать пользователям испорченные данные), но как учебная демонстрация он закрывает главный концептуальный вывод всего раздела про rollback: ни одна операция Time Travel или rollback не удаляет ни единого байта данных - они лишь управляют тем, какой снапшот считается «текущим» в эту секунду, а физическое удаление выполняется только явным вызовом expire_snapshots, который в этом сценарии вообще не вызывался.
Производственный кейс: часовой пояс, который сломал воспроизводимость ML-эксперимента¶
Команда ML построила модель оценки риска отмены заказа, обучив её на снимке таблицы lakehouse.analytics.orders по состоянию на 1 мая. Аналитик, готовивший датасет в интерактивном ноутбуке, использовал ровно тот паттерн, что был показан в разделе про синтаксис Spark SQL этого урока:
SELECT * FROM lakehouse.analytics.orders
TIMESTAMP AS OF '2025-05-01 00:00:00';
Модель показала хорошие метрики на отложенной выборке и была одобрена для прода. Три месяца спустя комплаенс-команда запросила точное воспроизведение обучающего датасета для аудита - стандартная процедура для моделей, влияющих на финансовые решения. Дата-инженер, не имевший доступа к исходному интерактивному ноутбуку аналитика, реализовал то же самое чтение внутри регламентного Airflow-пайплайна, который запускает PySpark-скрипт на отдельном кластере:
df_for_audit = spark.sql("""
SELECT * FROM lakehouse.analytics.orders
TIMESTAMP AS OF '2025-05-01 00:00:00'
""")
print(df_for_audit.count())
# 412 311 - на 1 873 строки больше, чем было в исходном обучающем датасете
Количество строк не совпало. Дальше команда потратила почти целый рабочий день, подозревая что угодно - от багов в самом Iceberg до повреждения метаданных - прежде чем кто-то заметил очевидную, но легко упускаемую деталь: интерактивная сессия аналитика была сконфигурирована с spark.sql.session.timeZone = 'Europe/Moscow' (стандартная настройка, унаследованная от шаблона корпоративного JupyterHub), тогда как кластер, на котором запускался регламентный Airflow-пайплайн, использовал настройку по умолчанию UTC. Строковый литерал '2025-05-01 00:00:00', синтаксически идентичный в обоих случаях, был интерпретирован Spark как два разных момента физического времени, разнесённых на три часа (2025-04-30 21:00:00 UTC против 2025-05-01 00:00:00 UTC) - а в этом трёхчасовом окне в таблицу успел закоммититься плановый ночной batch-загрузчик, добавивший именно те самые 1873 строки.
Диаграмма явно разводит по двум параллельным веткам один и тот же исходный текст SQL-запроса, демонстрируя, что результат разошёлся не из-за разницы в коде, а из-за разницы в окружении, в котором этот код выполнялся - ровно тот риск, о котором предупреждал раздел про адресацию по timestamp в начале урока, но здесь показанный как реальные потерянные часы работы, а не абстрактная оговорка.
| Параметр | Сессия аналитика (исходное обучение) | Сессия Airflow (попытка аудита) |
|---|---|---|
spark.sql.session.timeZone |
Europe/Moscow (UTC+3) |
UTC (значение по умолчанию) |
Резолвинг '2025-05-01 00:00:00' |
2025-04-30 21:00:00 UTC |
2025-05-01 00:00:00 UTC |
| Найденный снапшот | Snapshot X (до ночной загрузки) | Snapshot Y (после ночной загрузки) |
| Строк в датасете | 410 438 | 412 311 |
| Время на диагностику расхождения | - | ~7 часов |
Меры по итогам инцидента:
-
Немедленная мера: исходный
snapshot_id, использованный для обучения модели, был восстановлен вручную черезhistoryпо приблизительному времени запуска исходного эксперимента (зафиксированному в логах MLflow), и именно он зафиксирован как официальный источник для дальнейшего аудита - ровно тот сценарий, ради которого нужна была функцияresolve_snapshot_as_of, показанная в разделе про динамическую параметризацию. -
Структурное исправление: для всех ML-пайплайнов введено обязательное правило - обучающий датасет фиксируется только через
snapshot_id, никогда через timestamp напрямую в коде эксперимента, по образцуrun_manifest, разобранного ранее в этом уроке. Самsnapshot_idлогируется в MLflow как параметр эксперимента наравне с гиперпараметрами модели. -
Архитектурное уточнение: на уровне кластерных шаблонов зафиксирован единый
spark.sql.session.timeZone = 'UTC'для всех сред - интерактивных ноутбуков и регламентных кластеров одинаково, чтобы устранить сам источник несовпадения для любого кода, который по какой-то причине всё же продолжит использовать timestamp-литералы (например, в ad hoc расследованиях, где snapshot_id ещё не известен).
Этот инцидент - не повод отказываться от адресации по timestamp в принципе: как было сказано в начале урока, она остаётся лучшим инструментом первого контакта с проблемой, когда snapshot_id ещё не известен. Урок инцидента точнее: timestamp хорош для исследования, но прежде чем что-либо зафиксировать как воспроизводимый артефакт (обучающий датасет, отчёт для аудита, конфигурацию пайплайна), приблизительный timestamp обязан быть конвертирован в точный snapshot_id - и именно этот snapshot_id, а не исходная человекочитаемая дата, должен путешествовать дальше по системе.
Типичные заблуждения¶
«Rollback - это как git reset --hard, старые снапшоты после него теряются» - неверно, и это, возможно, самое важное заблуждение всего урока. Раздел про механику rollback и Кейс 3 практического блока показали прямо: rollback_to_snapshot не удаляет ни одного снапшота из массива snapshots - он только переключает указатель current-snapshot-id и добавляет запись в snapshot-log. Снапшоты, оставшиеся «впереди» по времени относительно нового текущего, продолжают существовать и доступны через Time Travel и через set_current_snapshot, пока их явно не удалит expire_snapshots. Корректная git-аналогия - не reset --hard, а checkout на старый коммит без удаления более новых коммитов из истории репозитория.
«Time Travel создаёт дополнительную нагрузку на хранение - чем активнее его использовать, тем больше платим за storage» - перепутана причина со следствием. Сам факт использования Time Travel - выполнение запроса с VERSION AS OF/TIMESTAMP AS OF - не создаёт никаких новых файлов и не увеличивает занятое место ни на байт: это обычная read-операция над уже существующими данными. Реальная стоимость возникает не от использования Time Travel, а от решения не удалять историю как можно дольше, то есть от настройки горизонта хранения и частоты запуска expire_snapshots - вопрос, полностью независимый от того, читает ли кто-либо реально эту историю через Time Travel или нет.
«TIMESTAMP AS OF возвращает снапшот, ближайший по времени к указанному моменту» - неточно и потенциально опасно как ошибочное предположение. Раздел про адресацию по timestamp показал на конкретном примере: движок всегда ищет последний снапшот, не позже указанного момента, и никогда не вернёт снапшот, созданный после запрошенной точки, даже если он находится к ней ближе по модулю разницы во времени, чем подходящий более ранний снапшот.
«Опция as-of-timestamp в PySpark DataFrameReader принимает такую же строку, как TIMESTAMP AS OF в SQL» - неверно и является источником реальных ошибок при переносе кода между SQL и PySpark API. Раздел про DataFrameReader явно показал: as-of-timestamp ожидает целое число миллисекунд с начала эпохи Unix, а не человекочитаемую строку - необходимо явное преобразование через datetime.timestamp() * 1000, и желательно - в одной общей, протестированной функции, а не повторяемое вручную в каждом месте кода.
«rollback_to_snapshot может откатить таблицу как назад, так и вперёд - на любой существующий snapshot_id» - неверно. Процедура работает только с предками текущего снапшота и завершается явной ошибкой валидации при попытке указать снапшот из параллельной ветки истории либо снапшот, уже не являющийся предком после более раннего отката. Для перемещения указателя на произвольный снапшот без этого ограничения существует отдельная, менее защищённая процедура set_current_snapshot, разобранная отдельно именно из-за этой разницы в гарантиях.
«Time Travel - это фича только Spark SQL, в PySpark DataFrame API такого нет» - неверно. Раздел про PySpark API показал, что DataFrameReader поддерживает Time Travel через опции snapshot-id, as-of-timestamp, tag и branch, причём программный доступ из PySpark особенно важен именно для автоматизированных пайплайнов, для которых SQL-литералы внутри строки запроса менее удобны и менее тестируемы, чем явные параметры функции.
«После expire_snapshots ещё можно восстановить удалённый снапшот средствами самого Iceberg, просто это будет медленнее» - неверно и опасно как предположение, на которое можно по ошибке опереться в продакшене. Раздел про горизонт хранения прямо подчеркнул: после физического удаления файлов данных и манифестов expire_snapshots восстановление этого конкретного снапшота средствами Iceberg невозможно в принципе - снапшот не «временно недоступен», а навсегда удалён из графа метаданных, и единственный путь назад в такой ситуации - внешний бэкап object storage, то есть ровно тот медленный путь восстановления, которого Time Travel изначально позволяет избежать, пока история не вычищена.
Производственный чек-лист: что нельзя забыть до того, как положиться на Time Travel в production¶
-
Свойства
history.expire.max-snapshot-age-msиhistory.expire.min-snapshots-to-keepзаданы осознанно для каждой таблицы исходя из реальных требований её потребителей (аудит, ML-репродуцируемость, типичное время обнаружения инцидента), а не оставлены на дефолтные 5 дней «потому что так настроено по умолчанию». -
Воспроизводимые артефакты (обучающие датасеты, регламентные отчёты для аудита) фиксируют
snapshot_id, а не приблизительный timestamp - по образцуrun_manifest, разобранного в разделе про динамическую параметризацию, и независимо от того, насколько удобным кажется человекочитаемый timestamp на момент первого запуска эксперимента. -
spark.sql.session.timeZoneзафиксирован единообразно для всех сред, где может выполняться код сTIMESTAMP AS OFили ручным резолвингом исторических моментов - интерактивные ноутбуки и регламентные кластеры должны использовать одну и ту же настройку, что напрямую устраняет класс инцидентов, разобранный в производственном кейсе этого урока. -
Контрактно значимые точки времени (конец квартала, момент перед миграцией) зафиксированы через
CREATE TAG ... RETAIN N DAYS, а не полагаются на то, что кто-то когда-нибудь вспомнит нужный timestamp и успеет найти соответствующий снапшот до того, как его вычиститexpire_snapshots. -
Команда дежурных инженеров понимает разницу между
rollback_to_snapshotиset_current_snapshotи имеет под рукой runbook с готовыми шаблонами обеих команд - решение о том, какую процедуру использовать, не должно приниматься впервые в момент реального инцидента под давлением времени. -
Возраст самого старого активного снапшота и общее число снапшотов вынесены в мониторинг - алерт, срабатывающий до того, как ожидаемый горизонт Time Travel для конкретной таблицы (например, «нам нужна возможность аудита за последний квартал») окажется под угрозой из-за давно не запускавшегося или слишком агрессивного
expire_snapshots. -
Расписание
expire_snapshotsсверено с самым длинным требованием к retention среди реальных потребителей таблицы, а не настроено изолированно командой эксплуатации без учёта потребностей ML и аналитики - именно этот последний пункт прямо открывает тему следующего урока модуля.
Последний пункт чек-листа заслуживает отдельного развёрнутого комментария, потому что именно он формирует мостик к следующей теме модуля. Весь этот урок последовательно показывал, что возможность путешествовать в историю таблицы - это не данность, а результат сознательного решения не удалять старые снапшоты дольше определённого срока. Решение о том, когда и насколько агрессивно эту историю всё-таки вычищать - это отдельная инженерная задача, требующая баланса между стоимостью хранения (раздел про горизонт хранения этого урока, прямо опирающийся на Write Amplification из шестого урока) и реальными потребностями в Time Travel, аудите и rollback, разобранными здесь. Эта задача - предмет следующего урока модуля.
Мостик к следующим урокам¶
Этот урок показал, как накопленная история снапшотов - побочный продукт каждой операции записи, разобранной в шестом, седьмом и восьмом уроках модуля - превращается в первоклассный инструмент аудита, воспроизводимости и disaster recovery. Следующие уроки модуля строят непосредственно на этом материале.
-
Урок 10 («Table Maintenance») - прямое продолжение последнего раздела этого урока: чек-лист и раздел про горизонт хранения многократно подчёркивали, что Time Travel работает только до тех пор, пока
expire_snapshotsне вычистил нужный снапшот, а Write Amplification CoW-операций прямо умножает реальную стоимость хранения этой истории. Этот урок разберёт точный синтаксисCALL-процедурexpire_snapshotsиrewrite_data_files, стратегии планирования компакции и то, как безопасно балансировать между стоимостью хранения истории и реальными потребностями в Time Travel, описанными здесь. -
Уроки 11-12 («Delta Lake») покажут, как тот же принцип путешествия во времени реализован в альтернативном table format - команды
RESTORE TABLEи... VERSION AS OF/TIMESTAMP AS OFсинтаксически и концептуально очень близки к материалу этого урока, но физически опираются не на дерево манифестов, а на Delta Log (последовательность JSON/Parquet-файлов транзакций), что меняет детали механики поиска исторической версии, хотя не меняет пользовательскую модель «читать прошлое без копирования данных». -
Урок 13 («Выбор формата») включит зрелость и эргономику Time Travel и rollback как один из критериев сравнительной матрицы Iceberg vs Delta Lake vs Apache Hudi - в частности, то, как разные форматы по-разному балансируют стоимость хранения истории, гибкость адресации (snapshot_id, timestamp, tag/branch, версия Delta Log) и интеграцию с инструментами ML-репродуцируемости, разобранными в этом уроке.
Домашнее задание¶
-
Воспроизведите Кейсы 0-3 практического блока на собственном self-hosted стенде, но замените сценарий катастрофы: вместо дублирующейся строки справочника валют смоделируйте ошибочный
UPDATEбезWHERE, затрагивающий всю таблицу разом. Найдите снапшот до инцидента черезhistory, прочитайте его через Time Travel и выполнитеrollback_to_snapshot, зафиксировав в отчёте все промежуточныеsnapshot_id. -
Эмпирически подтвердите утверждение «rollback занимает O(1) независимо от объёма данных»: создайте три версии одной и той же таблицы (1 тыс., 1 млн, 100 млн строк), выполните на каждой одинаковую последовательность операций и замерьте время выполнения
rollback_to_snapshotна каждой. Постройте график зависимости времени отката от объёма данных и письменно объясните результат через механику, разобранную в разделе про rollback под капотом. -
Доработайте функцию
resolve_snapshot_as_ofиз раздела про динамическую параметризацию так, чтобы она явно и информативно обрабатывала три граничных случая: запрошенный момент раньше первого снапшота таблицы, таблица не содержит ни одного снапшота, и токен времени передан без явного часового пояса. Напишите unit-тесты на каждый из этих трёх случаев. -
Воспроизведите программно разницу между
rollback_to_snapshotиset_current_snapshot: создайте таблицу с веткой (CREATE BRANCH), сделайте независимые коммиты вmainи в новую ветку, затем попытайтесь выполнитьrollback_to_snapshot('main')на snapshot_id из ветки. Зафиксируйте текст возникающей ошибки валидации, а затем покажите, чтоset_current_snapshotс тем же snapshot_id отрабатывает без ошибки, и письменно объясните, почему это различие в гарантиях оправдано с точки зрения безопасности операции. -
Спроектируйте и реализуйте мини-runbook (Markdown или Jupyter-ноутбук) для дежурного инженера: по входному приблизительному времени инцидента находить кандидатов в
history, выводить diff между текущим состоянием и состоянием до каждого кандидата (по образцу запроса из Кейса 2), и только после явного подтверждения человеком вызыватьrollback_to_snapshot. Обсудите письменно, какие самостоятельные действия скрипта было бы опасно оставлять без подтверждения человеком. -
Воспроизведите производственный кейс этого урока локально: запустите два Spark-сессии с разными значениями
spark.sql.session.timeZone(например,Europe/MoscowиUTC), выполните в обеих один и тот же запросTIMESTAMP AS OFс одинаковым строковым литералом на таблице, в истории которой есть коммит внутри окна расхождения часовых поясов, и подтвердите программно (через сравнениеsnapshot_idрезультатов), что запросы резолвятся в разные снапшоты. -
Спроектируйте письменно (1-2 абзаца) конкретную метрику и пороговое значение для алертинга, которые предупредили бы команду о том, что возраст самого старого активного снапшота таблицы приближается к границе, требуемой для аудита (например, требование «хранить минимум 90 дней истории для compliance-таблиц»), до того как
expire_snapshotsпотенциально вычистит нужный диапазон. -
Напишите unit-тест (PySpark +
pytest, локальный Iceberg-каталог), который выполняет следующую последовательность и проверяет инвариант на каждом шаге: создаёт таблицу, делает три коммита, откатывается ко второму черезrollback_to_snapshot, читает данные и сравнивает их (exceptAllв обе стороны, ожидая пустой результат) с заранее закэшированнымDataFrame, прочитанным черезVERSION AS OFдо отката. -
Исследуйте на практике, как меняются
snapshotsиhistoryпосле нескольких подряд вызововset_current_snapshotв разные стороны (например: откат, redo, повторный откат). Подтвердите письменно и кодом, что число строк вsnapshotsостаётся постоянным на всём протяжении эксперимента, в то время как число строк вhistoryрастёт на одну запись при каждом вызове. -
Сравните письменно (без необходимости полной реализации) Time Travel как точечное чтение (
snapshot-id/as-of-timestamp) и инкрементальное чтение (start-snapshot-id/end-snapshot-id) для гипотетической задачи сверки данных между двумя системами раз в сутки: какой из двух инструментов более уместен и почему, и что изменится в выборе, если сверка должна обнаруживать не только новые строки, но иUPDATE/DELETE, прошедшие черезMERGE INTOвосьмого урока.
Полная картина: жизненный цикл одного запроса Time Travel и одного Rollback¶
Завершая урок, соберём весь материал в единую диаграмму - от накопления истории снапшотов обычными операциями записи до момента, когда эта история либо служит инструментом расследования, либо физически упирается в горизонт хранения, замыкая цикл на тему следующего урока модуля.
Диаграмма проходит ровно тот путь, что был разобран в практическом демо-блоке этого урока, но в обобщённом виде, применимом к любому инциденту, а не только к конкретному примеру с витриной выручки. Левая часть (WRITE → INCIDENT → DETECT → HIST) - это то, что происходит независимо от Time Travel: обычная работа таблицы, инцидент и его обнаружение. Развилка ADDR отражает практический паттерн из раздела про способы адресации: расследование почти всегда начинается с приблизительного timestamp и только потом сужается до точного snapshot_id, зафиксированного через history. После того как нужный снапшот найден, он сначала используется только для чтения (READ → VERIFY) - это самый безопасный шаг, не влияющий на других пользователей таблицы, и именно на этом шаге гипотеза о причине инцидента подтверждается или отвергается. Только после подтверждения наступает развилка ANCESTOR, отражающая ограничение rollback_to_snapshot, разобранное в разделе про команды управления: если целевой снапшот - предок текущего (типичный случай простого «вернуться немного назад»), используется защищённая процедура rollback_to_snapshot; если нужно более нестандартное перемещение, используется менее защищённый set_current_snapshot. Обе процедуры сходятся в одной и той же физической операции - переключении указателя в metadata.json (PTR), которая не трогает сам массив снапшотов (NOCHANGE) - и именно это отсутствие изменений в snapshots[] объясняет, почему вся цепочка от обнаружения инцидента до его исправления может занять секунды, а не часы. Замыкающий диаграмму блок RETENTION - явное напоминание о том, что весь этот механизм работает только в пределах горизонта хранения истории, который требует осознанного управления - тема следующего урока модуля.
Итоги¶
Time Travel не требует отдельной инфраструктуры хранения истории - он целиком построен на том, что Snapshot Model уже не удаляет старые файлы и манифесты в момент коммита. Цена этой фичи - не в её использовании (любое количество запросов с AS OF не создаёт новых файлов), а в решении не удалять историю дольше необходимого, что является отдельным вопросом конфигурации горизонта хранения.
Iceberg поддерживает три способа адресации исторических данных - snapshot_id, timestamp и именованные ссылки (tag/branch) - с разными trade-off'ами между удобством для человека и надёжностью в автоматизации. Практический паттерн: timestamp для первого контакта с проблемой, snapshot_id для точной и воспроизводимой фиксации найденного решения.
Адресация по TIMESTAMP AS OF всегда ищет последний снапшот, ставший текущим не позже запрошенного момента - никогда не «ближайший» и никогда из будущего относительно запрошенной точки. Эта механика зависит от spark.sql.session.timeZone, рассинхронизация которого между средами - реальный, а не теоретический источник производственных инцидентов воспроизводимости, разобранный в этом уроке.
PySpark DataFrameReader зеркалирует возможности SQL-синтаксиса через опции snapshot-id, as-of-timestamp, tag, branch, но с важной несовместимостью форматов: as-of-timestamp ожидает epoch-миллисекунды, а не человекочитаемую строку. Перенос кода между SQL и PySpark API без учёта этой разницы - частый источник ошибок при автоматизации исторических выгрузок.
Rollback - это write-операция, принципиально отличная от read-only Time Travel, но под капотом она не создаёт нового снапшота, а лишь переключает указатель current-snapshot-id в новой версии metadata.json и добавляет запись в snapshot-log. Массив snapshots не меняется; именно поэтому rollback занимает константное время независимо от объёма данных в таблице, что прямо подтверждено замерами в практическом блоке урока.
rollback_to_snapshot/rollback_to_timestamp работают только с предками текущего снапшота и завершаются ошибкой при попытке нарушить это ограничение; для произвольного перемещения указателя без этой защиты существует set_current_snapshot. Различие в гарантиях между этими процедурами должно быть явно зафиксировано в runbook дежурной команды, а не выбираться импровизированно во время реального инцидента.
Стоимость хранения истории напрямую зависит от режима записи: для Copy-on-Write с частыми точечными изменениями длинный горизонт хранения умножает Write Amplification из шестого урока на число удерживаемых снапшотов, тогда как для Merge-on-Read та же история занимает кратно меньше места. Это инженерное решение, требующее осознанного выбора history.expire.max-snapshot-age-ms, а не дефолтного значения по умолчанию.
После выполнения expire_snapshots Time Travel и rollback к удалённому снапшоту становятся невозможны навсегда средствами самого Iceberg. Это прямой мостик к следующему уроку модуля: горизонт хранения, разобранный здесь как политика, там превращается в конкретную, исполняемую процедуру со своими стратегиями планирования.
Краткий глоссарий терминов урока¶
-
Time Travel - возможность прочитать таблицу Iceberg в состоянии любого из сохранённых снапшотов её истории, без копирования данных и без влияния на текущее состояние таблицы для других читателей.
-
snapshot_id- уникальный 64-битный идентификатор конкретного снапшота, самый строгий и однозначный способ адресации к историческому состоянию таблицы. -
VERSION AS OF/TIMESTAMP AS OF- синтаксис Spark SQL для Time Travel внутри блокаFROM, адресующий снапшот по идентификатору/имени ссылки или по времени соответственно. -
FOR SYSTEM_VERSION AS OF/FOR SYSTEM_TIME AS OF- ANSI SQL-эквивалентный синтаксис тех же операций, поддерживаемый Spark начиная с версии 3.3, без какой-либо разницы в результате или физическом плане. -
as-of-timestamp- опция PySparkDataFrameReaderдля адресации по времени; в отличие от SQL-литерала, ожидает целое число миллисекунд с начала эпохи Unix, а не человекочитаемую строку. -
Инкрементальное чтение (
start-snapshot-id/end-snapshot-id) - запрос строк, добавленных между двумя снапшотами через операции типаappend; решает иную задачу, чем точечный Time Travel, и не подходит для воспроизведения полного состояния таблицы. -
Rollback - операция возврата таблицы к одному из предыдущих снапшотов в качестве нового текущего состояния для всех читателей; в отличие от Time Travel, изменяет состояние таблицы и требует коммита.
-
rollback_to_snapshot/rollback_to_timestamp- хранимые процедуры Spark SQL, выполняющие rollback к снапшоту, заданному явно или через резолвинг timestamp; работают только с предками текущего снапшота. -
set_current_snapshot- хранимая процедура, переключающаяcurrent-snapshot-idна произвольный существующий снапшот таблицы, включая снапшоты, не являющиеся предками текущего - в отличие отrollback_to_snapshot. -
Снапшот-предок (ancestor) - снапшот, находящийся в цепочке
parent-snapshot-idмежду корнем истории таблицы и текущим снапшотом; единственный допустимый тип цели дляrollback_to_snapshot. -
Горизонт хранения (retention) - период, за который снапшоты считаются кандидатами на удаление при следующем запуске
expire_snapshots, задаваемый свойствамиhistory.expire.max-snapshot-age-msиhistory.expire.min-snapshots-to-keep. -
is_current_ancestor- колонка системной таблицыhistory, показывающая, входит ли данный исторический снапшот в линию предков текущего состояния таблицы на момент запроса; меняется наfalseдля снапшотов, «обойдённых» откатом. -
Run-манифест (run manifest) - практика фиксации
snapshot_idисходных таблиц в виде структурированной записи, сохраняемой рядом с артефактом пайплайна (модель, отчёт), для точного воспроизведения входных данных в будущем без зависимости от timestamp и часовых зон. -
Tag / Branch - именованные неизменяемая (
tag) и изменяемая (branch) ссылки на снапшоты, представленные во втором уроке модуля; в контексте Time Travel служат третьим, человекочитаемым способом адресации наравне сsnapshot_idиtimestamp.