Команда VK Cloud перевела материал Netflix о двух подходах к автомасштабированию Apache Flink. Сначала компания разработала собственный автоскейлер, который анализировал внешние метрики кластера и эффективно экономил ресурсы на простых потоковых конвейерах. Но с ростом числа stateful-задач и сложных графов обработки этого стало недостаточно.

В статье — о том, чем автомасштабирование на уровне отдельных операторов отличается от масштабирования всего кластера, как Flink Autoscaler использует показатель True Processing Rate, зачем Netflix запускает отдельный Temporal workflow для каждой задачи и почему ради стабильности иногда выгоднее сознательно оставить запас вычислительных ресурсов.

Почему автомасштабирование не опционально при нашем масштабе

Netflix использует потоковую обработку на Apache Flink с 2017 года. По состоянию на 2026 год мы эксплуатируем более 30 000 задач Flink в нескольких регионах AWS. Большинство из них разворачивается не вручную; они генерируются нашей управляемой платформой Data Mesh, поэтому большинство пользователей никогда не работают с задачами Flink напрямую. Меньшая, но растущая часть — это кастомные задачи, которые создают и эксплуатируют команды по всей компании для таких сценариев использования, как персонализация, рекламные технологии (Ads) и прямые трансляции событий (Live events). Они варьируются от задач с одним оператором, которые перекладывают записи между топиками Kafka, до stateful-конвейеров с ветвлениями, соединениями и терабайтами состояния. Их нагрузка колеблется в зависимости от суточных циклов, релизов и региональных аварийных переключений.

Выделять ресурсы под пиковую нагрузку для каждой задачи расточительно, а под среднюю — во время всплесков это приводит к задержкам. А в нашей платформе действие масштабирования не бесплатно. По умолчанию это означает создание savepoint, плавную остановку задачи и ее перезапуск с новым размером. Для крупной stateful-задачи такой цикл может занять несколько минут. Это оставляет по-настоящему сложный вопрос: как дать каждой задаче ресурсы, которые ей нужны, именно тогда, когда они ей нужны? И как сделать это без участия человека, не сломав при этом ничего?

Первый автоскейлер: наблюдение снаружи

Наш первый ответ мы создали примерно в 2019 году: это был автоскейлер, устроенный как задача потоковой обработки. Он работал на Mantis, потребляя живой поток метрик уровня кластера из Atlas, нашей платформы телеметрии, включая сигналы CPU, network, Kafka lag, input-rate и consume-rate для каждой задачи. Скейлер объединял время наверстывания отставания (catch-up time), рассчитанное на основе lag, пороги использования CPU и сети, историю наблюдаемой производительности и регрессию по недавнему input rate, чтобы решить, когда масштабироваться вверх или сможет ли кластер меньшего размера справиться с окном lookahead. Поскольку автоскейлер работает независимо от платформы Flink, на него не влияют проблемы внутри самого Flink. Мы построили его как потоковую задачу, и это тоже облегчало масштабирование. Каждая нода автоскейлера обрабатывала метрики для подмножества задач Flink, и нам ни разу не пришлось писать собственную логику шардирования или координации, чтобы успевать за растущим парком Flink. Он стабильно сокращал использование ресурсов на 25–45% в тысячах управляемых конвейеров. Смотрите наш прошлый доклад на Flink Forward 2020.

Но у наблюдения снаружи есть потолок. Система рассуждала о кластере целиком через грубые метрики контейнеров и масштабировала один параметр, общее число TaskManager, поэтому каждый оператор в задаче двигался вместе с остальными. Это подходило для простых конвейеров с одним оператором, для которых он и создавался. Но не годилось для multi-оператора stateful DAG, которые команды всё чаще приносили нам для Ads, рекомендаций и игр. Именно с такими задачами он не мог справиться: для каждого нового случая приходилось писать больше специального кода, а не создавать общую возможность.

Автоскейлер настолько хорош, насколько хороши метрики, которые предоставляют лежащие под ним внешние системы. Эти метрики могли пропускать реальные проблемы. Задача могла быть полностью занята, хотя это никак не отражалось на использовании CPU. Из-за этого задача застревала в деградированном состоянии, которое скейлер никак не мог увидеть. Недавно миграция сети незаметно изменила то, как часть трафика передавала данные о себе, и часть метрик Atlas, на которые полагался скейлер, перестала фиксировать всё точно. Этот пробел оставался незаметным, пока не всплыл в production намного позже.

Настало время пересмотреть build и buy.

Второй автоскейлер: рассуждение изнутри

Когда мы начинали, у сообщества Flink не было зрелого автоскейлера, который оно могло бы предложить. К тому времени, когда мы пересматривали решение, он появился: Apache Flink Autoscaler. Вместо того чтобы наблюдать за контейнерами снаружи, он рассуждает изнутри задачи.

