Moteurs de traitement des données : Polars, DataFusion, Ray et Spark

Traduction automatique Cet article a été traduit automatiquement depuis la version originale en anglais.

Pendant des années, Pandas a pris en charge le traitement tabulaire en mémoire, tandis qu’Apache Spark s’occupait du traitement distribué. Cette séparation fonctionnait bien lorsque les données étaient structurées.

Les pipelines modernes traitent également des images, de l’audio et de la vidéo. Dans ces workloads, le décodage sur CPU peut affamer l’inférence sur GPU, tandis que le garbage collection de la JVM et le Global Interpreter Lock de Python limitent le débit. De nouveaux moteurs utilisent Rust et Apache Arrow pour réduire ces coûts.

Pour comparer les différentes options, j’ai benchmarké Pandas, Polars, DataFusion, Daft et Rust natif sur deux jeux de données réels. Spark et Ray sont couverts dans des notebooks distribués distincts. Les trajets de taxi à New York représentent le traitement tabulaire ; les images Food-101 représentent un pipeline multimodal.

Le dépôt engine-comparison-demo contient le code. Exécutez-le sur votre propre matériel, car les performances d’un moteur dépendent de la machine, du jeu de données et du workload.


Types de workloads de traitement des données

Les moteurs sont spécialisés : la première question à résoudre est donc le type de workload auquel vous êtes réellement confronté. La distinction principale oppose les données structurées aux données multimodales.

Deux mondes de la donnée : structurée vs. multimodaleDeux mondes de la donnée : structurée vs. multimodale

Les données structurées ou tabulaires impliquent des opérations de filtrage, d’agrégation et de jointure. Elles sont généralement limitées par le CPU et tiennent souvent en mémoire. De nombreuses requêtes de ce type peuvent s’exécuter sur une seule machine moderne, mais la mémoire, les E/S, la concurrence, la reprise sur erreur et la forme du workload déterminent toujours si un cluster est approprié.

L’IA multimodale traite des images, de l’audio ou de la vidéo. L’inférence peut s’exécuter sur un GPU tandis que le décodage sur CPU, les lectures distantes ou le batching limitent le pipeline. L’étape qui limite le débit varie selon le type d’instance et les opérateurs ; mesurez donc l’utilisation sur l’ensemble du chemin au lieu de dimensionner le GPU isolément.

Les moteurs plus récents ciblent différentes parties de cet espace. L’exécution native peut réduire l’overhead de Python, les interfaces compatibles avec Arrow peuvent rendre les échanges de données moins coûteux, et l’exécution en streaming peut borner la mémoire sur des jeux de données plus volumineux que la RAM. Aucun de ces avantages n’est automatique pour chaque opérateur ou conversion.


Partie 1 : traitement sur un seul nœud

Pandas privilégie un modèle de programmation eager, en mémoire. De nombreuses opérations courantes créent des résultats intermédiaires, et l’exécution parallèle n’est pas coordonnée par un optimiseur de requêtes. Ce compromis rend le code exploratoire direct, mais peut devenir coûteux lors de scans analytiques plus importants. Dans le benchmark PDS-H de Polars, avec un facteur d’échelle de 10, Pandas a mis environ 365 secondes pour exécuter la suite, contre 3,89 secondes pour le moteur streaming de Polars. Le résultat d’environ 94x décrit ce benchmark et cette configuration ; il ne constitue pas un facteur de conversion général de Pandas vers Polars.

Exécution eager vs. lazyExécution eager vs. lazy

Polars pour les pipelines tabulaires locaux

Polars est un choix pratique par défaut lorsqu’un pipeline DataFrame local dépasse le cadre d’une exécution eager, monop processus. Son API lazy construit un plan de requête avant l’exécution, ce qui permet des optimisations telles que le predicate pushdown et l’élagage des projections. Le moteur peut ensuite exécuter les opérateurs en parallèle et, lorsque cette fonctionnalité est prise en charge, par lots en streaming.

L’écart se creuse à mesure que les données augmentent. À l’échelle 100 (~100 Go), le moteur streaming de Polars a terminé en 23,94 secondes, contre 152,27 secondes pour son propre moteur in-memory — soit environ 6 fois plus rapide, sur des données dépassant la capacité de la RAM.

