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