Оптимальный способ обработки чрезвычайно больших массивов данных (> 100 ТБ)
Главный data engineer influencer написал пост как он обрабатывал data set 100TB + (для особо современных подписчиках у него доступна версия в ТикТок). На картинке изображено дерево, как он обрабатывал каждый час и merge по 2 часа, потом по 4 часа и тп. Он классно обозначал проблемы: - data retention - cost of storage (у нас кстати на одном проекте в S3 образовалось 700TB данных, а используем только 40) - IO проблема - data shuffle Это еще повезло, что не надо исторически обновлять данные, так как это просто события. А прикиньте у вас данные по клиентам, и там например нужно взять сессию клиента, которая больше часа или 24 часов и использовать оконную функцию, тут уже так красиво не сделать.
Впервые о том, что такое конвейеры >100 ТБ, я узнал, когда присоединился к команде Core Growth в Facebook в 2016 году. Первые три месяца были наполнены мороженым, поездками на велосипеде и множеством развлечений!
Потом мой райский сон внезапно рассеялся - мой босс сказал мне, что теперь я отвечаю за push-рассылки, email и SMS-уведомления для всего Facebook.
Я был потрясен, когда узнал о том, что:
- Каждый день Facebook отправляет 50 миллиардов (!) уведомлений;
- У Facebook есть пять основных каналов, по которым отправляются уведомления (Jewel, Push, Email, SMS и Logged Out Push);
- Для оптимальной организации машинного обучения нужна низкая задержка (машинное обучение фильтрации уведомлений значительно ухудшалось, если задержка составляла более 1 дня);
- Это довольно сложная область, балансирующая между спамом и интересной информацией!
Если Вам больше нравится видео-формат, тогда держите ролик который я разместил в TikTok.
Мне было поручено создать гораздо более эффективную версию "notif_delivery_deduped", которая представляла бы собой дедуплицированный журнал всех уведомлений за день.
Старый конвейер представлял собой простой большой запрос GROUP BY, который выполнялся в конце дня. Только на выполнение этого запроса уходило 9 ЧАСОВ, процесс часто давал сбои и задерживал работу алгоритмов обучения ML-уведомлений, расположенных ниже в конвейере!
Дедупликация десятков миллиардов записей в микросерийном режиме
Изначально мой руководитель предложил мне создать потоковое задание во Flink, которое бы дедуплицировало данные в режиме реального времени из Scribe (Scribe - это версия Kafka от Facebook).
Я попытался это сделать, но загвоздка заключалась в том, что для выполнения данной задачи требовались десятки ТБ оперативной памяти. Инфраструктура Facebook в то время не могла себе этого позволить.
Потом мы решили дедуплицировать записи каждый час в пакетном режиме, а затем объединить эти часы вместе.
Сначала я попробовал следующий конвейер:
- Дедуплицировать час 1 → Дедуплицировать час 2 → объединить часы 1 и 2 → Дедуплицировать час 3 → объединить часы 1,2 и 3 и т. д., пока все 24 часа не будут дедуплицированы
Проблема с этим подходом заключалась в том, что он предполагал тонны операций IO и практически не оптимизировал время выполнения запроса (вместо 9 часов 6 часов)!