Ce qui suit est un exemple illustratif d’API, et non un extrait littéral de la version actuelle de engine_comparison_examples.ipynb ou 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()

La différence importante réside dans le modèle d’exécution. Dans la requête lazy ci-dessus, Polars peut pousser le filtre et la sélection de colonnes jusqu’au scan Parquet. Cela peut permettre d’ignorer des row groups et d’éviter de lire les colonnes inutilisées. Un pipeline Pandas eager classique ne dispose pas d’un plan global de requête à optimiser, même si l’utilisation rigoureuse des filtres Parquet, des colonnes sélectionnées et de backends alternatifs permet de réduire une partie de l’écart.

DataFusion pour les moteurs de requêtes embarqués

Alors que Polars est une bibliothèque utilisée directement, DataFusion est le moteur sur lequel on construit d’autres moteurs. Il alimente InfluxDB 3.0, GreptimeDB et l’accélérateur Spark Comet d’Apple.

Lors d’un run ClickBench publié par le projet DataFusion en novembre 2024, DataFusion est arrivé en tête des configurations Parquet testées sur un seul nœud. L’équipe Embucket a ensuite publié une étude de cas fournisseur pour un facteur d’échelle TPC-H de 1000, en utilisant datafusion-cli sur r6gd.metal avec 64 vCPU, 512 Go de RAM, du NVMe local et des mesures réalisées avec un cache chaud. Q18 et Q21 ont échoué ou se sont exécutées pendant une durée excessive lors du run initial, puis ont été réécrites structurellement pour obtenir les résultats complets. Cela montre une montée en puissance sur un seul nœud dans cette configuration, et non un résultat TPC-H général.

Ce qui suit est un exemple illustratif d’API, et non un extrait littéral de la version actuelle de engine_comparison_examples.ipynb ou 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()

La force de DataFusion réside dans sa modularité. Ses API d’extension couvrent les catalogues personnalisés, les fournisseurs de tables, les règles d’optimisation et les plans d’exécution. C’est pourquoi on le retrouve dès qu’il s’agit de construire une data platform personnalisée ou d’embarquer un moteur de requêtes dans son propre produit.

Daft pour les données multimodales

Daft intègre les images, l’audio, la vidéo et les embeddings au workflow DataFrame. Il fournit des expressions natives pour des opérations telles que le décodage d’images et les téléchargements depuis des URL, de sorte que les parcours de prétraitement courants ne nécessitent pas de boucle Python ligne par ligne. Cette intégration est la principale raison de le préférer à un moteur tabulaire généraliste.

Ce qui suit est un exemple illustratif d’API, et non un extrait littéral de la version actuelle de engine_comparison_examples.ipynb ou 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()
)

L’API vous semblera familière si vous connaissez Pandas, mais le moteur repose sur la même stack Rust + Arrow que Polars et DataFusion. Daft prend tout son intérêt lorsque les « lignes » de votre DataFrame sont des images, des PDF ou des tenseurs.

Benchmark sur un seul nœud

J’ai exécuté Pandas, Polars, DataFusion et Daft, ainsi que du Rust natif via Polars-rs, sur environ 41 millions de trajets de taxis jaunes de NYC couvrant toute l’année 2024. Spark et Ray apparaissent dans les exemples distribués, mais pas dans ce benchmark sur un seul nœud. Les résultats complets sont disponibles dans le dépôt de démonstration :

Résultats du benchmark : performances sur un seul nœud (périmètres de mesure mixtes)Résultats du benchmark : performances sur un seul nœud (périmètres de mesure mixtes)

Le graphique sert de contexte au benchmark associé, mais ne constitue pas un classement équitable des moteurs selon une comparaison strictement homogène : Pandas charge les données des trajets et des zones avant sa fonction chronométrée, tandis que Polars, DataFusion, Daft et Rust incluent dans la section mesurée des opérations différentes de lecture de fichiers, d’enregistrement ou de préparation de requête. Le script répète chaque opération trois fois et indique la médiane une fois les imports chargés ; l’état du cache du système de fichiers, la phase de warm-up, le processeur et d’autres conditions d’exécution peuvent modifier les valeurs. Alignez les limites de mesure et relancez le benchmark dans des conditions contrôlées avant d’interpréter l’ordre obtenu comme un résultat de performance.

