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

Apache Spark 2.0 в production: Structured Streaming вместо hourly batch

Первый production-пилот Spark 2.0: Structured Streaming заменил batch-обработку hourly-файлов, задержка упала с часа до минут. Что получилось на практике.

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

Apache Spark 2.0 вышел в июле 2016 года с Dataset API, унифицированным SQL и Structured Streaming - стриминговым движком поверх batch-ядра

В июле вышел Spark 2.0 с несколькими крупными изменениями сразу: Dataset API как единая абстракция вместо разрыва между RDD и DataFrame, унифицированный SparkSQL и, главное для нас, Structured Streaming - потоковая обработка, встроенная в Spark SQL. В декабре мы закрываем пилот, где всё это применили не на тестовых данных, а на реальной клиентской задаче.

Задача выглядела просто: клиент получает файлы от внешней системы каждый час, их надо загружать в хранилище, трансформировать и класть в аналитические витрины. Раньше это был cron + Python-скрипт + PostgreSQL. Задержка - от 60 до 75 минут в зависимости от объёма и загрузки сервера. Клиенту хотелось видеть данные быстрее: в идеале через 5-10 минут после появления файла.

Почему Spark, а не просто ускорить ETL

Первая мысль была очевидная: ускорить существующий скрипт. Но у задачи были особенности. Файлы приходят с переменным объёмом - иногда 50 MB, иногда 2 GB. Трансформации нетривиальные: джойны с несколькими справочниками, дедупликация по составному ключу, расчёт производных метрик. При 2 GB Python-скрипт занимал всё время до следующего файла, и margin практически не было.

Spark на клиентской задаче мы не использовали. PostgreSQL с параллельными запросами в 9.6 помог с аналитической нагрузкой, но для тяжёлых ETL-трансформаций это другая история. Решили пробовать 2.0: задача хорошо ложится на Structured Streaming - файл появился, его надо забрать, трансформировать и загрузить.

Что такое Structured Streaming в Spark 2.0

Принципиальное отличие от Spark Streaming 1.x в том, что Structured Streaming не про микробатчи как самоцель. С точки зрения API - это запрос к unbounded DataFrame: ты пишешь обычный SparkSQL или Dataset-преобразование, а движок сам разбивает поток на порции и обрабатывает инкрементально. Семантика exactly-once из коробки, state management встроен.

Для нашего случая это выглядело так: вместо «запускать скрипт раз в час» - постоянно запущенный стриминг-джоб, который наблюдает за директорией и обрабатывает новые файлы по мере появления.

spark = SparkSession.builder.appName("etl-stream").getOrCreate()

schema = StructType([...])  # явная схема, не inferSchema

stream = (
    spark.readStream
        .schema(schema)
        .option("maxFilesPerTrigger", 1)
        .csv("/data/incoming/")
)

transformed = stream \
    .join(ref_table, "product_id") \
    .dropDuplicates(["order_id", "event_ts"]) \
    .withColumn("revenue_rub", col("amount") * col("rate"))

query = (
    transformed.writeStream
        .outputMode("append")
        .format("jdbc")
        .option("url", jdbc_url)
        .option("dbtable", "analytics.events_stream")
        .option("checkpointLocation", "/data/checkpoints/etl")
        .trigger(processingTime="2 minutes")
        .start()
)

query.awaitTermination()

maxFilesPerTrigger=1 - сознательное ограничение: лучше обработать один файл за 2 минуты, чем ждать пока накопится партия. trigger(processingTime="2 minutes") - движок проверяет новые файлы каждые 2 минуты.

Что получилось

Задержка с момента появления файла до попадания данных в витрину - в среднем 4-7 минут. При пиковых объёмах бывает до 12, но не час. Это и было целью.

Несколько наблюдений по ходу пилота:

  • Dataset API приятнее RDD. Типизированные трансформации с compile-time проверкой (если Scala) или хотя бы явными схемами (Python) - это другой уровень отладки по сравнению со Spark 1.x.
  • Checkpoint - это хрупкость. Если схема данных меняется, checkpoint несовместим. Пришлось однажды чистить чекпоинт вручную после изменения структуры файла на стороне клиента. Это не баг, это логика инкрементального состояния - но нужно иметь процедуру.
  • Spark 2.0 и JDBC в writeStream - работает, но медленнее чем хотелось бы. Каждый микробатч открывает соединение с PostgreSQL. При маленьких файлах накладные расходы заметны. Решили промежуточно писать в Parquet, а отдельный джоб загружает из Parquet в PostgreSQL батчами.
  • Ресурсы. Для постоянно работающего стриминг-джоба нужен постоянно работающий кластер или хотя бы YARN. У клиента пока standalone Spark на одной машине - три executor'а по 4 GB. Хватает, но margin небольшой.

Что не так

Dataset API в Python - это DataFrame, типизация только в Scala/Java. В Python вы всё равно пишете через col() и доверяете схеме, которую сами задали. Это лучше чем inferSchema в продакшне, но не то же самое что Scala Case Class.

Мониторинг Structured Streaming в 2.0 минимальный. Есть query.status и query.lastProgress - JSON с метриками последнего триггера. Для интеграции с Grafana пришлось писать маленький экспортёр, который дёргает lastProgress и пишет в Prometheus. Это работает, но это рукоделие.

Документация на Structured Streaming в 2.0 хорошая в теории и неполная в практических деталях - особенно по части JDBC-синков и обработки ошибок при записи. Часть ответов нашли в JIRA и Stack Overflow.

Где сейчас

Пилот работает второй месяц в production. За это время - два инцидента: один раз упал джоб из-за corrupted-файла (добавили валидацию перед стримингом), один раз checkpoint не смог записаться из-за места на диске (алерт теперь есть). Оба инцидента решились быстро, данные не потерялись - exactly-once семантика отработала.

Клиент доволен задержкой. Spark 2.0 для этой задачи оказался рабочим выбором, не экспериментальным. Но это не волшебство: настройка ресурсов, обработка ошибочных файлов, мониторинг и процедура ротации чекпоинтов - всё это требует внимания отдельно.

Следующий шаг - смотреть на Kafka как источник вместо файловой директории, когда клиентская система будет готова его поддерживать. Пока - файлы.

Для проектов с ETL и аналитическими витринами Spark 2.0 интересен прежде всего Dataset API и Structured Streaming - это первая версия где стриминг выглядит как зрелый инструмент, а не как отдельный продукт с другой моделью программирования.

Контакт

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

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