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

Apache Spark 1.5: первый пилот на телеметрии - 40 минут против 4

Запустили пилот Spark 1.5 для клиента с телеметрией: задача, которая на Hadoop MapReduce занимала 40 минут, на Spark выполняется за 4. Spark SQL дал аналитикам SQL вместо Java.

Контекст момента

Apache Spark 1.5 стал основным инструментом обработки больших данных, превосходя Hadoop MapReduce по скорости за счёт вычислений в памяти и оптимизированного планировщика

Несколько месяцев смотрели на Spark издали - читали, слушали на конференциях, но в реальную задачу не совали. В ноябре появился повод: клиент с телеметрией промышленного оборудования, кластер Hadoop, ETL-задача которая давно бесит всех причастных. Взяли и попробовали.

Задача

Производственный клиент собирает телеметрию с нескольких тысяч датчиков - несколько сотен гигабайт в день в HDFS. ETL-джоба агрегирует сырые события за сутки: нормирует единицы измерения, считает скользящие средние, детектирует аномалии по пороговым правилам, кладёт результат в DWH. Написана на Java, запускается на MapReduce.

Время выполнения - около 40 минут. Для ночного батча это приемлемо, но задача растёт вместе с парком оборудования, и аналитики просят возможность гонять те же расчёты ad hoc за произвольный период. На MapReduce адхок-запрос за неделю - это несколько часов, что практически исключает итеративный анализ.

Почему Spark, а не тюнинг Hadoop

MapReduce сохраняет промежуточные результаты на диск между стадиями map и reduce. Для нашей задачи это несколько последовательных стадий - нормировка, затем агрегация, затем аномалии. Каждый переход - запись в HDFS и чтение обратно.

Spark держит данные в памяти между трансформациями (RDD - Resilient Distributed Dataset). Если данные умещаются в оперативку кластера, промежуточных записей на диск нет вообще. Для нашей ночной задачи дневная телеметрия в памяти умещается с запасом - отсюда и разница в скорости.

Второй момент - планировщик задач. MapReduce запускает новый JVM-процесс на каждую задачу. Spark использует долгоживущие executor-процессы, и накладные расходы на запуск существенно ниже. Для задачи с сотнями шагов агрегации это заметно.

Как переписывали

Первую версию написали на Scala - это родной язык Spark, API там первичный. Для команды это был небольшой барьер: у нас не Scala-шоп, но люди с Java справились за пару недель. Сам порт MapReduce-логики на Spark RDD API оказался не таким болезненным - трансформации map, filter, reduceByKey читаются понятно.

Ключевой момент по архитектуре: в MapReduce каждая джоба - отдельный граф. В Spark вся обработка описывается как один граф трансформаций (DAG), и планировщик сам решает что можно параллелить, где нужен shuffle, где нет. Это сняло несколько склеек которые в MapReduce приходилось делать руками.

Время выполнения той же задачи на том же кластере - около 4 минут. Примерно десятикратное ускорение, что совпадает с тем что описывают в benchmarks Databricks. Мы не ожидали что цифра ляжет так точно, но вот.

Spark SQL: аналитики пишут SQL

Отдельная история - Spark SQL, компонент который есть с 1.0, но в 1.3 получил DataFrame API (вместо SchemaRDD) и в 1.5 заметно подрос. Он позволяет зарегистрировать DataFrame как временную таблицу и писать к ней SQL-запросы - стандартный SQL, ничего экзотического.

# PySpark - аналитик читает телеметрию напрямую из HDFS
from pyspark.sql import SQLContext
sqlCtx = SQLContext(sc)
df = sqlCtx.read.parquet("hdfs:///telemetry/2015-12-*/")
df.registerTempTable("telemetry")

result = sqlCtx.sql("""
    SELECT
        sensor_id,
        DATE(ts)        AS day,
        AVG(value)      AS avg_value,
        MAX(value)      AS max_value
    FROM telemetry
    WHERE ts >= '2015-12-01'
    GROUP BY sensor_id, DATE(ts)
    ORDER BY sensor_id, day
""")
result.show()

Для клиента это важнее чем скорость ETL. У них три аналитика, ни один не пишет на Java. Раньше любой нестандартный срез шёл через нас - мы писали MapReduce-джобу, запускали, отдавали CSV. Теперь аналитик сам пишет SQL в PySpark shell или Jupyter notebook, получает результат за минуты.

Ограничения есть: Spark SQL не умеет в оконные функции в полном объёме (частично добавили в 1.4, но не все диалекты). Аналитику, который привык к LAG() и LEAD() в SQL Server, придётся искать обходные пути или уходить в DataFrame API. Мы на это наткнулись при переносе нескольких аналитических запросов - часть пришлось переписать на RDD вручную.

Что не понравилось

Управление памятью. Spark агрессивно использует heap JVM, и при неаккуратном коде получить OutOfMemoryError несложно. Особенно если забыть вызвать .cache() там где нужно или наоборот закешировать слишком много. Настройка spark.executor.memory и spark.driver.memory - это отдельный навык, который нарабатывается на ошибках.

Мониторинг. Spark UI (порт 4040) показывает DAG, стадии, задачи - всё это есть. Но интегрировать это в наш Zabbix-мониторинг пришлось отдельными скриптами через REST API. Из коробки - только UI, метрики надо тянуть руками или через Ganglia/Graphite.

Версионирование зависимостей. Spark 1.5 тянет свои версии Scala, Hadoop, и конкретных библиотек. Если кластер уже под что-то заточен - возможны конфликты. У нас был один неприятный час с несовместимостью версии Parquet-библиотеки между Spark и тем что ожидал HDFS-клиент.

Где сейчас

ETL-задача клиента переехала на Spark, ночной батч работает. Аналитики получили доступ к PySpark shell и уже самостоятельно считают несколько срезов, которые раньше заказывали у нас.

До полноценного production-развёртывания с HA и нормальным operations-процессом ещё работа. Но то что Spark убедительно решает задачу - уже ясно. Работы по построению хранилищ и аналитических пайплайнов ведём в рамках DWH и бизнес-аналитики.

Контакт

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

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