À cette échelle (environ 660 Mo de Parquet), les trois moteurs Python les plus récents terminent rapidement. Le chiffre de 94× entre Polars et Pandas présenté ci-dessus provient de PDS-H ; ma charge de travail sur 41 millions de trajets de taxi a produit un écart plus faible. Dans cette exécution aux périmètres de mesure mixtes, Polars, DataFusion, Daft et Polars-rs se situaient dans la même fourchette générale, DataFusion enregistrant le temps rapporté le plus faible. Cet ordre est propre à cette requête et à cette limite de mesure, et ne caractérise pas les moteurs en général. Lorsque les moteurs obtiennent des résultats aussi proches, il faut s’appuyer sur l’adéquation de leur API et sur des mesures répétées avec des données représentatives pour prendre une décision.


Partie 2 : données multimodales

Les pipelines ETL tabulaires réduisent généralement le volume de données : on filtre, on agrège, puis on écrit moins de données qu’on en a lues. Les pipelines d’IA multimodale font l’inverse. Un simple chemin vers un document peut produire des dizaines de segments de texte et de vecteurs d’embedding.

Dans un pipeline PySpark conventionnel, les images et les fichiers audio arrivent souvent sous forme de données binaires, puis sont transmis à des bibliothèques Python pour être décodés ou transformés. Le franchissement de la frontière JVM/Python et la sérialisation des résultats peuvent représenter une part importante du temps écoulé. Spark prend en charge des chemins fondés sur Arrow, des UDF vectorisées et des intégrations avec des accélérateurs, mais ces choix nécessitent une conception et des mesures explicites.

Exécution pipelinée et utilisation du GPU

Spark planifie le travail en étapes séparées par des frontières de shuffle. Une implémentation directe peut enchaîner téléchargement, décodage et inférence de manière à laisser attendre une partie de la capacité CPU ou GPU. Spark prend en charge la planification des ressources GPU et les plugins associés ; les temps d’inactivité ne sont donc pas inévitables. Toutefois, le chevauchement des tâches et la gestion de la contre-pression ne sont pas des propriétés automatiques d’un job DataFrame PySpark standard.

Modèle d’exécution pipelinéeModèle d’exécution pipelinée

L’exécution pipelinée constitue l’alternative. Plutôt que d’exécuter les étapes l’une après l’autre, le moteur superpose les entrées-sorties, le travail CPU et l’inférence GPU. Lorsque les étapes et les buffers sont équilibrés, ce chevauchement peut réduire les temps d’inactivité et améliorer le débit ; le démarrage et la vidange du pipeline, la contre-pression ou une étape plus lente peuvent néanmoins continuer à laisser des ressources en attente.

Benchmark du traitement d’images

J’ai benchmarké le traitement d’images sur 500 photographies réelles issues du jeu de données Food-101. Résultats obtenus avec le dépôt de démonstration :

Résultats du benchmark : performances multimodalesRésultats du benchmark : performances multimodales

Polars et DataFusion sont absents, car ce test cible les opérations natives sur les images plutôt que les expressions tabulaires. Sur cette exécution portant sur 500 images, le chemin de traitement d’images de Daft était environ 3,7 fois plus rapide que la référence Pandas + Pillow. La valeur exacte de multimodal_results.json 1.0148517909983639s divisée par la valeur de rust_multimodal_results.json 0.235865209s est 4.30267692 ; l’implémentation Rust était donc 4,3 fois plus rapide. Le README associé affiche les valeurs arrondies 1.01s et 0.24s, dont le ratio est de 4,2 ; il s’agit d’un arrondi d’affichage, tandis que les valeurs JSON font autorité. Cet exemple est utile pour mettre en évidence le surcoût d’orchestration, mais il est trop limité pour prédire à lui seul le débit en production.


Partie 3 : traitement distribué

Lorsqu’une seule machine ne suffit plus, il faut choisir comment répartir le travail.

Apache Spark

