Движки обработки данных: Polars, DataFusion, Ray и Spark

Автоматический перевод Эта статья была автоматически переведена с оригинальной английской версии.

На протяжении многих лет Pandas использовали для табличной обработки данных в памяти, а Apache Spark — для распределённой обработки. Такое разделение хорошо работало, пока данные оставались структурированными.

Современные пайплайны также обрабатывают изображения, аудио и видео. В таких нагрузках декодирование на CPU может не успевать за инференсом на GPU, а сборка мусора в JVM и Global Interpreter Lock в Python ограничивают throughput. Новые движки используют Rust и Apache Arrow, чтобы снизить эти накладные расходы.

Чтобы сравнить варианты, я протестировал Pandas, Polars, DataFusion, Daft и нативный Rust на двух реальных датасетах. Spark и Ray рассматриваются в отдельных распределённых ноутбуках. Поездки на такси в Нью-Йорке представляют табличную обработку, а изображения Food-101 — мультимодальный пайплайн.

Код находится в репозитории engine-comparison-demo. Запустите его на собственном железе: производительность движка зависит от машины, датасета и характера нагрузки.


Типы нагрузок обработки данных

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

Два мира данных: структурированные и мультимодальныеДва мира данных: структурированные и мультимодальные

Структурированные, или табличные, данные предполагают фильтрацию, агрегацию и join-ы. Обычно такая обработка упирается в CPU и часто помещается в память. Многие подобные запросы можно выполнять на одной современной машине, однако необходимость в кластере по-прежнему определяют объём памяти, I/O, конкуренция, восстановление после сбоев и форма нагрузки.

Мультимодальный AI обрабатывает изображения, аудио или видео. Инференс может выполняться на GPU, тогда как декодирование на CPU, удалённое чтение или батчинг ограничивают пайплайн. Узкое место меняется в зависимости от типа инстанса и операторов, поэтому измеряйте утилизацию на всём пути, а не подбирайте GPU изолированно.

Новые движки ориентированы на разные участки этого пространства. Нативное выполнение может снизить накладные расходы Python, совместимые с Arrow интерфейсы — удешевить обмен данными, а streaming execution — ограничить потребление памяти на датасетах, превышающих объём RAM. Ни одно из этих преимуществ не возникает автоматически для каждого оператора или преобразования.


Часть 1: обработка на одной машине

Pandas использует eager-модель программирования с хранением данных в памяти. Многие распространённые операции создают промежуточные структуры, а параллельное выполнение не координируется query optimizer-ом. Такая модель делает исследовательский код прямолинейным, но на больших аналитических сканированиях может стать дорогой. В бенчмарке PDS-H для Polars при scale factor 10 Pandas выполнял весь набор примерно за 365 секунд, тогда как streaming engine Polars — за 3,89 секунды. Результат примерно 94x описывает именно этот бенчмарк и конфигурацию, а не является универсальным коэффициентом перехода с Pandas на Polars.

Eager- и lazy-выполнениеEager- и lazy-выполнение

Polars для локальных табличных пайплайнов

Polars — практичный вариант по умолчанию, когда локальный DataFrame-пайплайн перерастает eager-выполнение в одном процессе. Его lazy API строит query plan до запуска, благодаря чему доступны оптимизации вроде predicate pushdown и projection pruning. Затем движок может выполнять операторы параллельно и, если это поддерживается, обрабатывать данные streaming-батчами.

С ростом объёма данных разрыв увеличивается. При scale factor 100 (примерно 100 ГБ) streaming engine Polars завершил работу за 23,94 секунды против 152,27 секунды у собственного in-memory engine — примерно в 6 раз быстрее, причём данные превышали объём RAM.

Ниже приведён иллюстративный эскиз API, а не дословный фрагмент текущего engine_comparison_examples.ipynb или tabular.py:

import polars as pl