Его ключевая идея — оценивать TPR каждого оператора: пропускную способность, которую оператор мог бы выдерживать при полной загрузке. Flink сообщает для каждой подзадачи долю каждой секунды, потраченную на реальную работу, отдельно от времени, проведённого в состоянии backpressure или простоя. Деление наблюдаемой пропускной способности на этот доля занятости оператора (busyness) экстраполирует ёмкость до полной загрузки: оператор, обрабатывающий 700 records/sec, будучи занят 70% времени, имеет TPR 700 / 0,7 = 1 000 records/sec. Начиная с источников, автоскейлер обходит граф задач и использует TPR каждого оператора, его input/output ratios и target utilization, чтобы вычислить параллелизм, нужный каждой вершине. Так ни один оператор не становится узким местом, и кластер не приходится менять целиком, как единое целое.

Эти два подхода задают разный контракт, сведённый в таблицу ниже.

Самодельный автоскейлер

OSS-автоскейлер

Единица масштабирования

Кластер Flink целиком

Отдельные вершины графа задач

Границы масштабирования

Мин./макс. число TaskManager

Мин./макс. параллелизм по вершинам

Источник метрик

Поток метрик Atlas (push)

REST API JobManager Flink (pull)

Сигнал

CPU, сеть, прогноз пропускной способности

Busy time по вершинам / true processing rate

Где лучше всего

Stateless-конвейеры с одним оператором

Сложные multi-оператор DAG (stateless и stateful)

Конфигурация

На уровне платформы, ограниченная настройка под задачу

Переопределения и профили на уровне задачи

Решающее для нас отличие — в двух последних строках: OSS-автоскейлер может масштабировать именно те stateful multi-оператор задач, с которыми не справлялась наша самодельная система, а ещё каждая задача несёт собственную конфигурацию: stabilization period, пороги и другое поведение масштабирования, настроенное под нагрузку. Это сделало его естественным выбором для кастомных задач, которые команды до этого масштабировали вручную.


Lakehouse-платформа для аналитики и ML

Объединяйте данные из разных систем и снижайте расходы на хранение в 7–10 раз

Получить консультацию

Как заставить это работать в масштабе Netflix

Принять алгоритм было несложно: сообщество уже сделало самую сложную часть работы. Наша работа заключалась в том, чтобы надёжно запускать его на наших собственных задач, и именно здесь наша система больше всего отличается от стандартного open-source-развёртывания.

Во-первых, OSS-автоскейлер изначально проектировался так, чтобы находиться внутри Kubernetes Operator для Flink. Но наша платформа Flink работает на собственном control plane, а не в этом операторе(см. наш прошлый доклад на Current Conference 2024). Позже сообщество приняло прекрасное решение оставить основную логику в виде отдельной библиотеки. Они провели рефакторинг четырёх обобщённых интерфейсов, благодаря которым мы легко встроились напрямую в нашу внутреннюю экосистему: context, несущий метаданные задачи и информацию REST API, state store, event handler и realizer, который применяет решения о масштабировании.

Этот сервис — приложение на Spring Boot, оркестрация которого выполняется на Temporal, устойчивом движке workflow. Workflow-оркестратор опрашивает наш control plane Flink примерно раз в минуту на предмет задач с включённым автомасштабированием и запускает по одному долгоживущему workflow на каждую задачу. Каждый workflow конкретной задачи забирает её метрики по вершинам из её Flink JobManager, запускает алгоритм оценки OSS и, если тот принимает решение о масштабировании, передаёт его realizer, который применяет изменение через наш control plane Flink.

Дизайн «один workflow на задачу» стал прямым ответом на боль. Сначала мы запускали оценки в едином batch-цикле по всему набору задач, и это было хрупко. Одна медленная или ведущая себя неправильно задача мог остановить сбор метрик и масштабирование для всех задач, стоящих за ним в очереди. Мы дали каждой задаче собственный устойчивый workflow, и это изолировало blast radius. Теперь одна проблемная задача падает и повторяет попытку самостоятельно, а runtime масштабируется горизонтально по мере подключения новых задач.

