Перевод статьи From Spark to Apache Doris: How Kwai Made A/B Testing Metrics 145x Faster at Scale. Автор — Siwei Zeng, архитектор движков обработки данных (Data Engine Architect) в Kwai. Оригинал опубликован в блоге VeloDB 12 августа 2026 года; перевод публикуется с разрешения VeloDB. Китайский первоисточник — публикация SelectDB от 6 июля 2026. Примечания переводчика вынесены в блоки «прим. перев.», в конце — короткий раздел от переводчика.

Kwai (Kuaishou), разработчик генеративной видеомодели Kling AI, перенесла общекорпоративный пайплайн расчёта метрик A/B-экспериментов со Spark на Apache Doris. Переработка системы ускорила расчёт метрик до 145 раз и снизила потребление ресурсов на 72% на кластере из 2000 узлов и 100 000 вычислительных единиц (CU).

Мы уже дважды писали о внедрении Apache Doris в Kwai, и оба раза речь шла о слое обслуживания запросов в реальном времени (real-time serving) для рекламных нагрузок: сначала о замене ClickHouse на единый lakehouse, обслуживающий почти миллиард запросов в день, затем о замене Elasticsearch.

На этот раз речь о другой команде и другой нагрузке: пакетные вычисления по расписанию. Департамент платформы данных Kwai отвечает за общекорпоративную платформу A/B-экспериментов. От неё зависит каждое бизнес-направление Kwai: в компании любое изменение стратегии должно показать положительный результат A/B-теста до выхода в production.

Раньше платформа A/B-экспериментов Kwai работала на Spark, и с ней было две проблемы. Первая — время выполнения: расчёт на одном пути исполнения занимал около 21 минуты, и продуктовая команда, которой результаты эксперимента были нужны во второй половине дня, обычно получала их только на следующее утро. Принятие каждого решения откладывалось на день. Вторая проблема — стоимость. По мере роста числа A/B-экспериментов вычислительные затраты на эксперименты с полным трафиком росли линейно.

С учётом объёма данных команда платформы данных Kwai подошла к переходу со Spark на Apache Doris как к перестройке системы вычислений. Команда переработала размещение данных в хранилище, агрегацию с дедупликацией (distinct), выполнение пользовательских функций (UDF), планирование задач и управление метаданными на frontend-узлах (FE).

Результаты переработки системы:

  • В 145 раз быстрее на одном пути исполнения при расчёте A/B-метрик: с 21 минуты до 8,7 секунды.

  • В 30 раз быстрее на самом длинном пути исполнения: с 65 минут до 2,12 минуты.

  • На 72% снизилось потребление ресурсов.

  • Крупнейший единый production-кластер Apache Doris: 2000 backend-узлов (BE), 100 000 CU и пул памяти в несколько сотен терабайт.

  • Окно восстановления FE сокращено с ~27 до ~10 минут, объём метаданных в памяти — с 4,16 до 1,49 ГБ.

Масштаб A/B-кластера Doris в Kwai: 2000 BE, 5 FE, 100 000 CU, 12 вычислительных групп; суточный профиль подачи задач
Масштаб A/B-кластера Doris в Kwai: 2000 BE, 5 FE, 100 000 CU, 12 вычислительных групп; суточный профиль подачи задач

Масштаб A/B-кластера Doris: 2000 BE и 5 FE, 100 000 CU, сотни TiB памяти, 12 вычислительных групп. В таблице — среднесуточные показатели по 14-дневному аудиту: сотни тысяч INSERT-задач в день, пик подачи 08:00–09:00, соотношение пика к минимуму 6,5:1.

Как считаются метрики A/B-экспериментов

Пайплайн A/B-метрик имеет стабильную структуру. На входе всегда две таблицы. Таблица накопленного распределения (cumulative assignment table) фиксирует, в какой эксперимент и в какую группу попал каждый пользователь. Широкая таблица метрик фиксирует поведение пользователей. Структура результата тоже фиксирована: результаты агрегируются по эксперименту, группе и бакету, а итоговый набор обычно содержит от нескольких сотен до десятков тысяч строк.

Полный расчёт состоит из четырёх шагов:

  1. Сканировать таблицу накопленного распределения, отфильтровав записи по дате и названию эксперимента.

  2. Сканировать широкую таблицу метрик, предварительно агрегируя данные по идентификатору пользователя (UID) и применяя определения метрик.

  3. Соединить две таблицы по UID и агрегировать на уровне бакета.

  4. Свернуть результаты уровня бакета и записать их в таблицу результатов эксперимента.

Основные затраты приходятся на шаги 2 и 3: соединение таблиц (join), агрегация с группировкой и передача данных по сети создают высокую нагрузку на CPU и сеть.

Определяющая характеристика этой нагрузки: ключ соединения — всегда UID, а измерения агрегации стабильны. Команда могла глубоко оптимизировать вычисления, потому что шаблон запроса был известен заранее: один шаблон вычислений, один путь исполнения.

Шаблон задачи: все A/B-задания используют один SQL-скелет из четырёх стадий — expHitScan, stdPreAggrOrScan, expJoinStdAggr, bucketResultAggr
Шаблон задачи: все A/B-задания используют один SQL-скелет из четырёх стадий — expHitScan, stdPreAggrOrScan, expJoinStdAggr, bucketResultAggr

Один SQL-шаблон для всех A/B-заданий: сканирование распределения → предварительная агрегация метрик по UID → join по UID с агрегацией по world × exp × group × bucket → свёртка до результатов эксперимента. Около 400 000 экземпляров в день.

Прим. перев. В статье слово «бакет» (bucket) встречается в двух разных смыслах, и их важно различать. В этом разделе речь о бакете эксперимента — единице разбиения пользователей внутри группы, по которой агрегируются результаты (столбец bkt на схеме). В разделе про Colocate Join ниже речь о хеш-бакете распределения данных в Doris (DISTRIBUTED BY HASH(uid) BUCKETS N) — физической единице распределения строк по узлам. Здесь важно не смешивать статистическое разбиение эксперимента с физическим распределением данных.

Почему Apache Doris подходит под эту нагрузку

Большая часть прироста производительности получена за счёт согласования модели вычислений этой нагрузки с механизмами исполнения Doris.

Кроме того, у Apache Doris есть три архитектурных преимущества перед Spark:

1. Pipeline-исполнение сокращает простои вычислительных потоков

Когда Spark перераспределяет данные между двумя стадиями (shuffle), данные обычно сбрасываются на диск, а последующие стадии ждут завершения всех предшествующих вычислений, поэтому потоки блокируются. Apache Doris использует pipeline-модель: поток, ожидающий данные, уступает место другой задаче, и CPU меньше простаивает.

2. Векторизованное исполнение повышает эффективность пакетной обработки

Spark SQL обрабатывает данные построчно через формат InternalRow. Apache Doris использует колоночные блоки (Block) на всём пути исполнения, обрабатывая по 4096 строк за пакет. Колоночное представление подходит для SIMD-обработки на CPU, а функции вызываются один раз на пакет вместо одного раза на строку, что снижает накладные расходы.

3. Рантайм на C++ убирает накладные расходы JVM

Spark работает на виртуальной машине Java (JVM), где stop-the-world-паузы сборщика мусора (GC) и прогрев JIT-компилятора требуют времени. Apache Doris выполняет скомпилированный машинный код C++ без JVM GC и выходит на высокую эффективность исполнения сразу после старта.

Прим. перев. Сравнение архитектур здесь упрощено. Spark SQL также поддерживает колоночный кэш и векторизованное чтение Parquet; конкретный путь исполнения зависит от операторов, источников и настроек. Версия и конфигурация Spark в кейсе не приведены. Утверждение об отсутствии JVM GC у Doris относится к вычислительному движку BE на C++: frontend Doris работает на JVM, чему посвящён раздел ниже.

Ключевые оптимизации: хранение, вычисления и планирование

Часть выигрыша дал сам движок Doris. Остальное — результат адаптации хранения, вычислений и планирования под шаблон A/B-метрик. A/B-кластер Doris включает 5 FE-узлов и 2000 BE-узлов; его ресурсы — 100 000 CU и пул памяти в несколько сотен терабайт — распределены по 12 логическим вычислительным группам (compute groups). Аудит за 14 дней насчитал миллионы задач при суточном объёме в сотни тысяч. Эти задачи сканируют 40 ТБ и сотни миллиардов строк в день.

При таком масштабе накладные расходы одного неэффективного оператора повторяются в каждой задаче кластера. Завершив миграцию движка, команда пошла дальше и системно оптимизировала четыре слоя: распределение данных в хранилище, вычислительные операторы, низкоуровневое исполнение и управление планированием — с учётом особенностей шаблона A/B-метрик. Пятая область, управление метаданными FE, стала отдельным направлением работы после того, как производительность запросов стабилизировалась.

Хранение: убираем межнодовый shuffle с помощью Colocate Join

Первая ключевая оптимизация, Colocate Join, напрямую опирается на фиксированный ключ соединения UID. Идея: распределять данные по хеш-бакетам на основе UID уже при записи, чтобы строки с одним UID всегда попадали на одну машину и в один бакет.

Тогда при выполнении запроса таблица распределения и таблица поведения могут соединяться локально — без межнодового shuffle и без массового перемещения данных. Каждый узел обрабатывает только данные своих локальных бакетов, завершает локальную агрегацию и передаёт на финальное объединение небольшой набор результатов на уровне бакета.

Данные из production показывают эффект. Один бакет сканирует десятки миллионов строк. Локальный join отбрасывает примерно 95% из них, оставляя несколько сотен тысяч совпавших строк, а агрегация сводит их к нескольким тысячам. Большая часть данных отфильтровывается на локальном узле, и по сети идёт только небольшой результат после join и агрегации.

При внедрении Colocate Join важны два требования:

  • Данные в таблицах должны распределяться по хеш-бакетам на основе ключа соединения — здесь это UID.

  • Обе соединяемые таблицы должны иметь строго одинаковое число бакетов.

После изменения структуры таблиц команда проверила план выполнения через EXPLAIN. Оптимизация активна только тогда, когда в плане отображается Colocated. Если в плане всё ещё присутствует Shuffle, обычно это означает, что у двух таблиц не совпадают число бакетов, ключи распределения или настройки Colocate.

Согласованное размещение при записи, локальные join при запросе: обе таблицы бакетированы по hash(uid), одинаковые uid попадают на один BE; по сети идут только агрегированные результаты бакетов
Согласованное размещение при записи, локальные join при запросе: обе таблицы бакетированы по hash(uid), одинаковые uid попадают на один BE; по сети идут только агрегированные результаты бакетов

Обе таблицы бакетированы по hash(uid) при записи, поэтому строки с одинаковым UID всегда попадают на один BE. Воронка одного бакета в production: сканирование — десятки миллионов строк, после локального join — сотни тысяч, на выходе — тысячи.

Прим. перев. Механизм описан в документации Doris: Colocation Join. Помимо двух требований из статьи, таблицы нужно поместить в одну colocation group через свойство colocate_with, а тип и число столбцов распределения должны совпадать. Документация также требует одинакового числа реплик партиций внутри группы. При восстановлении или переносе реплик Colocate Join может временно не использоваться, пока группа не вернётся в стабильное состояние.

Вычисления: Local Distinct для Grouping Sets

Colocate Join убрал межнодовый shuffle из join. Но в вычислительном слое операторы дедупликации всё ещё могли породить новый глобальный shuffle.

SQL для A/B-метрик активно использует GROUPING SETS, потому что один запрос должен выдать агрегаты сразу для нескольких комбинаций измерений. В нативной среде исполнения Doris стандартная двухфазная оптимизация для Distinct не поддерживала сочетание Grouping Sets + Distinct, поэтому такие запросы всё равно вызывали глобальный shuffle — с высокой нагрузкой на сеть и пиковым потреблением памяти.

Команда спроектировала и реализовала для этого переписывание запроса Local Distinct Grouping Sets. Colocate-бакетирование уже гарантирует, что строки с одним UID находятся на одном узле, поэтому дедупликацию можно выполнять локально на каждом вычислительном узле. Каждый узел сначала завершает локальный Distinct, затем локальные результаты проходят глобальную агрегацию — стоимость shuffle снижается, семантика сохраняется.

У оптимизации два режима:

  • Прозрачное переписывание: оптимизатор распознаёт паттерн Grouping Sets + Distinct и переписывает его в локальный путь вычислений без изменений в прикладном SQL.

  • Явный вызов: команды могут указать локальную дедупликацию напрямую через синтаксис Distinct Local — для случаев, которые оптимизатор не покрывает, но для которых команда подтвердила корректность такой дедупликации.

Выигрыш в production:

Метрика

Улучшение

Пиковое потребление памяти на CU

на 69% меньше

Объём данных при shuffle

на 92% меньше

Число строк при shuffle

на 76% меньше

Задержка запроса

на 21% меньше

Время CPU

на 13% меньше

Одна оговорка: эта оптимизация не всегда лучше исходного плана выполнения. Нативный COUNT DISTINCT умеет вычислять и передавать данные во время shuffle, то есть частично совмещает вычисления и передачу по сети. Переписанный Local-вариант должен завершить локальную дедупликацию до начала глобальной агрегации, что добавляет около 6 секунд ожидания на барьере синхронизации.

Поэтому переписывание лучше всего работает там, где shuffle — основное узкое место: например, Grouping Sets высокой кардинальности в сочетании с Distinct. Решение о его включении зависит от профиля запроса и реального плана выполнения.

Local Count Distinct: пиковая память на SQL-запрос 3,79 → 1,18 ГБ (−69%); shuffle bytes −92%, shuffle rows −76%, CPU −13%, задержка −21%
Local Count Distinct: пиковая память на SQL-запрос 3,79 → 1,18 ГБ (−69%); shuffle bytes −92%, shuffle rows −76%, CPU −13%, задержка −21%

Данные до и после полного внедрения (базовый период 26–29 марта и период стабильной работы 4–6 апреля): пиковое потребление памяти на один SQL-запрос 3,79 → 1,18 ГБ; shuffle 62,7 → 5,2 ГБ/день; CPU 1111 → 972 ядро-часа/день; средняя задержка запроса 90,3 → 71,3 с. Рекомендованные условия включения: cd_count ≥ 3 и returnRows ≤ 1 000 000.

Прим. перев. На момент перевода (сентябрь 2026) я не нашёл публичного PR, который можно уверенно сопоставить с описанным переписыванием Local Distinct для Grouping Sets и синтаксисом Distinct Local. Доступность именно этой реализации в стандартной поставке Doris остаётся неподтверждённой. В таблице оригинала память указана на CU, а на иллюстрации — на SQL-запрос; это расхождение источника, поэтому напрямую приравнивать показатели нельзя. Внизу иллюстрации также приведена тестовая выборка из 100 SQL-запросов по 95 таблицам: 35% запросов ускорились, а 65% замедлились из-за фиксированной стоимости перехода от pipeline к барьеру. Отсюда рекомендованные условия включения.

Вычисления: нативные C++ UDF

Colocate Join и переписывание Local Distinct устранили большую часть пересылок данных. Следующим узким местом, выявленным при профилировании, стали вычисления на CPU. UDF распределения (assignment UDF) реализует основную логику A/B-пайплайна: для каждой строки лога поведения пользователя она определяет, относится ли этот UID к тестовой группе, контрольной группе или группе стратегии.

Один вызов почти ничего не стоит — наносекунды. Но на десятках миллиардов строк лога UDF превратилась в типичный горячий участок кода, потребляющий ресурсы CPU.

Анализ профиля показал, что Java UDF потребляет примерно 80% ресурсов CPU, а главные узкие места — накладные расходы на вызовы JVM и создание объектов. Команда переписала её как нативную C++ UDF, устранив накладные расходы на вызовы через Java Native Interface (JNI), а затем оптимизировала наиболее затратные операции на уровне стандартной библиотеки шаблонов (STL): выделение памяти, доступ к хеш-таблицам и создание объектов.

Основные оптимизации:

  • P0, конкатенация строк. Исходная реализация создавала объект String для каждой строки, что занимало около 30% CPU. Команда перешла на переиспользуемый буфер фиксированного размера в ThreadLocal, сократив построчное выделение памяти и давление на GC.

  • P1, доступ к конфигурации эксперимента. Исходная реализация на каждой строке искала конфигурацию эксперимента в unordered_map, затрачивая время на вычисление хеша и переходы по указателям. Команда развернула конфигурацию в массив на этапе инициализации, получив доступ по индексу за O(1) во время выполнения.

  • P2, конструирование объекта пользователя. Исходная реализация создавала умный указатель и экземпляр объекта на каждую строку, добавляя затраты на выделение памяти в куче и подсчёт ссылок. Команда перешла на чтение исходных значений напрямую из колоночных данных Block, обойдясь без объектной обёртки.

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

В сумме три оптимизации улучшили общую производительность выполнения A/B-шаблона примерно в 3 раза.

Прим. перев. В P0 используются термины ThreadLocal и GC, хотя раздел посвящён C++ UDF. Из статьи неясно, относится этот пункт к Java-версии UDF или к другому этапу оптимизации.

Планирование: удерживаем SLA изоляцией, backpressure и очередями приоритетов

Когда Apache Doris начал выполнять сотни тысяч задач в день, оптимизаций на уровне SQL и вычислений уже не хватало, чтобы гарантировать общую стабильность. Команда добавила системное управление на уровне планировщика.

Первый слой — физическая изоляция. Команда разделила A/B-кластер на 12 независимых вычислительных групп, разместив нагрузки разных приоритетов в разных группах, каждая со своей квотой ресурсов. Всплеск трафика в низкоприоритетной задаче больше не может забрать ресурсы у высокоприоритетной работы, и взаимное влияние между бизнес-направлениями исчезает.

Второй слой — управление внутри вычислительной группы. Здесь сочетаются три механизма:

  • Лимиты параллелизма: ограничить число одновременно выполняющихся задач, а остальные поставить в очередь, чтобы всплеск нагрузки не перегружал кластер.

  • Обратное давление на источник задач (backpressure): динамически подстраивать скорость подачи задач под текущую нагрузку, согласуя скорость записи и подачи задач с пропускной способностью кластера.

  • Планирование с очередями приоритетов: разделить трафик на четыре очереди, от P1 до P4, и сначала выполнять задачи из высокоприоритетных очередей, чтобы очередь в низкоприоритетной группе не задерживала планирование высокоприоритетных задач.

Физическая изоляция разграничивает ресурсы между группами, а управление внутри группы обеспечивает её стабильную работу. Вместе они позволяют высокоприоритетным пайплайнам соблюдать SLA в часы пик.

Стабильность: управление метаданными в крупном кластере

Когда оптимизация производительности была завершена, проявилось менее очевидное узкое место: стабильность подсистемы метаданных.

Моделирование ёмкости

FE в Apache Doris управляет метаданными: информация о базах данных, таблицах, партициях, транзакциях и таблетах хранится в памяти JVM. Чтобы восстанавливаться после сбоя, система непрерывно записывает изменения в Edit Log и периодически сохраняет на диск снимок метаданных (Checkpoint Image). При перезапуске она загружает снимок и воспроизводит небольшой хвост Edit Log. В production выявились два риска:

  • Риск на уровне одной таблицы. Когда число партиций в одной таблице превысило 10 000, FE при формировании чекпоинта должен был обходить весь список партиций. Внутренняя структура использовала индекс типа Java int, который упёрся в предел и вызвал аварийное завершение процесса, нарушив production-пайплайн.

  • Риск на уровне кластера. С ростом числа таблиц и достижением десятков миллионов таблетов объекты метаданных FE накапливались в куче. Пиковое потребление памяти Master FE приблизилось к пределу в 400 ГБ, что часто приводило к Full GC.

Команда перешла от реактивных исправлений к активному управлению, введя моделирование ёмкости метаданных и мониторинг. В среднем один таблет занимает в памяти FE около 11 КБ. Эта оценка стала основой модели: по росту числа таблетов можно прогнозировать увеличение потребления памяти FE и заранее планировать расширение.

Оптимизация FE: меньше метаданных, ровнее GC, быстрее восстановление

Помимо управления ёмкостью, команда сжала сами структуры метаданных.

Исходный Table Inverted Index был спроектирован под локальное хранение и содержал избыточные поля — например, локальный путь к данным и локальный путь к метаданным. В облачном хранилище эти поля не участвуют в вычислениях, а значит, лишь увеличивают расход памяти.

Команда ввела Cloud Table Inverted Index, удалив поля, не используемые в облачных развёртываниях, и заменив часть коллекций компактными массивами, чтобы уменьшить объём памяти на каждый объект. В стресс-тесте на миллионе таблетов оптимизация снизила объём метаданных в памяти с 4,16 до 1,49 ГБ — на 64%. На десяти миллионах таблетов она экономит десятки гигабайт в пиковые моменты создания чекпоинта.

Прим. перев. Разделение TabletInvertedIndex на реализации для локального и облачного режимов внесено в upstream в PR #59683 9 января 2026 года и перенесено в ветку 4.1 в составе PR #61455 18 марта. Это подтверждает наличие соответствующей оптимизации в ветке 4.1, но не полное совпадение с внутренней реализацией Kwai. Цифры 4,16 → 1,49 ГБ — данные стресс-теста Kwai на миллионе таблетов из статьи. «Облачное развёртывание» здесь — режим Doris с разделением хранения и вычислений (compute-storage decoupled); в нём существуют вычислительные группы, о которых шла речь выше. Версия Doris в кейсе не указана.

Уплотнение структур уменьшило объём занимаемой памяти. Но управление этой памятью в JVM тоже требовало доработки. Метаданные FE, каталог и объекты чекпоинтов Edit Log живут долго и находятся преимущественно в старом поколении. На большой куче нужно заново настраивать выбор регионов для Mixed GC, долю молодого поколения и порог запуска конкурентной маркировки — initiating heap occupancy percent (IHOP); иначе старое поколение растёт до тех пор, пока не вызовет Full GC.

Команда настроила параметры G1 GC специально для FE с кучей объёмом примерно 400 ГБ:

Параметр

Назначение

-XX:ParallelGCThreads=48

Соответствует числу физических ядер

-XX:ConcGCThreads=12

ParallelGCThreads / 4

-XX:G1MixedGCLiveThresholdPercent=85

Поднят с 65 до 85: допустить сборку региона даже при доле мусора всего 15%, чтобы регионы старого поколения с высокой долей живых объектов тоже могли попасть в Mixed GC

-XX:InitiatingHeapOccupancyPercent=45

Запускать конкурентную маркировку при заполнении кучи на 45% — задаёт момент старта

-XX:-G1UseAdaptiveIHOP

Отключить адаптивный IHOP: на большой куче длинный хвост трафика может приводить к завышению порога

-XX:G1NewSizePercent=5

Нижняя граница молодого поколения, 5%

-XX:G1MaxNewSizePercent=25

Верхняя граница молодого поколения, 25% (в большой куче долю молодого поколения ограничивают)

-XX:G1ReservePercent=20

Резерв 20% против исчерпания to-space и для безопасных humongous-аллокаций

После настройки 24-часовая проверка в production показала снижение пикового потребления памяти Master FE с 370 до 270 ГБ без единого Full GC за период.

Наконец, команда распараллелила восстановление FE при старте. Исходный процесс загружал метаданные последовательно на фазе loadDb, что при масштабе в миллион таблиц занимало около 17 минут, а полное восстановление при старте — около 27 минут. Параллельная загрузка метаданных таблиц существенно сократила эту фазу и уменьшила общее окно восстановления FE примерно до 10 минут.

Результаты канареечного тестирования в production:

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

Полный старт

Ускорение

Фаза loadDb

Последовательно (до оптимизации)

27,2 мин

17,2 мин

8 потоков

15,4 мин

1,77x

5,8 мин

16 потоков (внедрено полностью)

10,4 мин

2,61x

4,7 мин

Сама фаза loadDb сократилась с 17,2 до 4,7 минуты — ускорение в 3,6 раза, что уменьшило общее время старта на 62%.

По мере роста развёртываний Apache Doris управление метаданными становится архитектурной задачей, которую нужно учитывать с самого начала проектирования системы. Модели ёмкости, планирование жизненного цикла и способность к восстановлению — всё это часть надёжной эксплуатации большой production-системы.

Прим. перев. Публичного PR, который можно уверенно сопоставить с описанной параллельной загрузкой на фазе loadDb, я не нашёл (поиск на сентябрь 2026). Доступность именно этой реализации в стандартной поставке Doris остаётся неподтверждённой. Не найдено и публичного issue, объясняющего упомянутое переполнение int при более чем 10 000 партиций: само число партиций не объясняет достижение предела 32-битного счётчика.

Итоги

Переход со Spark на Apache Doris ускорил расчёт A/B-метрик в Kwai до 145 раз и снизил потребление ресурсов на 72%. Часть выигрыша обеспечили архитектурные преимущества Doris: асинхронное pipeline-исполнение, векторизованный вычислительный движок и высокопроизводительный рантайм на C++.

Остальное — собственная глубокая оптимизация Kwai под сценарий A/B: Colocate Join, переписывание Local Distinct для Grouping Sets, нативные C++ UDF, изоляция ресурсов через физические вычислительные группы и логические очереди приоритетов, а также управление метаданными, обеспечивающее стабильность FE при таком масштабе кластера.

Эксплуатация в production при таком масштабе подтверждает, что Apache Doris подходит для этой нагрузки, и даёт сообществу референсную реализацию платформы A/B-тестирования на Doris.

Kwai также установила новый рекорд крупнейшего единого production-кластера Apache Doris: 2000 узлов и 100 000 CU. Это снимает многие открытые вопросы о том, до каких размеров можно масштабировать один кластер Doris, и оставляет большой запас для большинства развёртываний.

Обсудить расчёт A/B-метрик, другие сценарии, развёртывание и оптимизацию Doris можно в Slack сообщества Apache Doris. Управляемый сервис на базе Apache Doris — VeloDB.


От переводчика

Несколько замечаний, которые мне кажутся важными при чтении этого кейса.

Главный переносимый урок — условия применимости оптимизаций. Фиксированный ключ соединения UID и стабильный набор измерений агрегации позволили спроектировать размещение данных под запрос, сократить shuffle и затем оптимизировать использование CPU в UDF. Для ad-hoc OLAP с непредсказуемыми запросами применимость всей этой цепочки придётся проверять отдельно.

Как читать цифры. Ускорение в 145 раз получено на одном пути исполнения после переработки пайплайна; на самом длинном пути — примерно в 30 раз. Потребление ресурсов снизилось на 72%. Данных о версии и настройках Spark в статье нет, поэтому делать из неё общий вывод «Doris быстрее Spark в 145 раз» нельзя — это сравнение конкретного переработанного пайплайна с его прежней реализацией. Показатели памяти 4,16 → 1,49 ГБ относятся к отдельному стресс-тесту на миллионе таблетов. Утверждение о крупнейшем кластере Doris — со слов VeloDB/SelectDB.

Что из этого доступно обычному пользователю Doris. Colocate Join, вычислительные группы и Workload Group с лимитами параллелизма и очередями документированы в Doris. Конкретная схема P1–P4 и обратная связь с источником задач описаны как часть системы Kwai. Разделение TabletInvertedIndex на локальную и облачную реализации перенесено в ветку 4.1 (PR #61455). Доступность именно описанных в кейсе Local Distinct для Grouping Sets и параллельного loadDb не подтверждена; если знаете соответствующие PR — напишите в комментариях, дополню.

Раздел про G1 показывает настройку под конкретную нагрузку: долгоживущие метаданные FE и кучу около 400 ГБ. Переносить эти параметры на другой JVM-сервис стоит после анализа GC-логов и проверки на своей нагрузке. Повышение допустимой доли живых объектов в регионе расширяет набор кандидатов для Mixed GC, но может увеличить паузы; отключение Adaptive IHOP тоже требует обоснования. Подробнее — в руководстве Oracle по настройке G1.

Раскрытие аффилированности: я сооснователь Datanomix, компании-партнёра VeloDB. Перевод согласован с VeloDB. Дополнения переводчика — в блоках «прим. перев.», подписях к иллюстрациям и этом разделе.

Ссылки

Комментарии (0)