Spark est une solution mature pour les ETL tabulaires de grande ampleur, les jointures coûteuses en shuffle et les organisations qui exploitent déjà son écosystème. Sa récupération fondée sur la lineage et son large support des plateformes sont importants lorsque la fiabilité et la familiarité opérationnelle priment sur la vitesse locale. Les pipelines mixtes CPU/GPU nécessitent une configuration des ressources et une conception du pipeline plus réfléchies que le chemin SQL ; c’est là que Ray Data et Daft proposent une abstraction plus ciblée.

Ce qui suit est un aperçu illustratif de l’API Spark, et non un extrait textuel de la version actuelle de distributed_spark.ipynb ou de 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 pour le calcul hétérogène

Ray Data a été conçu pour les workloads d’IA. Plutôt que les barrières entre les stages de Spark, il utilise un modèle streaming qui maintient les GPU alimentés. La fonctionnalité la plus importante pour les pipelines d’IA est la planification de ressources mixtes : vous pouvez déclarer qu’un actor veut « 1 GPU, 4 CPU » et qu’un autre ne veut que des CPU ; Ray se charge du reste.

Amazon a annoncé plus de $120 millions d’économies annuelles après la migration de certaines charges internes de traitement des données de Spark vers Ray. Le rapport indique une amélioration de 91 % de l’efficacité des coûts lors de la preuve de concept et de 82 % en production par GiB d’entrée S3. L’échelle est notable, mais il s’agit d’une étude de cas de migration portant sur des workloads spécifiques, et non d’une économie attendue pour une adoption classique de Ray.

Ce qui suit est un aperçu illustratif de l’API Ray Data, et non un extrait textuel de la version actuelle de distributed_ray.ipynb ou de 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/")

Le ratio CPU/GPU mérite ici une attention particulière. Dans les benchmarks d’Anyscale, Ray Data a continué à monter en charge à mesure que le nombre de CPU par GPU augmentait, tandis que les autres moteurs testés plafonnaient. La source fait état d’une accélération pouvant atteindre 3 fois pour l’inférence d’images lorsque la pénurie de CPU diminue, car les feeders CPU peuvent suivre le rythme du GPU.

Daft distribué

Daft se met à l’échelle via son moteur Flotilla : un worker Swordfish par nœud, avec Flotilla au-dessus pour assurer l’ordonnancement à l’échelle du cluster. Swordfish prend en charge l’exécution locale en Rust et pipeline les E/S avec le calcul au moyen d’un streaming par petits lots, afin que chaque nœud reste actif sans attendre l’étape suivante.

Dans les benchmarks publiés par Daft, Flotilla s’est montré 4 à 18 fois plus rapide que les implémentations Spark testées sur quatre workloads multimodaux. L’écart le plus important concernait un cas de détection d’objets dans des vidéos. Considérez ce résultat comme un indice de l’importance de la conception de l’exécution pour ces pipelines, puis comparez-le aux résultats concurrents d’Anyscale ci-dessous et à un test local.

Anyscale, l’entreprise à l’origine de Ray, a publié un benchmark concurrent dans lequel Ray Data a comblé ou inversé l’écart sur des instances à fort nombre de CPU après tuning. Ensemble, ces deux études de fournisseurs justifient de tester des opérateurs représentatifs et différents ratios CPU/GPU, plutôt que d’en déduire un classement durable.

Ce qui suit est un exemple illustratif d’API Daft, et non un extrait verbatim de la version actuelle de distributed_daft.ipynb ou de 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/")

Comparaison distribuée

FonctionnalitéSparkRay DataDaft (Flotilla)
Modèle d’exécutionUne tâche par cœur, fondé sur les partitionsTâches et actors en streamingUn Swordfish par nœud, lots en streaming
AtoutsSQL/ETL massif, tolérance aux pannesCalcul hétérogène, saturation des GPUPipelines multimodaux, mémoire bornée
Contrôles GPUOrdonnancement des ressources et plugins de l’écosystèmeRessources explicites pour les tâches et actorsOrdonnancement intégré des pipelines CPU/GPU
Optimisation couranteExecutors, mémoire, partitions, accélérateursLots, actors, object storeLots, ressources, concurrence des E/S
Bon cas d’évaluationTâches SQL/ETL volumineuses et riches en shuffleEntraînement ou inférence avec des ressources mixtesIngestion et transformation multimodales