# scan_parquet reads only the schema — no data loaded yet
q = (
    pl.scan_parquet("yellow_tripdata_2024-01.parquet")
    .filter(
        (pl.col("trip_distance") > 5.0)
        & (pl.col("total_amount") > 30.0)
    )
    .group_by("payment_type")
    .agg(
        pl.col("total_amount").mean().alias("avg_fare"),
        pl.col("trip_distance").mean().alias("avg_distance"),
        pl.len().alias("trip_count"),
    )
)

# The entire plan is optimized and executed here, in parallel
result = q.collect()

Главное различие — модель выполнения. В приведённом lazy-запросе Polars может протолкнуть фильтр и выбор столбцов в Parquet scan. Это позволяет пропустить row groups и не читать неиспользуемые столбцы. В типичном eager-пайплайне Pandas нет общего query plan, который можно было бы оптимизировать, хотя аккуратное использование Parquet-фильтров, выбор столбцов и альтернативные бэкенды позволяют частично сократить разрыв.

DataFusion для встраиваемых query engine

Если Polars — это библиотека, с которой работают напрямую, то DataFusion — движок, поверх которого строят другие движки. Он используется в InfluxDB 3.0, GreptimeDB и ускорителе Comet Spark от Apple.

В опубликованном проектом DataFusion в ноябре 2024 года прогоне ClickBench DataFusion показал лучший результат среди протестированных одноузловых конфигураций с Parquet. Позже команда Embucket опубликовала vendor case study для TPC-H scale factor 1000, используя datafusion-cli на r6gd.metal с 64 vCPU, 512 ГБ RAM, локальным NVMe и измерениями на прогретом кэше. Q18 и Q21 в исходном прогоне завершились ошибкой или выполнялись чрезмерно долго, после чего для финальных результатов их структурно переписали. Это демонстрация scale-up на одной машине в указанной конфигурации, а не общий результат TPC-H.

Ниже приведён иллюстративный эскиз API, а не дословный фрагмент текущего engine_comparison_examples.ipynb или tabular.py:

from datafusion import SessionContext

ctx = SessionContext()
ctx.register_parquet("taxi", "yellow_tripdata_2024-01.parquet")

# SQL executed directly against Parquet — no intermediate copies
df = ctx.sql("""
    SELECT payment_type,
           COUNT(*)           AS trip_count,
           AVG(trip_distance) AS avg_distance,
           AVG(total_amount)  AS avg_fare
    FROM taxi
    WHERE trip_distance > 5.0 AND total_amount > 30.0
    GROUP BY payment_type
    ORDER BY trip_count DESC
""")
result = df.to_pandas()

Сильная сторона DataFusion — его модульность. Extension API охватывают пользовательские каталоги, table providers, правила оптимизатора и execution plans. Поэтому DataFusion часто выбирают при создании собственной data platform или встраивании query engine в собственный продукт.

Daft для мультимодальных данных

Daft делает изображения, аудио, видео и эмбеддинги частью DataFrame workflow. Он предоставляет нативные expressions для таких операций, как декодирование изображений и загрузка по URL, поэтому для распространённых preprocessing-путей не нужен построчный Python loop. Именно эта интеграция — главная причина выбрать его вместо универсального табличного движка.

Ниже приведён иллюстративный эскиз API, а не дословный фрагмент текущего engine_comparison_examples.ipynb или multimodal.py:

import daft

df = daft.read_parquet("yellow_tripdata_2024-01.parquet")

result = (
    df.where(
        (daft.col("trip_distance") > 5.0)
        & (daft.col("total_amount") > 30.0)
    )
    .groupby("payment_type")
    .agg(
        daft.col("total_amount").mean().alias("avg_fare"),
        daft.col("trip_distance").mean().alias("avg_distance"),
        daft.col("trip_distance").count().alias("trip_count"),
    )
    .collect()
)

Если вы знакомы с Pandas, API покажется привычным, но под капотом используется тот же стек Rust + Arrow, что и в Polars и DataFusion. Daft особенно полезен, когда «строки» в вашем DataFrame — это изображения, PDF-файлы или тензоры.

Бенчмарк на одной машине

Я протестировал Pandas, Polars, DataFusion и Daft, а также нативный Rust через Polars-rs, на примерно 41 млн поездок на жёлтых такси в Нью-Йорке за полный 2024 год. Spark и Ray представлены в распределённых примерах, но не в этом одноузловом бенчмарке. Полные результаты доступны в demo repo:

Результаты бенчмарка: производительность на одной машине (разные области измерения времени)Результаты бенчмарка: производительность на одной машине (разные области измерения времени)

График служит контекстом для сопутствующего бенчмарка, а не честным сравнением движков по принципу apples-to-apples: Pandas загружает данные поездок и зон до запуска измеряемой функции, тогда как Polars, DataFusion, Daft и Rust включают разный объём чтения файлов, регистрации или настройки запроса в измеряемую секцию. Скрипт повторяет каждую операцию три раза и сообщает медиану после загрузки импортов; состояние файлового кэша, прогрев, CPU и другие условия запуска могут изменить значения. Перед тем как считать порядок результатом сравнения производительности, выровняйте границы измерений и повторите тест в контролируемых условиях.

При таком объёме (примерно 660 МБ Parquet) три новых Python-движка завершают работу быстро. Приведённая выше цифра 94x для Polars относительно Pandas получена в PDS-H; на моей нагрузке с 41 млн строк такси разрыв оказался меньше. В этом прогоне с разными границами измерений Polars, DataFusion, Daft и Polars-rs показали результаты одного порядка, причём DataFusion зафиксировал минимальное измеренное время. Такой порядок относится к этому запросу и границе измерения, а не к движкам в целом. Когда результаты настолько близки, решение следует принимать по удобству API и повторным измерениям на репрезентативных данных.


Часть 2: мультимодальные данные

Табличный ETL обычно уменьшает объём данных: фильтрует, агрегирует и записывает меньше, чем прочитал. Мультимодальные AI-пайплайны делают обратное. Один путь обработки документа может развернуться в десятки текстовых чанков и векторов эмбеддингов.

В обычном пайплайне PySpark изображения и аудио часто поступают как бинарные данные, а затем передаются в Python-библиотеки для декодирования или преобразования. Переход через границу JVM/Python и сериализация результатов могут занимать значительную долю wall-clock time. Spark поддерживает Arrow-based paths, vectorized UDF и интеграции с акселераторами, но эти возможности требуют явного проектирования и измерений.

Пайплайновое выполнение и утилизация GPU

Spark планирует работу по стадиям, разделённым границами shuffle. Прямая реализация может выстроить загрузку, декодирование и инференс последовательно, из-за чего CPU или GPU будут простаивать. Spark поддерживает планирование ресурсов GPU и плагины, поэтому простой не неизбежен. Однако overlap и backpressure не являются автоматическими свойствами обычной задачи PySpark DataFrame.

Модель пайплайнового выполненияМодель пайплайнового выполнения

Пайплайновое выполнение — альтернатива такой схеме. Вместо последовательного запуска стадий движок перекрывает I/O, работу CPU и инференс на GPU. Если стадии и буферы сбалансированы, такое перекрытие сокращает простой и повышает throughput; при запуске и завершении, backpressure или наличии медленной стадии ресурсы всё равно могут простаивать.

Бенчмарк обработки изображений

Я протестировал обработку изображений на 500 реальных фотографиях из датасета Food-101. Результаты доступны в demo repo:

Результаты бенчмарка: мультимодальная производительностьРезультаты бенчмарка: мультимодальная производительность

Polars и DataFusion отсутствуют, поскольку этот тест проверяет нативные операции над изображениями, а не табличные expressions. В прогоне на 500 изображениях image path Daft оказался примерно в 3,7 раза быстрее baseline Pandas + Pillow. Точное значение multimodal_results.json 1.0148517909983639s, делённое на значение rust_multimodal_results.json 0.235865209s, даёт 4.30267692, поэтому реализация на Rust оказалась быстрее в 4,3 раза. В сопутствующем README отображаются округлённые значения 1.01s и 0.24s, их отношение равно 4,2x; это результат округления при отображении, тогда как значения в JSON являются эталонными. Пример полезен для выявления накладных расходов оркестрации, но сам по себе слишком мал, чтобы прогнозировать production throughput.


