Запуск экспериментов Apache Spark во сне (и в самолете)
Автор: Прашант Гедде Нараянасвами (Prashanth Gedde Narayanaswamy)
Будучи одним из крупнейших и наиболее критически важных конвейеров данных Netflix, member-sessionizer объединяет десятки миллиардов почасовых клиентских событий, включая действия пользователей, воспроизведения и сигналы вовлеченности, в отдельные сеансы просмотра для каждого профиля.
Быстрая двухнедельная оптимизация затянулась на два месяца, когда стандартные подходы зашли в тупик. Поскольку сбои происходили только при полном производственном масштабе, каждый тест занимал часы, ограничивая меня всего тремя-четырьмя попытками в день. Поворотным моментом стало изменение цели: вместо того чтобы пытаться ускорить один запуск, я создал способ запускать и мониторить несколько тестов одновременно и отслеживать каждый результат.
Этот подход привел к нестабильной настройке: одновременный запуск пяти заданий Spark и отслеживание их через нестабильный Wi-Fi во время долгого перелета домой, в то время как облачные серверы выполняли всю тяжелую работу. В этой статье рассказывается о том, как я настроил этот цикл, а также о «бутылочных горлышках» Spark и особенностях работы с памятью, с которыми я столкнулся по пути.
Дилемма масштаба: малые выборки не выявляют ошибки масштабирования
Стандартная инженерная мудрость гласит, что для быстрой итерации следует изолировать ошибки на небольших наборах данных. Это работает для стандартных программных ошибок, но совершенно не подходит для «узких мест» производительности больших данных.
При обновлении нашего конвейера сеансов с устаревшей архитектуры логирования на современную платформу новое задание столкнулось с критическими сбоями, включая ошибки нехватки памяти (OOM) у исполнителей (executor), серьезное переполнение области перемешивания (shuffle spillage) и тысячи повторных попыток этапов (stage retries). Ни одну из этих проблем не удалось воспроизвести с использованием выборочных данных. Сбои возникали только при полном производственном масштабе из-за перекоса данных (data skew): вкладка Netflix, оставленная включенной на ночь для воспроизведения потокового видео, генерирует десятки тысяч телеметрических событий, в то время как типичный сеанс производит лишь горстку. Стандартная 1%-ная выборка полностью пропускает эти экстремальные сеансы из «длинного хвоста».
Чтобы обнаружить реальные ошибки, каждый тестовый запуск должен был выполняться на полном часе невыборочных производственных данных со всеми сопутствующими накладными расходами на развертывание.
Стратегия: изоляция на основе веток (Branch-Driven Isolation)
Самым большим преимуществом стала функция в Dataflow — инструменте командной строки Netflix для сборки, тестирования и развертывания конвейеров данных. Он поддерживает разработку на основе веток (branch-driven development): наши рабочие процессы выполняются на Maestro (оркестраторе рабочих процессов Netflix), который планирует и выполняет задания, и они уже используют очищенное имя ветки (заполнитель ${dataflow.sanitizedbranch}). Поэтому, если я отправляю ветку вроде pgn/cap-experiment, то развернутый рабочий процесс, его таблицы и пространство имен выполнения получают метку с именем этой ветки. Это исключает конфликты с членами команды или между моими собственными параллельными запусками.
При такой изоляции модель была простой: одна ветка — это один независимый эксперимент. Я параметризовал планировщик (ограничение размера, разделы перемешивания, порог для тяжелых устройств, флаги функций) и присвоил каждому варианту собственную ветку и выходную таблицу. Вместо того чтобы ставить в очередь пять последовательных тестов на 10 часов, я одновременно запустил пять вариантов для одного и того же временного окна. Общее время выполнения сократилось с суммы всех запусков до длительности всего одного самого медленного задания.
Управление таблицами для каждой ветки не требовало никаких дополнительных затрат. Dataflow автоматически подготавливает и обновляет таблицы на основе объявленных схем без ручного DDL. Когда я изменял структуру пакета событий между экспериментами, схема таблицы обновлялась беспрепятственно, без необходимости выполнения CREATE TABLE или ALTER TABLE.
Единственным реальным ограничением было соперничество за кластер. Поскольку эти тяжелые задания Spark разделяли ресурсы с остальной частью компании, часы пик иногда ограничивали меня двумя одновременными запусками вместо пяти. Но даже запуск двух или трех заданий параллельно был колоссальным шагом вперед по сравнению с однопоточной отладкой.
Главный вывод: не сосредотачивайтесь на ускорении одного многочасового запуска, создайте среду, в которой вы сможете запускать множество экспериментов одновременно, не мешая друг другу.