- Посовещавшись с инженерами, я предположил, что лучше организовать этот процесс в виде дерева, и тогда можно добиться колоссальной экономии времени!
Эта схема сработала лучше всего! Вот репозиторий кода с подробным описанием.
Итак, это решение действительно сработало! Оно позволило почти полностью дедуплицировать все данные в течение дня и минимизировала количество операций ввода-вывода! Теперь уведомления приходили через час после полуночи, а не через девять часов!
Однако, этот паттерн включал в себя множество FULL OUTER JOIN и GROUP BY, что влекло за собой слишком большое количество перетасовок (shuffling)… Как же минимизировать перетасовку (shuffling)?
Убедитесь, что таблицы merge, а также итоговая таблица отсортированы и бакетированы по уникальному идентификатору. Это означает, что единственное время, которое Вы фактически платите за перетасовку (shuffling), - это время работы с GROUP BY, а все FULL OUTER JOIN происходят вообще без перетасовки (shuffling)! Bucket Join - одна из самых мощных вещей, которые только можно сделать при работе с очень большими данными!
Ввод дедуплицированный событий в процесс обучения ML
После того, как мы оптимизировали поток уведомлений, я обратил свой взор на другие процессы. Между таблицей характеристик ML - уведомлений и таблицей дедупликации существовало неэффективное соединение. Кроме того, у последней таблицы был свой собственный шаг дедупликации.
Как раз в это же время Facebook начала внедрять Spark. Таким образом, я стал одним из первых инженеров по данным в Facebook, применивших Spark в своей работе!
Я сделал две вещи:
- Сделал генерацию таблицы характеристик сортированной и бакетированной по тем же ключам, что и таблица событий;
- Перевел GROUP BY на Spark, чтобы дедупликация и сортировка были в 10 раз эффективнее. (Разница в производительности Spark и Hive наиболее заметна в случае GROUP BY и JOIN с высокой кардинальностью)
Это позволило команде разработчиков генерировать и оценивать функции ML гораздо надежнее и быстрее! Кроме того, мы сократили использование вычислительных ресурсов для уведомлений на 30 %!
Что стоит, а чего не стоит делать в случае очень Больших данных
При работе с очень большими данными настоятельно рекомендую следующее:
- Для обработки данных используйте Apache Spark или Apache Flink;
- Используйте Scala Spark или SparkSQL, а не PySpark, потому что при таких огромных масштабах Вы не почувствуете особой разницы в производительности;
- Разделите данные и по дням и по часам (интервалы поменьше непродуктивными из-за медленного времени запуска Spark);
- Рассмотрите возможность разбиения данных на подразделы:
В Facebook мы сделали разбиение по дням, часам и "каналам", где канал - это метод, используемый для отправки уведомления;
Идеальными вариантами для разбиения являются разделы с низкой кардинальностью (< 15 значений);
- Если есть возможность, то сэмплируйте, чтобы не работать с Большими данными;
-
Отслеживайте косты по следующим направлениям:
- Стоимость вычислений в Spark;
- Стоимость операций ввода-вывода из S3 (или любого другого облачного провайдера, который Вы используете);
- Стоимость хранения данных в S3 (или у любого другого облачного провайдера, которого Вы используете).
Вообще, конечно, наибольшие затраты приходятся на операции входа и выхода данных из S3.
- Для минимизации количества join логируйте данные заранее, используя такие инструменты, как прокси Sidecar;
- Старайтесь, чтобы Ваши операции JOIN помещались в памяти и чаще используйте BROADCAST JOIN
Чего стоит избегать при работе с очень Большими данными:
- Длительное хранение = большие расходы;
- Попробуйте использовать Trino/Presto, Snowflake или любые другие похожие решения!
- Помните, что перетасовка (shuffle), вызванная JOIN, стоит гораздо дороже, чем shuffle, вызванная GROUP BY;
- Использование в SparkSQL ORDER BY, а не sortWithinPartitions в целях упорядочивания данных - несбыточная мечта в случае очень Больших данных!
Пожалуйста, не злоупотребляйте следующим (на первый взгляд может показаться, что операции, перечисленные ниже, могут помочь, но на самом деле это не так):
- Редко изменяйте какие-либо настройки Spark, кроме spark.executor.memory, spark.sql.shuffle.partitions, spark.driver.memory, spark.default.parallelism, spark.sql.adaptive.enabled, spark.sql.autoBroadcastJoinThreshold;
- Редко увеличивайте порог BROADCAST JOIN свыше 8 Гбайт, это не работает и вызывает проблемы со стабильностью;
- Редко меняйте типы сжатия, например lz4 vs snappy. На мой вщгляд, это совершенно бесполезное занятие.



