TGViewer
дата инженеретта дата инженеретта @data_engineerette · 3.43K subscribers
Post #494 2.26K
Управляем ошибками. Часть 2

📼 Late Data

Посмотрим на 3 паттерна для работы с данными, которые пришли позже, чем ожидалось

1️⃣Pattern: Late Data Detector

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

Интересная мысль - "shifting the late data problem"
Пример: мы пишем партиции по времени обработки. В партицию 21:00 к нам залетел кусок данных за 20:00 и 19:00. А наши пользователи используют партиции по времени события. Тогда мы перекладываем ответственность ковыряться в этих партициях на них 😁

2️⃣Pattern: Static Late Data Integrator

Как вообще можно перегрузить данные за прошлое?

1. Создать кучу дагранов, где каждый перегружает 1 день. Если упало - перезапускаем конкретный день

2. Создать один дагран, где в коде генерируется список нужных дат. И по каждой дате запускается загрузка. Если упало - просто перезапускаем, пойдет считаться с упавшего дня. Это и есть Static Late Data Integrator. А статическое - потому что мы сами задаем 14 дней или сколько угодно

И тут я поняла, что неосознанно это и делала. У нас часто была проблема, что данные в источники просто не приходили 😁 Потом мы шли разбираться с владельцами, и данные заливались, но позднее. Чтобы это учитывать, в моем подходе был такой алгоритм:

1. Задаем стартовую и конечную даты расчета
2. Создаем диапазон значений


full_range = pd.date_range(start=str(start_dt), end=str(end_dt)).strftime("%Y-%m-%d")


3. Из меты достаем существующие партиции


def get_existing_partitions(table_name):
partitions = (
spark.sql(f"show partitions {table_name}")
.select(F.split(F.col("partition"), "=")[1].alias("dt"))
.collect()
)

return [p[0] for p in partitions]


4. Находим разницу


lost_range = full_range.difference(existing_partitions_pdf)


5. Итерируемся по потеряшкам


for dt in lost_range:
calc_mart(dt)


Если в будущем снова будет пустая дата, нам не придется перезапускать определенный день - он пойдет считаться сам

3️⃣Pattern: Dynamic Late Data Integrator

Предлагается завести табличку с 4 полями:
🤩партиция
🤩время обработки
🤩время добавления новых записей
🤩флаг обработано или нет

Так мы запросом можем найти партиции, которые уже обрабатывались, но в которые попали новые данные. А в iceberg есть удобное свойство last_updated_at на уровне таблицы

🤩 Filtering

Pattern: Filter Interceptor

Как будто это антипаттерн. Предлагается создать доп колонки с фильтрами id_is_not_null, status_is_not_failed и выводить количество отфильтрованных записей, чтобы понимать, на каком этапе ошибка в коде или в данных. Но прям пробегаться по каждой записи в датафрейме… Как будто это все-таки dq

🌳 Fault Tolerance

Pattern: Checkpointer


Просто нужно создавать чекпоинты и хранить последний оффсет обработанной записи и состояние, если оно есть

Еще раз напомнили про семантики доставки:
🤩exactly once - нужны другие паттерны, расскажу, когда дойду
🤩at least once - чекпоинт после обработки, могут быть дубликаты при перезапуске после падения
🤩at most once - чекпоинт до обработки, данные потеряются при падении

#depatterns
  • ❤ 11
  • 👍 7
  • 🤔 1
More from @data_engineerette
  1. Sep 25, 2026Как прошла SmartData 2026? Я вот перечитываю свои впечатления от прошлого года и понимаю,…
  2. Sep 24, 2026Исследование data-people Тут ребята из DevCrowd запустили ежегодное исследование специалис…
  3. Sep 22, 20265 октября начнется 19-й поток программы Data Engineer от Newprolab Программа для junior- и…
  4. Sep 15, 2026Каким должен быть хороший DE? Меня однажды спросили на собесе: 🤩Какие 3 качества важны дл…
  5. Sep 12, 2026Mermaid-диаграммы Наконец-то дошли руки поковыряться в mermaid-диаграммах, это что за имба…
  6. Sep 2, 2026Iceberg — это внезапный бум или планомерная подготовка? Заметили, как с определенного моме…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →