Как мы расширили Apache DataFusion для выполнения одного запроса на множестве машин

Источник: Datadog

Как мы расширили Apache DataFusion для выполнения одного запроса на множестве машин

Источник: Datadog

Узнайте, как Datadog создала Distributed DataFusion для масштабирования интерактивных запросов Apache DataFusion на несколько машин.

•Обновлено: 2 октября 2026 г.

По мере перехода Datadog к открытым стандартам мы выбрали Apache DataFusion в качестве единого компонуемого движка запросов, чтобы заменить ранее использовавшиеся специализированные системы. Однако при масштабах Datadog некоторые запросы должны сканировать огромные объемы данных, сохраняя при этом низкую задержку при возврате результатов. Хотя DataFusion предоставил нам необходимую расширяемость и компонуемость, он мог выполнять запрос только на одной машине. Чтобы удовлетворить наши требования к масштабируемости, мы создали Distributed DataFusion — фреймворк с открытым исходным кодом, который расширяет возможности DataFusion для выполнения одного запроса на нескольких машинах.

В этой статье мы расскажем, почему мы создали Distributed DataFusion, как он работает и какими были проектные решения при масштабировании интерактивных запросов на несколько машин.

Создание единого движка запросов

По мере того как Datadog расширялся, охватывая метрики, трассировки, логи, профили, мониторинг реальных пользователей и сигналы безопасности, инженерная культура «снизу вверх» привела к тому, что команды создавали специализированные конвейеры приема данных и движки запросов для различных рабочих нагрузок. Эти системы обеспечивали производительность, необходимую для каждого конкретного случая использования, но они также привели к фрагментации: разные движки, разные интерфейсы и ограниченная компонуемость между наборами данных.

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

Чтобы реализовать это видение, мы остановились на Apache Arrow в качестве формата данных в памяти, Substrait в качестве формата плана выполнения и Apache DataFusion в качестве нашего компонуемого движка запросов. Но DataFusion выполняет каждый запрос на одной машине, что ограничивает его способность поддерживать некоторые из наших крупнейших интерактивных рабочих нагрузок.

Distributed DataFusion расширяет модель выполнения DataFusion на несколько машин, позволяя одному запросу использовать дополнительные вычислительные ресурсы, сохраняя при этом задержку примерно на одном уровне по мере роста объема запрашиваемых данных. Проект имеет открытый исходный код, написан на Rust и доступен на GitHub и Crates.io под лицензией Apache 2.0.

Почему существующих распределенных движков было недостаточно

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

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

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

  • Интерактивное выполнение с низкой задержкой: На другой стороне экрана находится человек, ожидающий результата. Мы не искали систему, которая запускает фоновые отказоустойчивые задания. Нам нужно было решение с потоковой передачей без копирования (zero-copy).

Интерактивное выполнение с низкой задержкой: На другой стороне экрана находится человек, ожидающий результата. Мы не искали систему, которая запускает фоновые отказоустойчивые задания. Нам нужно было решение с потоковой передачей без копирования (zero-copy).

  • Интеграция с нашим стеком запросов: Arrow и Substrait являются контрактами между движком запросов и сервисами, которые его вызывают, поэтому любой распределенный фреймворк должен был хорошо интегрироваться с обоими.

Интеграция с нашим стеком запросов: Arrow и Substrait являются контрактами между движком запросов и сервисами, которые его вызывают, поэтому любой распределенный фреймворк должен был хорошо интегрироваться с обоими.

  • Эффективность при масштабировании: Накладные расходы на распределение должны быть как можно меньше. Тяжелые запросы должны эффективно использовать доступные ресурсы, в то время как легкие запросы не должны подвергаться штрафам.

Эффективность при масштабировании: Накладные расходы на распределение должны быть как можно меньше. Тяжелые запросы должны эффективно использовать доступные ресурсы, в то время как легкие запросы не должны подвергаться штрафам.