Часть 3: распределённая обработка

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

Apache Spark

Spark — зрелый вариант для крупного табличного ETL, join-ов с большим объёмом shuffle и организаций, уже эксплуатирующих его экосистему. Восстановление на основе lineage и широкая поддержка платформ важны, когда надёжность и операционная привычность перевешивают локальную скорость. Смешанные CPU/GPU-пайплайны требуют более тщательной настройки ресурсов и дизайна пайплайна, чем SQL-path, и здесь Ray Data и Daft предлагают более специализированную абстракцию.

Ниже приведён иллюстративный эскиз Spark API, а не дословный фрагмент текущего distributed_spark.ipynb или spark_etl.py:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder \
    .appName("TaxiETL") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

orders = spark.read.parquet("s3a://lake/taxi/*.parquet")
zones = spark.read.csv("s3a://lake/taxi_zones.csv", header=True)

# Spark excels at this: joining massive tables with a shuffle
result = (
    orders
    .filter(F.col("trip_distance") > 5.0)
    .join(zones, orders.PULocationID == zones.LocationID, "inner")
    .groupBy("Borough")
    .agg(
        F.sum("total_amount").alias("total_revenue"),
        F.avg("total_amount").alias("avg_fare"),
    )
    .orderBy(F.desc("total_revenue"))
)

Ray Data для неоднородных вычислений

Ray Data изначально проектировался для AI-нагрузок. В отличие от stage barriers в Spark, он использует streaming model, который поддерживает постоянную загрузку GPU. Самая важная для AI-пайплайнов возможность — планирование смешанных ресурсов: можно объявить, что одному actor нужны «1 GPU и 4 CPU», а другому — только CPU, после чего Ray распределит ресурсы.

Amazon сообщил о более чем $120 млн долларов ежегодной экономии после переноса выбранных внутренних задач обработки данных со Spark на Ray. В отчёте указана эффективность затрат на 91% выше в proof of concept и на 82% выше в production в расчёте на GiB входных данных из S3. Масштаб впечатляет, но это migration case study для конкретных нагрузок, а не ожидаемая экономия при типичном внедрении Ray.

Ниже приведён иллюстративный эскиз Ray Data API, а не дословный фрагмент текущего distributed_ray.ipynb или ray_inference.py:

import ray

ray.init()

ds = ray.data.read_images("s3://my-bucket/food101/")

class ImageClassifier:
    def __init__(self):
        import torch
        from torchvision.models import resnet18, ResNet18_Weights
        self.model = resnet18(weights=ResNet18_Weights.DEFAULT).cuda()
        self.model.eval()
        self.preprocess = ResNet18_Weights.DEFAULT.transforms()

    def __call__(self, batch):
        import torch
        tensors = torch.stack([
            self.preprocess(img) for img in batch["image"]
        ]).cuda()
        with torch.no_grad():
            preds = self.model(tensors)
        return {
            "prediction": preds.argmax(dim=1).cpu().numpy(),
            "confidence": preds.max(dim=1).values.cpu().numpy(),
        }

# ActorPoolStrategy creates persistent GPU workers
predictions = ds.map_batches(
    ImageClassifier,
    compute=ray.data.ActorPoolStrategy(size=4),
    num_gpus=1,
    batch_size=64,
)
predictions.write_parquet("s3://output/predictions/")

Здесь стоит обратить внимание на соотношение CPU и GPU. В бенчмарках Anyscale Ray Data продолжал масштабироваться при увеличении числа CPU на GPU, тогда как другие протестированные движки вышли на плато. Источник сообщает об ускорении инференса изображений до 3x по мере снижения CPU starvation: CPU-feeders успевают поддерживать загрузку GPU.

Распределённый Daft

Daft масштабируется через движок Flotilla: по одному воркеру Swordfish на узел, а Flotilla сверху выполняет планирование на уровне кластера. Swordfish отвечает за локальное выполнение на Rust и объединяет I/O с вычислениями через streaming небольшими батчами, поэтому каждый узел остаётся загруженным и не ждёт следующую стадию.

