ETL на Apache Spark: первый production-кластер и Dataset API вместо RDD
Запускаем Spark-кластер для обработки журналов событий ~50 ГБ/день. Время агрегации за сутки сократилось с 4 часов на SSIS до 18 минут. Dataset API вместо RDD - про типобезопасность.
Apache Spark 2.0 анонсирован с Dataset API и Structured Streaming для унифицированной обработки данных
На одном из проектов по DWH и аналитике мы несколько месяцев жили с честным, но медленным пакетным ETL: SSIS-пакет ночью собирал журналы событий из нескольких систем, трансформировал, грузил в хранилище. Объём подрос до ~50 ГБ в сутки, и агрегация за день стала занимать четыре часа. Аналитики смотрели данные утром и иногда ждали до обеда. Терпение закончилось.
Выбор был между тем, чтобы оптимизировать SSIS-пакеты, и тем, чтобы пересесть на что-то принципиально другое. Оптимизировать - значит бороться с инструментом, который проектировался под другие объёмы. Решили смотреть в сторону Spark.
Почему Spark, а не что-то попроще
На слуху были разные варианты. Kafka Streams мы уже щупали для потоковых агрегаций - там другой профиль задачи. Здесь же типичный batch: собрать за сутки, пересчитать, записать. Spark для этого и создавался.
Конкретный триггер - анонс Spark 2.0 со Structured Streaming и унифицированным API. Dataset API появился ещё в 1.6 - это типизированный DataFrame: у тебя есть схема, компилятор проверяет типы полей, ошибки ловишь не при выполнении задания на кластере, а при сборке. Для команды, которая привыкла к строгой типизации, это важно.
RDD API мы тоже смотрели - и поняли, что работа с ним похожа на написание MapReduce руками, только синтаксис приятнее. Dataset же даёт что-то вроде типизированных SQL-запросов с компилятором за спиной.
Как поднимали кластер
На пилот взяли четыре физических сервера - один master, три worker-а. HDFS поверх тех же машин: replication factor 2, namenode на master-е. Spark в standalone-режиме без YARN - это проще операционно, когда нет существующего Hadoop-кластера, с которым надо интегрироваться.
Схема получилась такая:
исходные логи -> HDFS (raw) -> Spark job -> HDFS (processed) -> DWH (PostgreSQL)
Логи падают в HDFS через небольшой коллектор на каждом источнике. Spark-задание читает partitioned-директории по дате, агрегирует, результат кладёт обратно в HDFS в Parquet. Отдельный загрузчик перекладывает итоговые таблицы в PostgreSQL для аналитиков.
Настройка HDFS заняла больше времени, чем ожидали. namenode - единая точка отказа в нашей конфигурации, HA-режим (два namenode с ZooKeeper) решили не поднимать на пилоте, но на продакшне это будет обязательно.
Dataset API на практике
Пример агрегации, с которой начали: считаем количество событий каждого типа по каждому источнику за сутки.
case class EventLog(ts: Long, source: String, eventType: String, userId: Long)
val events: Dataset[EventLog] = spark.read
.parquet(s"hdfs://namenode/raw/events/date=$date")
.as[EventLog]
val daily = events
.groupBy($"source", $"eventType")
.agg(count("*").as("cnt"))
.orderBy($"source", $"eventType")
daily.write
.mode(SaveMode.Overwrite)
.parquet(s"hdfs://namenode/processed/daily/$date")
Выглядит просто, но за этим стоит важная вещь: если поменять схему EventLog (например, переименовать поле), компилятор сразу покажет, где этот класс используется и что сломалось. С RDD такую ошибку ловишь в рантайме на кластере, после нескольких минут выполнения задания.
На практике это сэкономило несколько циклов отладки - именно в тот момент, когда схема источника неожиданно поменялась: один из поставщиков логов добавил поле и переименовал другое. С Dataset компилятор поймал несоответствие при сборке.
Что получили
Время агрегации суточного объёма сократилось с четырёх часов до восемнадцати минут. Это не магия распределённых вычислений в вакууме - здесь совпали несколько факторов: параллельная обработка на трёх worker-ах, чтение Parquet (колоночный формат, не читаешь лишнее), и то, что исходный SSIS-пакет работал последовательно и не был заточен под такой объём.
Несколько наблюдений по ходу:
- Parquet vs CSV. Исходные логи в CSV занимали примерно вдвое больше места и читались заметно медленнее. Перевели сырые данные тоже в Parquet при записи в HDFS - чтение сильно ускорилось. Первое очевидное улучшение.
- Размер executor. По умолчанию Spark выставлял executor memory в 1 ГБ. На наших заданиях это вызывало GC-давление и occasional OOM. Подняли до 4 ГБ на executor, число executor-ов уменьшили соответственно - стало стабильнее.
- Shuffle. groupBy на большом объёме - это shuffle по сети между worker-ами. При плохом партиционировании данных один worker перерабатывает, остальные ждут. Мы выбрали ключ партиционирования по source - распределение оказалось неравномерным (один источник даёт 60% объёма). Придётся переделывать.
- Spark History Server. Без него отлаживать задания тяжело. Поднять несложно, и он сразу даёт DAG задания, время каждой стадии, метрики shuffle. Рекомендуем не откладывать.
Где находимся
Пилот работает в параллельном режиме рядом со старым SSIS-пакетом. Расхождений в итоговых числах нет. На следующем этапе - переключить аналитиков на данные из Spark как основной источник и начать отключать SSIS-пайплайн по частям.
Hadoop-стек операционно тяжелее, чем PostgreSQL или даже Kafka: JVM на каждом узле, отдельный мониторинг heap и GC, HDFS со своей логикой балансировки блоков. Это не аргумент против - при нашем объёме альтернатив немного. Но это аргумент за то, чтобы считать TCO честно, не только по лицензии.
Spark 2.0 ещё не вышел в финале - работаем на 1.6.1, Dataset API доступен начиная с 1.6. Посмотрим, что в 2.0 изменится в Structured Streaming - там заявлен единый API для batch и stream, что выглядит интересно для следующего этапа.