Запусти дважды
Чему научишься
- собирать настоящий DAG:
@dag,@task, зависимость, заданная вызовом функции - читать протокол прогона: порядок задач, состояния, что попало на склад
- воспроизводить аварию тикета №01 у себя: один и тот же день, загруженный дважды
- делать шаг идемпотентным через перезапись партиции и доказывать это проверкой, а не обещанием
Тикет №01, часть третья. «Сделай так, чтобы не повторилось»
15:20. Суточный отчёт пересчитан руками, Академия успокоилась. В тикете появляется последняя строка — от смотрителя архива:
«Хорошо. Теперь сделай так, чтобы ночной ретрай больше не удваивал день. Ленту ведёшь ты».
Отчёт уже исправлен. Но это ещё не конец работы: нужно изменить код так, чтобы авария перестала быть возможной.
Для этого сначала воспроизведём её у себя. Не отдельной функцией, запущенной вручную, а маленькой настоящей лентой — с задачами, зависимостями и повторным прогоном.
Ленту запускает оркестратор. Он знает расписание, порядок задач и правила перезапуска. Его записи ты сегодня уже читал: журнал прогонов с попытками и ретраями ведёт именно он. Теперь соберёшь его ленту в коде.
Возьмём Airflow: его часто спрашивают на собеседованиях, а в уроке нам нужен его современный интерфейс.
В DAG-файле для Airflow 3 импорт выглядит так:
from airflow.sdk import dag, task
В наших ячейках вместо него стоит from arena_airflow import dag, task, run_dag. Причина одна: настоящий Airflow в браузере не поднять — ему нужны база , и веб-сервер.
Декораторы и связи задач ниже работают так же, как в Airflow. Когда перенесёшь такой DAG в реальный dags/, заменить нужно именно импорт.
Теперь разберём две основные сущности.
Таска — обычная функция под декоратором @task. DAG — функция под @dag, внутри которой таски ВЫЗЫВАЮТСЯ. Из этих вызовов Airflow строит зависимости.
Строка load(extract()) означает: «load идёт после extract и получает его результат». Отдельная стрелка здесь не нужна — передача данных уже задаёт порядок.
Явная стрелка first >> second нужна в другом случае: порядок между задачами есть, а передачи данных нет. Например: «сначала посчитай проверки, потом публикуй витрину».
Осталось одно понятие, без которого нельзя правильно перезапускать прошлые даты: logical_date.
Это ДАТА, ЗА КОТОРУЮ считаем, а не момент фактического запуска. Прогон за 12 марта может случиться ночью 13-го или через неделю при разборе аварии. В обоих случаях он обязан посчитать одно и то же.
Теперь соберём ленту из двух задач и запустим её за одну дату дважды — ровно то, что ночью сделал ретрай.