Хотя на предыдущей диаграмме показана экономия времени, архитектура, обеспечивающая такую изоляцию, выглядит следующим образом:

Как помог агент ИИ
Запуск пяти параллельных экспериментов быстро становится хаотичным. Отслеживание идентификаторов экземпляров, конфигураций, утечек памяти и сбоев этапов для разных вариантов усложняет контроль.
Именно здесь на помощь пришел агент ИИ — не для того, чтобы писать хитрый код, а чтобы выполнять рутинную операцию. Я описывал варианты на обычном английском языке, а агент обновлял конфигурации планировщика, выполнял развертывание в ветку, запускал задания и опрашивал метрики. После завершения он сравнивал количество утечек и сбоев по вариантам в сводной таблице, позволяя мне задать вопрос: «Какой запуск ближе всего к работе без OOM при пиковой нагрузке?» — и получить немедленный ответ.
Ценность заключалась не в глубоких экспертных знаниях Spark, а в оперативной памяти. Он превратил утомительные перекрестные проверки пользовательского интерфейса Spark в один промпт.
В конце концов, я упаковал этот рабочий процесс в многократно используемый навык агента (.agents/skills/sandbox-deploy-and-run/), запускаемый через /sandbox-deploy-and-run [PARAM=value]. Он обрабатывает весь жизненный цикл — предварительную проверку, развертывание, фоновый опрос, автоматическую диагностику и поддержку реентерабельного возобновления — превращая многочасовые эксперименты в одну команду.
Почему tmux имел значение
Все это работало на Workbench, в облачной среде разработки Netflix, с использованием моего ноутбука в качестве тонкого клиента. Без сохранения сеансов обрыв SSH-соединения или переход ноутбука в спящий режим прервали бы многочасовые задания Spark прямо во время их выполнения.
Эту проблему решил tmux, запущенный непосредственно на Workbench. Каждый эксперимент жил в именованном облачном сеансе: я мог закрыть ноутбук, переподключиться спустя часы и продолжить с того места, где остановился. В сочетании с изоляцией веток это сделало возможным подход «запусти пять и отойди».
Теперь это моя стандартная настройка: локальный терминал (Ghostty), выделенные рабочие деревья git (git worktrees) для каждой ветки и постоянные сеансы tmux на Workbench. Локальный ноутбук одноразовый — его закрытие или потеря сети никогда не прерывают работу заданий.
Эта настройка превратила долгий день в пути в время активной итерации. Я запускал пакеты у выхода на посадку, проверял результаты с помощью агента ИИ через нестабильный Wi-Fi в самолете и втиснул 4–5 полных циклов экспериментов в один 24-часовой перелет.

Что выявили все эксперименты
Первопричина сводилась к одному оператору: collect_list. Поскольку Spark буферизует каждый полный пакет событий в памяти, сбои по нехватке памяти (OOM) в длинных сеансах были вызваны размером отдельного пакета, а не общим объемом данных. Устаревший конвейер обрабатывал идентичное количество строк без сбоев, доказывая, что каждое событие в новом пакете стало значительно тяжелее.
Вес событий увеличивали два фактора:
- Типы полей и боксинг Option: Идентификатор сеанса изменился с 8-байтового Long на строку UUID. Хуже того, обертывание полей в Scala Option заставляло Spark десериализовать типизированные датасеты в объекты JVM (обертки Some), что увеличивало распределение объектов и нагрузку на сборщик мусора (GC) для десятков тысяч событий на каждый сеанс.
- Раздутые строковые атрибуты: Малоценные поля (например, URL-адреса рефереров, распределения тестов) добавляли по нескольку сотен байт на каждое событие. Умножение этих строк на collect_list впустую тратило терабайты в час, хотя нисходящие задачи редко их использовали.

Исправления:
- Примитивные типы без боксинга: Удаление оберток Option из полей, не допускающих значения NULL (non-nullable), снизило пиковое использование памяти исполнителем со 112% до 58% от слота при нулевом сбросе на диск.
- Уменьшенные пакеты событий: Удалены неиспользуемые атрибуты для каждого события и подняты инвариантные для сеанса метаданные (версия приложения, страна) до заголовка сеанса.
Вывод: Снижение веса отдельного события обеспечило гораздо большую стабильность, чем выделение кластеру дополнительной памяти.
Два эксперимента, которые изменили мое мнение
Из десятков запусков два результата полностью перевернули мои представления об управлении памятью в Spark:
- Находка 1: Сужение строк может спровоцировать более быстрые OOM. Уменьшение размера нашей таблицы опустило ее ниже порога автошироковещательной рассылки (auto-broadcast) Spark, побудив планировщик запросов переключиться с соединений перемешивания (shuffle joins) на широковещательные соединения (broadcast joins). Поскольку широковещательные соединения закрепляют таблицы на стороне сборки непосредственно в памяти исполнителя, более узкие строки фактически увеличили общую нагрузку на память.
- Находка 2: Широковещательные хэш-соединения молча истощали память. Адаптивное выполнение запросов (AQE) переписало девять сортировочно-слитных соединений (sort-merge joins) в широковещательные хэш-соединения во время выполнения. Хотя каждое соединение по отдельности оставалось в пределах нашего лимита широковещательной рассылки в 400 МБ, AQE закрепило около 3,1 ГБ совокупных данных одновременно в одном исполнителе, вытолкнув контейнеры за пределы допустимого.
Хирургическое исправление: Вместо того чтобы отключать широковещательные соединения глобально, я снизил только порог выполнения AQE (spark.sql.adaptive.autoBroadcastJoinThreshold) до 10 МБ, оставив порог времени компиляции нетронутым. Это единственное изменение высвободило около 3,1 ГБ памяти исполнителя без влияния на исходное планирование запросов.
Вывод: Широковещательные соединения ускоряют запросы, но закрепленная память не бесплатна. Следите за накопленными широковещательными рассылками в вашем плане выполнения, а не только за размерами наборов данных.