Partie 4 : le rôle de Rust et d’Arrow

Derrière cette compétition, la convergence constitue l’aspect le plus intéressant. Polars, Daft et DataFusion utilisent largement Rust, tandis que Ray intègre des composants natifs en complément de ses API Python. Tous peuvent échanger des données via certaines parties de l’écosystème Apache Arrow.

Convergence de Rust et d’ArrowConvergence de Rust et d’Arrow

L’interface Arrow PyCapsule (protocole __arrow_c_stream__) fournit aux bibliothèques compatibles une méthode standard pour échanger des flux Arrow. Les transferts peuvent éviter la sérialisation ligne par ligne et réutiliser les buffers lorsque les schémas et les dispositions mémoire sont compatibles. La matérialisation, le rechunking, la conversion de types ou le transfert entre périphériques peuvent néanmoins entraîner des copies ; vérifiez donc le handoff avec du profiling au lieu de supposer qu’il est gratuit.

Ce qui suit est une esquisse illustrative d’API, et non un script autonome : elle suppose que events.parquet existe déjà. Le engine_comparison_examples.ipynb associé crée .data/events.parquet dans ses cellules de configuration avant l’exemple 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")

Pourquoi ces moteurs plus récents utilisent-ils Rust pour l’exécution native ?

  1. Pas de garbage collector à traçage dans le code Rust. Le système d’ownership donne aux opérateurs natifs un contrôle plus direct sur l’allocation et la libération. Cela peut réduire la latence liée au GC, mais n’empêche ni la pression mémoire ni les erreurs de mémoire insuffisante.

  2. Représentations natives compactes. Les structures Rust ne contiennent pas d’en-têtes d’objets Java. Le bénéfice concret dépend de la disposition utilisée par le moteur : les systèmes JVM en colonnes évitent eux aussi de représenter chaque valeur comme un objet individuel.

  3. Vérifications de concurrence plus strictes. Les règles de Rust garantissent l’absence de nombreuses data races à la compilation. Le code du moteur peut toujours contenir des blocs unsafe et des bugs logiques de concurrence, mais le langage réduit la surface de défaillance.

Combinés à la disposition mémoire en colonnes d’Arrow, ces choix peuvent réduire les surcoûts d’allocation et de sérialisation. Ils ne les éliminent pas à toutes les frontières : les objets Python, le transport réseau, les schémas incompatibles et les déplacements entre périphériques restent déterminants.


Partie 5 : adopter plusieurs moteurs

Vous n’avez pas besoin de faire passer toute la plateforme par un seul moteur. La composition est utile lorsque le coût du handoff entre les moteurs est inférieur à celui qui consisterait à faire gérer toutes les charges par un seul système.

Le modèle de coexistence : pipeline multi-moteursLe modèle de coexistence : pipeline multi-moteurs

Un relay possible utilise Spark pour les jointures lakehouse déjà en place, écrit dans une table ou un format de fichier ouvert, puis transmet le résultat à Ray Data ou Daft pour l’inférence GPU ou les transformations multimodales. Polars ou DuckDB peuvent servir à l’analyse locale des mêmes fichiers. Une équipe réduite peut n’avoir besoin que d’un seul de ces moteurs ; ajoutez-en un autre lorsqu’un goulot d’étranglement mesuré justifie cette frontière opérationnelle.

Cette composition est rendue possible par l’interopérabilité. Parquet est le format de handoff fichier commun dans ces exemples, tandis qu’Arrow peut fournir un échange en mémoire. La prise en charge de Delta et d’Iceberg dépend du moteur et du connecteur ; vérifiez donc chaque frontière. Un handoff par fichier entraîne également des coûts d’I/O et de conversion : il n’est pas automatiquement gratuit.

Sélection du moteur selon le cas d’évaluation

Cas d’évaluationSélection initialeÀ valider
DataFrames analytiques locauxPolars / DuckDBAdéquation de l’API, mémoire, mix de requêtes
SQL/ETL distribué existantApache SparkComportement des shuffles, opérations, coûts
Traitements parallèles natifs PythonDask / RaySurcoût du scheduler, modèle de défaillance
Entraînement ou inférence mixte CPU/GPURay Data / DaftUtilisation des accélérateurs, backpressure
Ingestion multimodaleDaft / Ray DataOpérateurs natifs, comportement des retries
Moteur de requêtes embarquéDataFusionAPI d’extension, couverture SQL