extract забирает окно у порта, load кладёт строки на склад. После двух запусков сравни не только состояния задач, но и данные. Механику порта и учебного склада здесь разбирать не нужно: следи за результатами прогонов. У задач одинаковые состояния, но строк на складе вдвое больше. Значит, зелёный DAG и правильные данные — не одно и то же.Голый INSERT — это не загрузка
Сорок строк превратились в восемьдесят. При этом ни одна задача не упала: оба прогона DAG закончились успешно.
Ровно это и произошло ночью 13 марта. Только вместо твоей кнопки сработал автоматический ретрай, а вместо сорока строк задвоился день выручки.
КВЕРИ: Зелёный прогон означает, что шаг дошёл до конца. Не то, что он поступил правильно. Разницу между этими двумя утверждениями тебе и платят замечать.
Почему так произошло? Причина в одной строке: db.insert(...) дописывает данные к уже лежащим на складе.
Для события, которое гарантированно происходит один раз, это нормальная операция. Но загрузка за дату — не такое событие. Один и тот же день загрузят столько раз, сколько запустят шаг.
Повторные запуски неизбежны: оркестраторы умеют автоматически перезапускать задачи после сбоя. В Airflow ретраи можно включить строкой default_args={'retries': 2}. По умолчанию ретраев НЕТ: у таски retries = 0, пока их не настроили.
Но отсутствие автоматического ретрая проблему не решает. То, что ночью не перезапустит оркестратор, утром перезапустит человек, разбирающий аварию.
Поэтому нам нужно другое свойство шага — идемпотентность. Сколько бы раз шаг ни выполнили с одними и теми же входами, итоговое состояние должно остаться тем же.
Важно: идемпотентность не означает «повторный запуск не упадёт». Она означает, что состояние склада после второго прогона неотличимо от состояния после первого.
Чтобы получить это свойство, меняем единицу записи. Загрузчик должен мыслить не отдельными строками, а партицией — куском данных, целиком принадлежащим одной дате.
Тогда каждый запуск делает два действия:
- снять всё, что относится к этой дате —
delete_partition(table, day); - положить заново то, что посчитали сейчас —
insert(table, rows).
Посмотрим, что это меняет. Первый запуск удаляет ноль строк и кладёт сорок. Второй удаляет сорок и кладёт сорок. Третий делает то же самое.
Итог после каждого запуска одинаковый. При этом соседние даты не затрагиваются: удаляется только СВОЯ партиция, а не таблица целиком.
Теперь становится понятнее и лекарство из de1l1. Задвоенный день нельзя было исправить простым перезапуском: шаг, который только дописывает, положил бы третий комплект строк. Нужна была именно перезапись диапазона.
Осталось не просто назвать шаг идемпотентным, а доказать это.
Комментарий «шаг идемпотентен» ничего не гарантирует — именно такой комментарий висел над загрузчиком предшественника. Проверка должна воспроизвести повторный запуск и сравнить состояние до и после него.
Поэтому доказательство здесь простое: два прогона подряд и сравнение снимков склада. Снимок — db.snapshot(table), хеш содержимого по алгоритму SHA-256. Если снимки до и после повтора совпали, повторный запуск действительно не изменил состояние склада.
INSERT 40 → 80 → 120 строк, у перезаписи своей партиции 40 → 40 → 40, а truncate даёт тот же итог и стирает соседний день.Что здесь настоящее, а что песочница. arena_airflow повторяет ПУБЛИЧНЫЙ интерфейс Airflow 3: DAG, @dag, @task, PythonOperator, стрелки >>, XCom, TaskGroup, сенсоры, logical_date, ds, retries и retry_delay. Весь список сейчас разбирать не нужно: в этом уроке важны декораторы, зависимости и дата прогона.
Эта часть переносится в настоящий dags/: достаточно заменить строку импорта на from airflow.sdk import dag, task.
Но вокруг DAG в песочнице есть учебная обвязка, которой в Airflow НЕ существует. При переносе её нужно выкинуть целиком: run_dag(...), d.structure(), run.summary() и склад arena_source.db (create_table, insert, delete_partition, snapshot, partitions).
В их роль играют команда airflow dags test, веб-интерфейс и твоя настоящая база. То есть переносится тело тасок, а не пульт, с которого мы запускаем их в уроке.
Есть и более крупные упрощения. У шима нет и не будет , базы , веб-интерфейса, пулов и слотов, executor-ов (Celery, Kubernetes), отложенных задач с triggerer, обнаружения зависших процессов и настоящего cron.
Расписание здесь считается по датам, а не по часам сервера. Таски идут строго по одной, в порядке зависимостей: параллельный исполнитель давал бы разный вывод при разных запусках. Ретраи тоже не спят по-настоящему — retry_delay копится в виртуальных часах прогона.
Поэтому параллелизм, ресурсы и очереди в этом курсе придётся принять на слово. Но семантика DAG, порядок, ретраи, идемпотентность и backfill — последовательный запуск за прошлый диапазон дат — здесь воспроизводятся по-настоящему. Именно их и спрашивают на собеседовании.
load.
После починки должны выполняться три условия:
- два прогона за 12 марта подряд оставляют на складе одинаковое число строк и одинаковый
db.snapshot('stg_orders'); - прогон за 13 марта ДОБАВЛЯЕТ свою партицию и не трогает партицию 12-го;
- все прогоны остаются зелёными.
db.delete_partition(table, partition): функция удаляет одну партицию и возвращает число удалённых строк. Ключ партиционирования объявлен как partition_by='day', а нужная дата лежит в payload['day'].Вопрос с собеседования
Как это спрашивают на собеседовании. «Что такое идемпотентная загрузка и как её сделать?»
Начни с определения через ПОВТОР: после любого числа запусков с теми же входами результат такой же, как после одного.
После определения ждут механику. Здесь нужна конкретика: перезапись партиции (DELETE диапазона + INSERT, или INSERT OVERWRITE, или MERGE по ключу). Перезаписывать нужно партицию, выбранную по логической дате прогона.
Отдельный балл — за фразу «загрузка ведётся по logical_date, а не по «сегодня»». Иначе ретрай за прошлую неделю возьмёт текущую дату и посчитает уже другой набор данных.
Дальше обычно спрашивают: «как вы доказываете идемпотентность?» Правильный ответ — тестом, который запускает шаг дважды и сравнивает итоговое состояние. Комментарий в коде ничего не доказывает.
Ещё один частый вопрос: «а если задача упала ПОСЛЕ записи, но до отметки об успехе?» Именно ради этого случая идемпотентность и нужна. Оркестратор запустит ретрай с начала шага, поэтому повтор должен привести склад к тому же состоянию.
load упала ПОСЛЕ того, как записала строки на склад, но до отметки об успехе. Оркестратор запускает ретрай. Что произойдёт, если шаг умеет только дописывать строки?load идемпотентным так: перед вставкой вызывает db.truncate('stg_orders'). Два прогона за 12 марта дают одинаковый результат, но загрузка должна сохранять другие даты. Что не так?Главное из урока
- DAG собран:
@dagоборачивает функцию,@task— отдельный шаг, а зависимость возникает из вызоваload(extract()). Стрелка>>нужна, когда порядок есть, а передачи данных нет. logical_date— дата, ЗА которую считаем, а не момент запуска. Прогон за 12 марта должен дать один и тот же результат и ночью 13-го, и через неделю.- Голый
INSERT— не загрузка: два зелёных прогона за одну дату дали 80 строк вместо 40. Так в примере воспроизводится авария тикета №01. - Лекарство — перезапись СВОЕЙ партиции:
delete_partition(table, day)передinsert. Неtruncate: он сотрёт соседние даты и сломает первый же backfill. - Идемпотентность нужно доказывать: два прогона подряд и сравнение
db.snapshot(). Комментарий в коде доказательством не считается.
Дальше по ленте — последний урок главы: что записывать в журнал дежурства и какими запросами утром проверять, что ночь прошла нормально.