Ведение единого журнала экспериментов
При проведении десятков параллельных тестов самый большой риск — это не сбой задания, а потеря контроля над тем, что вы уже попробовали, или путаница в том, какой JAR-файл и какая конфигурация привели к результату.
Чтобы решить эту проблему, я настроил агент ИИ на добавление строки в таблицу Markdown в репозитории после каждого запуска:
- Запуск № | Версия JAR | Изменение по сравнению с предыдущим запуском | Конфигурация | Ключевые метрики (Сброс на диск, Сбойные задачи, Пиковая память, Длительность)
Если запуск давал ключевой результат, агент добавлял короткую заметку, объясняющую, что этот результат доказал или опроверг.
Ключевые выводы:
- Единый источник истины: Это создало проверяемую историю для всего исследования.
- Принудительная дисциплина: Требование указывать «Изменение по сравнению с предыдущим запуском» заставляло меня изолировать переменные. Всякий раз, когда я объединял несколько изменений, я не мог выделить первопричину и был вынужден перезапускать тест — ошибка, которую формат журнала сразу же выявлял.
Вторичный документ (events-join-spark-tuning.md) фиксировал обобщенные выводы в виде краткого справочника «симптом — рычаг управления», чтобы соратникам по команде не приходилось пробираться через всю хронологическую историю.
Ошибка, на которой я научился: Исправление, которое работало только в часы наименьшей нагрузки.
Чтобы ограничить размер сеанса, я добавил row_number ().over (Window…) и отфильтровал первые N событий. Тесты в часы наименьшей нагрузки прошли чисто, но при пиковой нагрузке приложение аварийно завершалось с ошибками OOM.
Проблема носила структурный характер: оператор Window в Spark буферизует весь сеанс в памяти перед вычислением номеров строк и фильтрацией. Ограничение подрезало вывод на последующих этапах, но оно не ограничивало использование памяти во время буферизации, поэтому один перекошенный сеанс перегружал память исполнителя.
Исправление и урок: Мы перенесли усечение перед этап Window, чтобы ограничивать события до того, как произойдет буферизация.
Вывод: Тестируйте исправления нахудших сценариях перекоса входных данных. Запуски в часы наименьшей нагрузки маскируют структурные «бутылочные горлышки», которые пиковый масштаб немедленно проявит.
Окончательное исправление
Подводя итог, было выпущено четыре изменения, каждое из которых было изолировано в отдельном эксперименте перед внедрением:
- Поля, не допускающие значения NULL, переведены из Option в обычные примитивы. Пиковая память исполнителя упала примерно со 112% от слота до примерно 58%, без какого-либо сброса на диск.
- Уменьшен пакет событий. Удалены ненужные никому поля для каждого события и подняты инвариантные для сеанса поля (версии, страна, тип подключения) на уровень сеанса, чтобы они сохранялись один раз за сеанс, а не один раз за событие.
- Ограничен размер сеанса до этапа Window. Слишком большие сеансы отбрасывались до буферизации, а пакет усекался до лимита, а не отбрасывались целиком сеансы.
- Снижен порог выполнения широковещательной рассылки для AQE. Параметр spark.sql.adaptive.autoBroadcastJoinThreshold установлен на уровне 10 МБ, при этом порог времени компиляции остался нетронутым.
Нерешенная проблема «узкого места»
Хотя эти исправления заставили конвейер работать чисто, они не решили основную проблему масштабирования: агрегирование пакета событий всего сеанса в одну строку имеет жесткий предел. По мере развития нашего источника телеметрии и добавления новых атрибутов каждое событие в пакете становится тяжелее, а это означает, что эта нагрузка на память со временем вернется без изменения ни одной строки кода с нашей стороны.
Уменьшение размера и отказ от боксинга дали нам запас прочности, но они не изменили фундаментальную форму данных. Главный вывод из этой работы заключается в том, что оптимизация покупает время, а не бесконечный масштаб. Если ваша базовая модель данных допускает неограниченный рост строк, одна лишь оптимизация вас не спасет — в конечном итоге сама архитектура должна измениться.
Резюме
- Если ошибка проявляется только при масштабировании, прекратите попытки делать выборки. Вместо этого создайте способ запускать множество полномасштабных запусков параллельно.
- Если ваш планировщик может использовать имя ветки, то одна ветка становится одним изолированным экспериментом, и одновременный запуск множества веток становится простым.
- Используйте агента ИИ для выполнения операций, а не только для написания кода. Развертывание, опрос, извлечение метрик, сравнение вариантов: именно на это уходит большая часть времени.
- Как только цикл повторяется, упакуйте его в навык. Цикл «развертывание — запуск — наблюдение — диагностика» превращается в одну команду.
- Используйте tmux для долгих запусков. Запуск должен пережить закрытие вашего ноутбука.
- Ведите один единый журнал экспериментов и требуйте, чтобы каждая запись отражала ровно одно изменение по сравнению с предыдущим запуском. Сам формат заставляет вас изменять только одну переменную за раз.
- При проблемах с памятью сначала смотрите на ширину одной строки, а не только на общий объем данных. Более крупные типы полей, боксинг Option и избыточные неиспользуемые строки — все это суммируется внутри пакета collect_list.
- Оптимизация покупает запас прочности, а не новый потолок. Если форма данных может расти без ограничений, дизайн в конечном итоге придется изменить, одна лишь оптимизация вас не спасет.