Эти требования исключили другие движки по разным причинам. ClickHouse не был расширяемым в той мере, в какой нам требовалось. Apache Spark и Apache DataFusion Ballista полагаются на материализацию промежуточных данных, что делает их плохо подходящими для наших интерактивных рабочих нагрузок с низкой задержкой. Trino расширяем, но он плохо интегрируется с Arrow и Substrait, а его эксплуатация в наших масштабах обходится дорого. В совокупности эти требования привели нас к решению расширить Apache DataFusion, а не внедрять существующий распределенный движок.

Результаты тестирования

Мы тестируем Distributed DataFusion, используя общедоступную инфраструктуру для бенчмаркинга, которую может воспроизвести любой пользователь с учетной записью AWS. В тестах используются стандартные наборы данных TPC-H и TPC-DS.

Наша тестовая среда состоит из кластера из 12 узлов Amazon EC2 c5n.2xlarge с 8 vCPU, 21 ГБ оперативной памяти и пропускной способностью сети до 25 Гбит/с на узел. Кластер считывает данные Parquet из Amazon S3. Каждый движок в сравнении работает на одном и том же оборудовании и наборах данных, с коэффициентами масштабирования TPC-H от 1 ГБ до 100 ГБ.

Более низкие значения указывают на более быстрое выполнение запроса. Результаты нормализованы относительно Distributed DataFusion (1.0×), поэтому столбцы выше 1.0× представляют более медленное выполнение.

Синтетические тесты показывают, как Distributed DataFusion соотносится с другими движками. В промышленной эксплуатации в Datadog мы увидели эффект в трех широких категориях:

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

Существующие тяжелые запросы: Некоторые длительные запросы, которые ранее с трудом выполнялись на одном узле DataFusion, теперь завершаются до 10 раз быстрее за счет выполнения на нескольких машинах при тех же общих затратах.

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

Существующие легкие запросы: Когда ожидаемая стоимость распределения перевешивает его преимущества, планировщик оставляет запрос на одной машине. Это позволяет избежать добавления накладных расходов на распределение для запросов, которые уже комфортно помещаются на одном узле.

  • Новые рабочие нагрузки на едином стеке запросов: По мере того как мы адаптируем все больше нашей инфраструктуры к Arrow, Substrait и DataFusion, мы переносим на этот стек все более тяжелые рабочие нагрузки. Ранее эти нагрузки опирались на специализированные движки запросов с индивидуально разработанным распределением, поскольку они не могли выполняться на одной машине.

Новые рабочие нагрузки на едином стеке запросов: По мере того как мы адаптируем все больше нашей инфраструктуры к Arrow, Substrait и DataFusion, мы переносим на этот стек все более тяжелые рабочие нагрузки. Ранее эти нагрузки опирались на специализированные движки запросов с индивидуально разработанным распределением, поскольку они не могли выполняться на одной машине.

Distributed DataFusion позволяет нам обслуживать эти рабочие нагрузки на общем стеке на базе DataFusion, и любая система, полагающаяся на него, получит преимущества от текущих и будущих улучшений производительности. Чтобы понять, как Distributed DataFusion достигает этого, давайте рассмотрим, как он расширяет существующую модель выполнения DataFusion.

Расширение DataFusion за пределы одной машины

DataFusion — это расширяемый фреймворк для создания пользовательских баз данных и аналитических систем. Чтобы понять, какое место занимает Distributed DataFusion, давайте кратко повторим, как запрос проходит через одноузловой движок DataFusion:

Distributed DataFusion работает исключительно с физическим планом, поэтому он оставляет предыдущие уровни планирования без изменений. Независимо от того, начинается ли запрос как SQL или Substrait, в конечном итоге он превращается в логический план, описывающий, что делает запрос, а затем в физический план, описывающий, как он выполняется. Distributed DataFusion расширяет эту модель физического выполнения с нескольких процессоров на одной машине до нескольких процессоров на нескольких машинах. Он делает это, опираясь на существующую модель секционирования DataFusion, которая обрабатывает непересекающиеся потоки данных параллельно.

Объединение результатов на нескольких машинах