В бенчмарках, опубликованных Daft, Flotilla работал в 4–18 раз быстрее протестированных реализаций Spark на четырёх мультимодальных нагрузках. Наибольший разрыв наблюдался в задаче детектирования объектов на видео. Считайте этот результат свидетельством важности дизайна выполнения для таких пайплайнов, а затем сопоставьте его с приведёнными ниже результатами Anyscale и локальным тестом.

Anyscale, компания, стоящая за Ray, опубликовала конкурирующий бенчмарк, в котором после тюнинга Ray Data сократил или устранил отставание на инстансах с большим числом CPU. В совокупности эти два vendor study показывают, что нужно тестировать репрезентативные операторы и соотношения CPU/GPU, а не строить устойчивую league table.

Ниже приведён иллюстративный эскиз Daft API, а не дословный фрагмент текущего distributed_daft.ipynb или daft_pipeline.py:

import daft
import numpy as np
from daft import col

# Daft's Flotilla engine distributes work across the cluster
df = daft.read_parquet("s3://data-lake/pdf_metadata/*.parquet")

# Download PDFs — Daft parallelizes downloads in Rust
df = df.with_column("pdf_bytes", col("pdf_url").download())

# Define a GPU class UDF for batched embedding generation
@daft.cls(gpus=1)
class TextEmbedder:
    def __init__(self):
        from sentence_transformers import SentenceTransformer
        self.model = SentenceTransformer("all-MiniLM-L6-v2", device="cuda")

    @daft.method.batch(
        return_dtype=daft.DataType.fixed_size_list(daft.DataType.float32(), 384)
    )
    def __call__(self, text_col):
        texts = text_col.to_pylist()
        embeddings = self.model.encode(texts, batch_size=32)
        return np.asarray(embeddings, dtype=np.float32)

# Daft schedules CPU downloads and GPU embeddings simultaneously
embedder = TextEmbedder()
df = df.with_column("embedding", embedder(col("text")))
df.write_parquet("s3://output/embeddings/")

Сравнение распределённых решений

ВозможностьSparkRay DataDaft (Flotilla)
Модель выполненияЗадача на ядро, на основе партицийStreaming-задачи и actorsПо Swordfish на узел, streaming-батчи
Сильные стороныМасштабный SQL/ETL, отказоустойчивостьНеоднородные вычисления, насыщение GPUМультимодальный pipelining, ограниченная память
Управление GPUПланирование ресурсов и плагины экосистемыЯвные ресурсы задач и actorsИнтегрированное планирование CPU/GPU-пайплайна
Типичный тюнингExecutors, память, партиции, акселераторыБатчи, actors, object storeБатчи, ресурсы, конкурентность I/O
Подходящий сценарий оценкиКрупные SQL/ETL-задачи и jobs с большим shuffleОбучение или инференс со смешанными ресурсамиМультимодальный ingestion и transformation

Часть 4: роль Rust и Arrow

За конкуренцией движков скрывается более интересная история — их сближение. Polars, Daft и DataFusion активно используют Rust, а Ray сочетает нативные компоненты с Python API. Все они могут обмениваться данными через части экосистемы Apache Arrow.

Сближение Rust и ArrowСближение Rust и Arrow

Arrow PyCapsule Interface (протокол __arrow_c_stream__) предоставляет совместимым библиотекам стандартный способ обмена Arrow streams. При этом можно избежать построчной сериализации и переиспользовать буферы, если схемы и layout памяти совместимы. Materialization, rechunking, преобразование типов или передача между устройствами всё ещё могут привести к копированию данных, поэтому проверяйте передачу профилированием, а не предполагайте, что она бесплатна.

Ниже приведён иллюстративный эскиз API, а не самостоятельный скрипт: он предполагает, что events.parquet уже существует. Сопутствующий engine_comparison_examples.ipynb создаёт .data/events.parquet в setup cells до примера с Arrow.

from datafusion import SessionContext
import polars as pl

