📼 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