Чтобы увидеть, как секционирование работает на практике, давайте начнем с простого запроса ORDER BY к таблице данных о погоде. Мы попросим DataFusion вернуть 10 самых высоких температур, когда-либо зарегистрированных:

SELECT "MaxTemp" FROM weather ORDER BY "MaxTemp" DESC LIMIT 10

Вы можете запустить его самостоятельно в DataFusion Fiddle.

На машине с двумя процессорами план физического выполнения выглядит следующим образом:

На одной машине с двумя процессорами DataFusion разбивает сканирование на два потока данных, которые выполняются параллельно перед объединением результатов. Distributed DataFusion расширяет ту же модель на рабочие узлы (workers).

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

SET distributed.max_tasks_per_stage = 2; -- ограничение максимального количества распределенных рабочих узлов до 2, для наглядности.

SET distributed.file_scan_config_bytes_per_partition = 1024; -- принудительное распределение этого запроса планировщиком.

SELECT "MaxTemp" FROM weather ORDER BY "MaxTemp" DESC LIMIT 10

Запустите этот пример в DataFusion Fiddle.

Этот план вводит три новых понятия:

  • Stage (Этап): Секция плана, разделенная сетевой границей. Другие движки запросов часто называют это фрагментом.

Stage (Этап): Секция плана, разделенная сетевой границей. Другие движки запросов часто называют это фрагментом.

  • Task (Задача): Задача для распределенного плана — это то же самое, что секция для одноузлового плана. Каждый этап делится на задачи, и каждая задача выполняется ровно на одном рабочем узле. Для простоты можно считать, что одна задача — это один рабочий узел.

Task (Задача): Задача для распределенного плана — это то же самое, что секция для одноузлового плана. Каждый этап делится на задачи, и каждая задача выполняется ровно на одном рабочем узле. Для простоты можно считать, что одна задача — это один рабочий узел.

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

DistributedLeafExec: Режим физического плана, представляющий различные варианты листовых узлов, выполняемых в разных распределенных задачах.

Несмотря на более сложный план, модель выполнения схожа:

  • Один узел: Два процессора обрабатывают два параллельных потока данных, которые объединяются в один поток перед возвратом результата.

Один узел: Два процессора обрабатывают два параллельных потока данных, которые объединяются в один поток перед возвратом результата.

  • Распределенная система: Два рабочих узла, каждый с двумя процессорами, обрабатывают четыре параллельных потока данных. NetworkCoalesce собирает потоки с удаленных рабочих узлов по сети перед их объединением в один поток.

Распределенная система: Два рабочих узла, каждый с двумя процессорами, обрабатывают четыре параллельных потока данных. NetworkCoalesce собирает потоки с удаленных рабочих узлов по сети перед их объединением в один поток.

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

Перемешивание (Shuffling) результатов на нескольких машинах

Рассмотрим среднее значение, сгруппированное по RainToday:

SELECT "RainToday", AVG("MinTemp") FROM weather GROUP BY "RainToday";

Запустите этот пример в DataFusion Fiddle.

На следующей диаграмме показано, как DataFusion вычисляет агрегацию на одной машине:

Изначально обе секции содержат смесь ключей группировки. После перераспределения каждая секция владеет непересекающимся набором ключей.

Планировщик разбивает агрегацию на три шага:

  • Частичная агрегация: Каждая секция независимо агрегирует строки, которые она содержит. Поскольку строки для заданного значения RainToday распределены по обеим секциям, ни одна из них не имеет полного результата. Для среднего значения каждая секция выдает текущую сумму и количество для каждой группы, а не готовое среднее значение.

Частичная агрегация: Каждая секция независимо агрегирует строки, которые она содержит. Поскольку строки для заданного значения RainToday распределены по обеим секциям, ни одна из них не имеет полного результата. Для среднего значения каждая секция выдает текущую сумму и количество для каждой группы, а не готовое среднее значение.

  • Перераспределение (Repartition): DataFusion хеширует каждый частичный результат по его ключу группировки и перемешивает данные так, чтобы все частичные результаты для одного и того же ключа попадали в одну секцию. После этого шага одна секция владеет всеми частичными результатами для RainToday = “Yes”, а другая — всеми частичными результатами для RainToday = “No”.

Перераспределение (Repartition): DataFusion хеширует каждый частичный результат по его ключу группировки и перемешивает данные так, чтобы все частичные результаты для одного и того же ключа попадали в одну секцию. После этого шага одна секция владеет всеми частичными результатами для RainToday = “Yes”, а другая — всеми частичными результатами для RainToday = “No”.

  • Финальная агрегация: Теперь каждая секция содержит все частичные результаты для ключей, которыми она владеет, поэтому она может вычислить финальные средние значения, объединив суммы и количества. Секции по-прежнему могут выполнять эту работу независимо и параллельно.

Финальная агрегация: Теперь каждая секция содержит все частичные результаты для ключей, которыми она владеет, поэтому она может вычислить финальные средние значения, объединив суммы и количества. Секции по-прежнему могут выполнять эту работу независимо и параллельно.

С двумя значениями RainToday (Yes и No) эта агрегация легко помещается на одной машине. Чтобы увидеть, что меняется, когда агрегацию нужно распределить между рабочими узлами, давайте вместо этого сгруппируем по Humidity9am, которое содержит процент относительной влажности, измеренный в 9 утра. Мы будем использовать два рабочих узла, каждый с двумя процессорами:

SET distributed.max_tasks_per_stage = 2; -- ограничение максимального количества распределенных задач до 2, для наглядности

SET distributed.file_scan_config_bytes_per_partition = 1024; -- принудительное распределение этого запроса планировщиком.

SELECT "Humidity9am", AVG("MinTemp") FROM weather GROUP BY "Humidity9am";

Запустите этот пример в DataFusion Fiddle.

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

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

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

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

На этом этапе данные правильно распределены внутри каждой машины, но еще не по всему кластеру. NetworkShuffle, узел, представленный в Distributed DataFusion, обменивается локально перераспределенными данными между машинами, чтобы они были глобально перераспределены:

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

Вычисления завершены, но результаты все еще разбросаны по нескольким машинам. NetworkCoalesceExec собирает эти потоки на одной машине, а CoalescePartitionsExec объединяет их в финальный результат, возвращаемый пользователю.

Этот пример показывает, как Distributed DataFusion распределяет агрегацию, но фреймворк также может распределять и другие операции, включая:

  • Разделенные соединения (Partitioned joins)

Разделенные соединения (Partitioned joins)

  • Широковещательные соединения (Broadcast joins)

Широковещательные соединения (Broadcast joins)

  • Распределенные объединения (Distributed unions)

Распределенные объединения (Distributed unions)

  • Агрегации типа «дерево-редукция» (Tree-reduce aggregations)

Агрегации типа «дерево-редукция» (Tree-reduce aggregations)

Расширяемость на практике

Все, что мы рассмотрели до сих пор — секционирование, сетевые границы, объединение и перемешивание — применимо к большинству распределенных движков запросов. Что отличает Distributed DataFusion, так это то, что он предоставляет эти возможности как расширяемый фреймворк, а не как готовый движок.

Вместо того чтобы предписывать, как данные должны считываться, распределяться или выполняться, Distributed DataFusion позволяет пользователям использовать собственные источники данных, узлы выполнения, сетевые решения и логику планирования. Документация проекта включает примеры адаптации кластера к сетевой инфраструктуре, распределения пользовательских узлов выполнения, распространения пользовательской конфигурации между рабочими узлами, создания пользовательских распределенных планов и routing data partitions to specific workers for cache-affinity scenarios.

Уроки, извлеченные при проектировании Distributed DataFusion

Два самых важных урока, которые мы извлекли при создании Distributed DataFusion, касались того, когда не стоит использовать распределение и как сохранить расширяемость DataFusion.

Знайте, когда не стоит распределять

Распределение запросов никогда не было самоцелью. Это был компромисс, необходимый для поддержки масштабов Datadog. Каждый распределенный запрос вносит дополнительные накладные расходы:

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

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

  • Сериализация данных: передача по сети требует сериализации и десериализации данных, что потребляет ресурсы процессора.

