ADG Оставить заявку
Блог Данные и аналитика 5 мин чтения

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, что выглядит интересно для следующего этапа.

Контакт

Нужна такая же инженерная работа?

Опишите задачу и контекст. Ответим в течение рабочего дня, при необходимости подпишем NDA.