Mise à l’échelle diagonale

Pendant longtemps, le choix consistait à faire du scale-up (une machine plus puissante) ou du scale-out (davantage de machines). Polars Cloud propose une approche intermédiaire, appelée mise à l’échelle diagonale. On scale horizontalement pendant la lecture depuis le stockage cloud afin de saturer les I/O, puis on revient à un seul nœud puissant une fois que les filtres et les agrégations ont réduit le volume de données. Le shuffle distribué n’a jamais lieu.

La leçon générale est de comparer le coût par job réussi, et non le prix horaire d’une instance. Un nœud plus puissant peut coûter moins cher lorsque la réduction du temps d’exécution compense son tarif supérieur, mais le résultat dépend des I/O, de la pression mémoire et de la quantité de travail éliminée avant un shuffle.


Reproduire la comparaison

Le dépôt de démonstration associé contient les scripts, les notebooks et l’environnement utilisés pour les comparaisons. La commande rapide ci-dessous utilise le dataset de taxis couvrant toute l’année 2024, fourni par défaut dans le dépôt, et correspondant au run de plus grande taille présenté précédemment.

# 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 ..

Le dépôt contient également des notebooks comparant les API côte à côte, ainsi que des configurations Docker Compose pour exécuter localement Spark, Ray et Daft en mode distribué.


Points clés

  1. Démontrez qu’une seule machine ne suffit pas. Les résultats du scale-up montrent que certains gros scans analytiques tiennent sur un seul nœud. Mesurez la mémoire, les I/O, le spill et les besoins de reprise avant d’accepter le surcoût d’un cluster.

  2. Les modèles d’exécution comptent davantage qu’une syntaxe familière. La planification lazy, le pushdown, le streaming et les opérateurs parallèles expliquent une grande partie de l’écart entre un script DataFrame eager et un moteur analytique.

  3. Mesurez l’ensemble du pipeline d’accélération. Le décodage, les I/O réseau, le batching, la sérialisation et le backpressure déterminent si un GPU reste pleinement utilisé. Ray Data et Daft font de ce chevauchement une abstraction centrale ; Spark peut le prendre en charge avec une conception et des outils supplémentaires.

  4. Faites une présélection selon le workload, puis benchmarkez. Utilisez l’adéquation à l’écosystème pour réduire le champ, puis un job représentatif de bout en bout pour choisir. Les résultats de benchmark des fournisseurs sont des hypothèses, pas des garanties.

  5. Les formats ouverts rendent le choix réversible. Les interfaces compatibles avec Arrow, Parquet et les formats de tables ouverts peuvent réduire les coûts de transfert. Vérifiez si une conversion donnée réutilise les buffers ou les copie.

Pour commencer, essayez Polars pour un pipeline tabulaire local, et Daft ou Ray Data pour un pipeline mixte CPU/GPU. Choisissez Spark lorsque son exécution distribuée, son écosystème ou votre base opérationnelle existante répond à une exigence. La seule taille du dataset est un argument peu convaincant.


Références

Benchmarks et données de performance

Moteurs et frameworks

  • Polars - Bibliothèque Rust de DataFrame avec exécution lazy
  • Apache DataFusion - Moteur de requêtes Rust embarquable (GitHub)
  • Daft - DataFrame distribué nativement multimodal
  • Ray - Framework de calcul distribué pour l’IA (Data Internals)
  • Apache Spark - Moteur distribué d’ETL et de SQL
  • DuckDB - Moteur analytique SQL in-process
  • Dask - Calcul parallèle natif Python

Architecture et écosystème

Dépôt de démonstration

  • engine-comparison-demo - Code associé à cet article : benchmarks, notebooks et configuration Docker Compose pour Spark/Ray/Daft

Jeux de données

  • NYC Taxi Trip Records - Données des taxis jaunes de la NYC TLC (Parquet, environ 2,9 M de lignes par mois)
  • Food-101 Dataset - 101 000 images d’aliments provenant de l’ETH Zurich (Bossard et al., ECCV 2014)