# Compute in DataFusion (Rust-native execution)
ctx = SessionContext()
ctx.register_parquet("events", "events.parquet")
df = ctx.sql("""
    SELECT user_id, COUNT(*) as event_count
    FROM events WHERE event_type = 'purchase'
    GROUP BY user_id HAVING COUNT(*) > 5
""")

# Convert through Arrow; compatible buffers may be reused
arrow_table = df.to_arrow_table()
df_polars = pl.from_arrow(arrow_table)

# Continue analysis in Polars
result = df_polars.with_columns(
    pl.col("event_count").rank().alias("rank")
).sort("rank")

Почему новые движки используют Rust для нативного выполнения?

  1. В коде на Rust нет tracing garbage collector. Ownership даёт нативным операторам более прямой контроль над выделением и освобождением памяти. Это может снизить latency, связанную со сборкой мусора, но не предотвращает давление на память и ошибки out-of-memory.

  2. Компактные нативные представления. Структуры Rust не содержат заголовки Java-объектов. Практическая польза зависит от layout движка: columnar-системы на JVM также избегают представления каждого значения как отдельного объекта.

  3. Более строгие проверки конкурентности. Правила safe Rust исключают многие data race на этапе компиляции. Код движка всё ещё может содержать unsafe blocks и логические ошибки конкурентного выполнения, но язык сужает поверхность отказов.

В сочетании с columnar layout памяти Arrow эти решения могут снизить накладные расходы на аллокации и сериализацию. Однако они не устраняют их на каждой границе: Python-объекты, сетевой транспорт, несовместимые схемы и перемещение данных между устройствами по-прежнему имеют значение.


Часть 5: использование нескольких движков

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

Модель совместного использования: мультидвижковый пайплайнМодель совместного использования: мультидвижковый пайплайн

Один из вариантов — использовать Spark для устоявшихся lakehouse join-ов, записывать результат в открытый табличный или файловый формат, а затем передавать его в Ray Data или Daft для инференса на GPU и мультимодальных преобразований. Polars или DuckDB могут выполнять локальный анализ тех же файлов. Небольшой команде может быть достаточно одного движка; добавляйте второй, только когда измеренный боттлнек оправдывает операционную границу.

Такую композицию делает возможной интероперабельность. В этих примерах общим файловым форматом передачи служит Parquet, а Arrow может обеспечивать обмен в памяти. Поддержка Delta и Iceberg зависит от движка и коннектора, поэтому проверяйте каждую границу. Передача через файл также имеет затраты на I/O и преобразования и не является бесплатной автоматически.

Выбор движка по сценарию оценки

Сценарий оценкиКороткий списокЧто проверить
Локальные аналитические DataFramePolars / DuckDBСоответствие API, память, набор запросов
Существующий распределённый SQL/ETLApache SparkПоведение shuffle, эксплуатация, стоимость
Параллельные задачи в PythonDask / RayНакладные расходы планировщика, модель отказов
Обучение или инференс со смешанными CPU/GPURay Data / DaftУтилизация акселераторов, backpressure
Мультимодальный ingestionDaft / Ray DataНативные операторы, поведение retry
Встраиваемый query engineDataFusionExtension API, покрытие SQL

Диагональное масштабирование

Долгое время выбор сводился к scale-up (более мощная машина) или scale-out (больше машин). Polars Cloud предлагает промежуточный вариант, который называет диагональным масштабированием. Сначала масштабируйтесь горизонтально при чтении из cloud storage, чтобы максимально загрузить I/O, а затем переключайтесь на один мощный узел, когда фильтры и агрегации уменьшат объём данных. Распределённого shuffle не происходит.

Более общий вывод: сравнивайте стоимость успешно выполненной задачи, а не почасовую цену инстанса. Более крупный узел может оказаться дешевле, если сокращение runtime перекрывает более высокий тариф, но результат зависит от I/O, давления на память и объёма работы, устранённого до shuffle.


Воспроизведение сравнения

В сопутствующем demo repository находятся скрипты, ноутбуки и окружение, использованные для сравнений. Быстрая команда ниже использует стандартный в репозитории датасет такси за полный 2024 год, соответствующий показанному выше большому прогону.

# Install dependencies
uv sync