Во-вторых, между «работает у сообщества» и «работает в масштабе Netflix» стояли три инженерных пробела:

  • Сбор метрик при высоком параллелизме. На крупных задачах получение метрик от JobManager становилось узким местом, и часть причины крылась в runtime самого Flink. Чтобы решить это, мы изменили JobManager так, чтобы он кэшировал временные имена метрик и очищал их один раз, вместо повторного сканирования при каждом запросе, и добавили серверную фильтрацию, чтобы автоскейлер запрашивал только нужные ему метрики. Благодаря этому автоскейлер стал работать на задачах с числом подзадач Flink до 3 000, тогда как раньше он начинал испытывать проблемы уже примерно после 1 000. Эти изменения находятся в нашем внутреннем форке релиза Flink, а некоторые переданы в апстрим, например FLINK-36172.

  • Сохранение forward connection. Две отдельные вершины, соединённые forward-соединением, обязаны работать с одинаковым параллелизмом, потому что записи передаются в памяти по фиксированному локальному каналу. Если масштабировать только одну из них, Flink не выдаёт ошибку: он молча превращает это ребро в сетевой shuffle. Наш форк обнаруживает forward-connected подграфы и масштабирует каждый из них как единое целое.

  • Учёт ограничений sink. У некоторых sink есть конечная пропускная способность записи. Поэтому мы добавили обнаружение backpressure асинхронных sink (тоже изменение в форке), чтобы автоскейлер не отмасштабировал задачу вверх, если sink не может принять больше.

Прежде чем что-либо применить, realizer выполняет набор проверок безопасности. Например, он отказывается уменьшать масштаб задач в регионе, который эвакуируется во время общекорпоративного регионального аварийного переключения. Он также проверяет, что диска хватит, чтобы новый кластер вместил checkpoint-состояние задачи, и добавляет небольшой резервный буфер для более крупных кластеров.

Путь к одному автоскейлеру

В прошлом году автоскейлер на основе OSS достиг статуса general availability для кастомных задач в Netflix. Первые результаты обнадёживают: например, наша команда клиентской телеметрии и логирования сократила годовые расходы на вычисления Flink на 58%, сэкономив примерно $1,1 миллиона в год. Эту эффективность обеспечивают три ключевых фактора. Во-первых, статическое provisioning всегда должно учитывать пиковые нагрузки, а автомасштабирование динамически подстраивается под суточные циклы: оно улавливает падение трафика ночью и по выходным по сравнению с пиками в будни. Во-вторых, автоскейлер непрерывно подстраивает ёмкость сам: раньше командам приходилось вручную оптимизировать ресурсы после улучшений производительности или спадов после праздников. Наконец, единые размеры контейнеров дают более качественный bin-packing и более гранулярные шаги масштабирования.

Слишком агрессивное уменьшение масштаба тоже оказывается ловушкой. Урежьте слишком сильно, и CPU насыщается, а lag подскакивает. Система при этом не может отреагировать мгновенно: её окно метрик и stabilization period должны выстраиваться заново после каждого перезапуска. Сейчас мы используем target utilization 0,45, ниже дефолтного для сообщества значения 0,7, сознательно жертвуя небольшой долей эффективности ради стабильности. Меньшее количество более спокойных рескейлов стоит предельных издержек для крупных stateful-задач.

Хотя наш скейлер даёт детализированные сигналы и принимает решения на уровне вершин для stateful DAG, быстрый рескейл всё ещё сильно зависит от производительности восстановления состояния Flink Core. Сегодня самая большая оставшаяся цена масштабирования stateful-задач — не логика скейлера, а сам процесс перезапуска и восстановления состояния. Flink 2 решает это через архитектуру disaggregated state, которая хранит состояние во внешнем хранилище, а не на локальном диске. Это может резко снизить зависимость рескейла и восстановления от общего размера состояния. Netflix начал поддерживать Flink 2.2, и мы планируем поэкспериментировать с этим новым State Backend. Хотим понять, поможет ли он устранить узкие места восстановления состояния при масштабировании крупных stateful-задач.

В перспективе мы стремимся перевести все внутренние сценарии использования скейлера на новый, основанный на OSS-автоскейлере, чтобы упростить нашу операционную поверхность.

Ключевые выводы

По пути мы вынесли три урока, которые применимы не только к Flink:

  • Выбор метрик важнее сложности алгоритма. Наша самая полезная отладка редко касалась математики масштабирования; дело было в том, какому сигналу доверять больше всего. Разберитесь в своих метриках, прежде чем настраивать алгоритм.

  • Задавайте разумные значения по умолчанию, но оставляйте пространство для настройки. Наши управляемые задачи достаточно похожи друг на друга: одно хорошее значение по умолчанию покрывает большинство из них без изменений. В этом и смысл платформы. Но если навязать единую конфигурацию каждой задаче, это накажет те, которые не вписываются. Поэтому мы сочетаем значения по умолчанию с переопределениями на уровне задач и намеренно скрываем настройки, которые требуют глубокой экспертизы. Большинству команд вообще не должно приходиться думать об автоскейлере.

  • Сначала внедряй, потом расширяй. Мы создавали решение своими силами, потому что в 2019 году ничего зрелого не подходило нашей платформе. Когда появился сильный проект сообщества, правильным шагом было не защищать наши инвестиции вечно и не выдирать систему в одночасье. Мы решили внедрить его для новых нагрузок, вносить исправления обратно в проект и спланировать осознанную миграцию.