Сериализация данных: передача по сети требует сериализации и десериализации данных, что потребляет ресурсы процессора.

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

Накладные расходы на координацию: несколько машин должны координировать выполнение, обмениваясь информацией о состоянии и другими метаданными.

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

  • Распознавание случаев, когда распределение нецелесообразно: если ожидается, что запрос просканирует лишь небольшой объем данных или выполнит относительно мало вычислений, Distributed DataFusion выполняет его на одном рабочем узле вместо распределения.

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

  • Поддержание легковесности распределения: данные перемещаются по сети с использованием потоков Arrow Flight, уровень координации остается легковесным, а система постоянно тестируется для поддержания эффективности распределения.

Поддержание легковесности распределения: данные перемещаются по сети с использованием потоков Arrow Flight, уровень координации остается легковесным, а система постоянно тестируется для поддержания эффективности распределения.

Создавайте фреймворк, а не сервис

Это следует той же философии, которая привела нас к созданию Apache DataFusion. Distributed DataFusion — это не готовый сервис. Это библиотека, которую организации могут использовать для создания движков запросов, адаптированных к их собственным источникам данных, с минимальным количеством предписанных деталей реализации.

Мы проектировали его с учетом потребностей многих внутренних пользователей. Разные команды в Datadog создают разные производственные сервисы, подключаются к разным источникам данных и поддерживают разные шаблоны запросов. Такая же гибкость полезна и для организаций вне Datadog, которым необходимо создавать распределенные движки запросов для своих собственных рабочих нагрузок.

Создание Distributed DataFusion как фреймворка для распределенных движков запросов позволяет нам:

  • Способствовать вкладу сообщества, который улучшает проект для всех

Способствовать вкладу сообщества, который улучшает проект для всех

Следующие шаги для Distributed DataFusion

Distributed DataFusion продолжает развиваться. В настоящее время мы сосредоточены на двух областях:

Адаптивное выполнение запросов

Поскольку Distributed DataFusion — это фреймворк, пользователи могут подключать свои собственные источники данных. Не каждый источник данных предоставляет статистику, на которую опирается планировщик для оценки того, сколько данных просканирует запрос — или даже стоит ли вообще распределять запрос.

Мы делаем выполнение адаптивным. Вместо того чтобы полностью фиксировать стратегию распределения заранее, движок будет корректировать ее по мере наблюдения за фактическим объемом данных, проходящих через план, что позволит принимать более обоснованные решения даже при отсутствии или ненадежности статистики.

Ускорение на GPU

Теперь, когда Distributed DataFusion может масштабировать один запрос на несколько машин, мы также исследуем другой тип оборудования: GPU.

GPU предлагают гораздо более высокое соотношение производительности к цене, чем CPU, но имеют меньше памяти. Распределение помогает решить это ограничение. Распределяя запрос между несколькими GPU, мы можем объединить их память, используя при этом их совокупную вычислительную мощность, что позволяет выполнять рабочие нагрузки, которые не поместились бы на одном устройстве.

Создаем распределенные движки запросов вместе

Distributed DataFusion уже является частью производственной инфраструктуры в Datadog и других компаниях. Мы создали его как фреймворк, а не как фиксированный сервис, чтобы он мог поддерживать гораздо больше сценариев использования, помимо наших собственных.

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

Если вы создаете интерактивный движок запросов, мы будем рады, если вы попробуете Distributed DataFusion и внесете свой вклад. Независимо от того, подключаете ли вы новый источник данных, добавляете стратегию распределенного соединения, запускаете тесты производительности на своих наборах данных или улучшаете документацию, каждый вклад помогает сделать Distributed DataFusion более прочной основой для создания движков распределенных запросов.

Хотите работать над такими проектами, как Distributed DataFusion? Изучите инженерные вакансии в Datadog и помогите создавать системы, лежащие в основе наблюдаемости в масштабе.

О чём эта статья

Что-то непонятно? Спросите по статье — объясню простыми словами.

Не хотите разбираться сами? Мы поможем.