# Run the tabular benchmark (~41M NYC taxi trips from full-year 2024)
uv run python -m engine_comparison.benchmarks.tabular

# Run the multimodal benchmark (500 real food photos)
uv run python -m engine_comparison.benchmarks.multimodal

# Run native Rust benchmarks (Polars-rs + image crate)
cd rust_benchmark && cargo run --release && cd ..

В репозитории также есть ноутбуки для side-by-side сравнения API и конфигурации Docker Compose для локальных распределённых запусков Spark, Ray и Daft.


Основные выводы

  1. Докажите, что одной машины недостаточно. Результаты scale-up показывают, что некоторые большие аналитические сканирования помещаются на одном узле. Перед тем как принимать накладные расходы кластера, измерьте требования к памяти, I/O, spill и восстановлению.

  2. Модели выполнения важнее привычного синтаксиса. Lazy-планирование, pushdown, streaming и параллельные операторы во многом объясняют разрыв между eager-скриптом для DataFrame и аналитическим движком.

  3. Измеряйте весь пайплайн работы с акселератором. Декодирование, сетевой I/O, батчинг, сериализация и backpressure определяют, будет ли GPU оставаться загруженным. Ray Data и Daft делают это перекрытие центральной абстракцией; Spark также поддерживает его при дополнительном дизайне и использовании инструментов.

  4. Сначала формируйте shortlist по нагрузке, затем запускайте бенчмарк. Используйте соответствие экосистемы, чтобы сузить выбор, а затем принимайте решение по репрезентативной end-to-end задаче. Vendor benchmarks — это гипотезы, а не гарантии.

  5. Открытые форматы сохраняют обратимость выбора. Совместимые с Arrow интерфейсы, Parquet и открытые табличные форматы могут снизить стоимость передачи. Уточняйте, переиспользует ли конкретное преобразование буферы или копирует их.

Если нужен отправной вариант, попробуйте Polars для локального табличного пайплайна, а Daft или Ray Data — для смешанного CPU/GPU-пайплайна. Выбирайте Spark, когда его распределённое выполнение, экосистема или существующая операционная база закрывают конкретное требование. Сам по себе размер датасета — слабое основание для выбора.


Ссылки

Бенчмарки и данные о производительности

  • Polars PDS-H Benchmarks — streaming Polars против Pandas при SF-10 (94x) и SF-100 (6,4x относительно in-memory)
  • DataFusion ClickBench Results (Nov 2024) — опубликованные проектом одноузловые результаты для Parquet
  • Embucket: TPC-H SF-1000 on DataFusion — vendor case study на r6gd.metal с измерениями на прогретом кэше; Q18 и Q21 были переписаны после того, как исходный прогон завершился ошибкой или выполнялся чрезмерно долго
  • Amazon Spark-to-Ray Migration — экономия $120 млн долларов в год, эффективность затрат в production на 82% выше
  • Daft Flotilla Benchmarks — в 4–18 раз быстрее Spark на мультимодальных нагрузках
  • Anyscale Ray vs. Daft Benchmarks — конкурирующий vendor benchmark для мультимодальных нагрузок

Движки и фреймворки

  • Polars — Rust-библиотека для DataFrame с lazy execution
  • Apache DataFusion — встраиваемый query engine на Rust (GitHub)
  • Daft — распределённый DataFrame с нативной поддержкой мультимодальных данных
  • Ray — фреймворк распределённых вычислений для AI (Data Internals)
  • Apache Spark — движок распределённых ETL и SQL
  • DuckDB — аналитический SQL-движок, работающий внутри процесса
  • Dask — параллельные вычисления в Python

Архитектура и экосистема

Demo repository

  • engine-comparison-demo — сопутствующий код статьи: бенчмарки, ноутбуки и Docker Compose для Spark/Ray/Daft

Датасеты

  • NYC Taxi Trip Records — данные NYC TLC о поездках на жёлтых такси (Parquet, примерно 2,9 млн строк в месяц)
  • Food-101 Dataset — 101 тыс. изображений еды от ETH Zurich (Bossard et al., ECCV 2014)