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

Pendant des années, Pandas s’est chargé du traitement de données tabulaires en mémoire, tandis que Apache Spark gérait les travaux distribués. Cette séparation fonctionnait bien tant que les données étaient structurées.

Les modèles modernes pipelines traitent également des images, de l’audio et de la vidéo. Dans ces types de charges de travail, le décodage CPU peut entraver les opérations d’inférence GPU, tandis que la collecte de déchets du JVM ainsi que le verrou d’interpréteur global de Python limitent le débit. De nouveaux moteurs exploitent Rust et Apache Arrow afin de réduire ces contraintes.

Afin de comparer ces solutions, j’ai effectué des benchmarks sur Polars, DataFusion, Daft, Ray Data et Spark à l’aide de deux datasets réels. Les trajets en taxi à New York constituent des données tabulaires, tandis que les images du dataset Food-101 représentent des pipeline multimodaux.

Le repoire engine-comparison-demo il contient le code. Exécutez-le sur votre propre matériel, car les performances du moteur dépendent de la machine, de dataset, ainsi que du volume de travail.

En résumé : Commencez par définir la structure du travail, et non par chercher un benchmark gagnant absolu. Polars constitue une excellente option par défaut pour les pipelines DataFrame locaux ; DataFusion convient aux équipes embedding qui ont besoin d’un moteur de requêtes ; Ray Data et Daft sont conçus pour des environnements mixtes de CPU/GPU pipelines ; quant à Spark, il reste un choix fiable pour le SQL distribué et les processus ETL. Les améliorations de performance présentées dans ce guide sont des résultats spécifiques à certains types de charges de travail, et non des promesses universellement valables. Exécutez à nouveau le benchmarks correspondant sur vos données et votre matériel avant de prendre une décision.


Types de charges de travail de traitement de données

Les moteurs sont spécialisés, ce qui signifie que la première question à se poser concerne le type de charge de travail réellement rencontrée. La distinction principale porte sur les charges structurées et les charges multimodales.

Deux mondes des données : structurées versus multimodales

Les données structurées ou tabulaires impliquent du filtrage, de l’agrégation et des jointures. Elles sont généralement limitées par CPU et tiennent souvent dans la mémoire. Selon Polars, environ 90 % des requêtes traitent moins de 1 TB., Ainsi, de nombreux agents peuvent s’exécuter sur une seule machine moderne au lieu d’un cluster.

Multimodal AI permet de traiter des images, de l’audio ou de la vidéo. L’inférence peut s’exécuter sur un GPU, tandis que la décodage CPU, les lectures à distance ou les contraintes liées au regroupement en lots limitent le pipeline. Ce rapport varie en fonction de l’instance et des opérateurs ; il est donc préférable d’évaluer l’utilisation sur l’ensemble du parcours plutôt que de définir uniquement les capacités du GPU de manière isolée.

Les moteurs plus récents ciblent des segments différents de ce domaine. L’exécution native permet de réduire la surcharge liée à Python, les interfaces compatibles avec Arrow facilitent l’échange de données en diminuant les coûts associés, et l’exécution en flux permet de limiter l’utilisation de la mémoire sur des datasets supérieurs à la RAM. Aucun de ces avantages n’est automatiquement disponible pour tous les opérateurs ou toutes les conversions.


Partie 1 : traitement sur nœud unique

Pandas privilégie un modèle de programmation « eager » en mémoire. De nombreuses opérations courantes génèrent des intermédiaires, et l’exécution parallèle n’est pas coordonnée par un optimiseur de requêtes. Ce compromis permet aux scripts d’exploration de rester simples, mais il peut devenir coûteux lors de scans analytiques de grande envergure. Dans Polars, PDS-H benchmark Avec un facteur d’échelle de 10, Pandas a nécessité environ 365 secondes pour traiter l’ensemble des données, tandis que le moteur en flux continu de Polars n’a demandé que 3,89 secondes. Ce résultat, soit une différence d’environ 94 fois, dépend spécifiquement de benchmark ainsi que de la configuration choisie ; il ne constitue donc pas un facteur de conversion général applicable de Pandas à Polars.

Exécution impatiente vs. exécution lente

Polars pour les données tabulaires locales pipelines

Polars Il s’agit d’une option par défaut pratique lorsque un DataFrame local pipeline dépasse les limites d’exécution en mode eager sur un seul processus. Son mode API différé permet de générer un plan de requête avant l’exécution, ce qui facilite des optimisations telles que le pushdown des prédicats et le tri des données à projeter. L’engine peut ensuite exécuter les opérateurs en parallèle, et le cas échéant, par lots en mode streaming.

L’écart s’agrandit à mesure que les données augmentent. À facteur d’échelle 100 (~100 GB), Le moteur de traitement en flux de Polars a terminé sa tâche en 23,94 secondes, contre 152,27 secondes pour son propre moteur en mémoire — soit environ 6 fois plus rapide, même sur des données dépassant la capacité de la RAM.

Depuis engine_comparison_examples.ipynb:

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 majeure réside dans le modèle d’exécution. Dans la requête différée présentée ci-dessus, Polars permet d’appliquer le filtre ainsi que la sélection de colonnes directement lors du balayage du fichier Parquet. Cela permet d’ignorer certains groupes de lignes et d’éviter la lecture de colonnes inutilisées. Un traitement immédiat typique avec Pandas pipeline ne dispose pas de plan global de requête à optimiser, bien que l’utilisation judicieuse des filtres Parquet, de la sélection de colonnes et d’alternatives de backend puisse combler en partie cette lacune.

DataFusion pour des moteurs de requêtes embarqués

Là où Polars est une bibliothèque que l’on utilise directement, DataFusion c’est le moteur sur lequel vous construisez d’autres moteurs. Il fournit l’énergie nécessaire InfluxDB 3.0, GreptimeDB, et celui d’Apple Accélérateur Comet Spark.

En novembre 2024 Exécution de ClickBench publiée par le projet DataFusion, DataFusion a été en tête des configurations Parquet à nœud unique testées. L’équipe Embucket a par la suite publié un Exécution de l’échelle TPC-H avec un facteur de 1000 sur un seul nœud de grande taille. Ces résultats illustrent les performances que peut atteindre une architecture à grande échelle dans des conditions spécifiques ; ils ne prouvent pas pour autant qu’un travail de charge de 1 TB doive systématiquement éviter l’utilisation d’un cluster.

Depuis engine_comparison_examples.ipynb:

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 son modularité. L’extension APIs prend en charge les catalogues personnalisés, les fournisseurs de tables, les règles d’optimisation ainsi que les plans d’exécution ; c’est pourquoi elle apparaît systématiquement lorsque quelqu’un développe une plateforme de données sur mesure ou embedding un moteur de requêtes pour son propre produit.

Approche naïve pour les données multimodales

Insensé Il permet d’intégrer des images, de l’audio, de la vidéo ainsi que embeddings au flux de travail du DataFrame. Il offre des expressions natives pour des opérations telles que la décodage d’images et le téléchargement depuis des URL, ce qui évite l’utilisation de boucles Python ligne par ligne pour les traitements préalables courants. C’est cette intégration qui constitue la raison principale de choisir cet outil plutôt qu’un moteur tabulaire généraliste.

Depuis engine_comparison_examples.ipynb:

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()
)

Le API vous semblera familier si vous connaissez Pandas, mais le moteur repose sur la même pile Rust + Arrow que celle utilisée par Polars et DataFusion. Daft se révèle particulièrement utile lorsque les « lignes » de votre DataFrame sont des images, des PDF ou des tenseurs.

Nœud unique benchmark

J’ai exécuté les quatre moteurs, ainsi que le code natif en Rust via Polars-rs, sur environ 41 M. Trajets de taxis jaunes à New York pour l’ensemble de l’année 2024. Les résultats complets se trouvent dans le repo de démonstration:

Benchmark Résultats : performances sur nœud unique

À cette taille (environ 660 Mo en format Parquet), les quatre moteurs plus récents terminent rapidement leur traitement. Le rapport 94 fois entre Polars et Pandas présenté ci-dessus provient de PDS-H ; mon ensemble de données représentant des trajets de taxis composé de 41 millions de lignes a donné un écart plus faible. Polars, DataFusion, Daft et Polars-rs se situent dans la même fourchette de performances, DataFusion affichant le temps total le plus bas lors de ce test. Il s’agit d’un résultat utile pour cette requête, mais ce n’est pas un classement fiable : c’est API l’adaptation du modèle et des mesures répétées sur des données représentatives qui devraient permettre de choisir parmi des moteurs aussi proches les uns des autres.


Partie 2 : données multimodales

L’ETL tabulaire réduit généralement les volumes de données : filtres, agrégations, et enregistrement d’une quantité inférieure à celle lue. Les processus multimodaux AI pipelines font le contraire. Un seul chemin de document peut se décomposer en dizaines de fragments de texte ainsi que en vecteurs embedding.

Dans un PySpark conventionnel pipeline, les images et les fichiers audio sont généralement introduits sous forme de données binaires avant d’être transmis à des bibliothèques Python afin d’y être décodés ou transformés. Le passage entre le environnement JVM et celui de Python, ainsi que la sérialisation des résultats, peuvent représenter une part importante du temps d’exécution total. Spark permet l’utilisation de chemins basés sur Arrow, de UDF vectorisées et d’intégrations avec des accélérateurs, mais ces choix exigent une conception et des mesures explicites.

Exécution en pipeline et utilisation de GPU

Spark planifie les tâches en étapes séparées par des limites de shuffle. Une implémentation simpliste peut exécuter le téléchargement, la décodage et l’inférence dans une séquence qui laisse inutilisée soit la capacité CPU, soit GPU. Spark prend en charge la planification des ressources via GPU ainsi que divers plugins, ce qui permet d’éviter que le temps d’inactivité ne soit inévitable ; il s’agit simplement de souligner que l’overlapping et la pression inverse ne sont pas des propriétés automatiques des jobs ordinaires basés sur PySpark DataFrame.

Modèle d’exécution en pipeline

Exécution en pipeline constitue l’alternative. Au lieu d’exécuter les étapes les unes après les autres, le moteur superpose les opérations d’entrée/sortie, CPU de traitement, et GPU d’inférence, de sorte que seul le composant le plus lent est en attente à un instant donné.

Traitement d’images benchmark

J’ai effectué des benchmarks de traitement d’images sur 500 photographies réelles provenant du Introduction aux systèmes d’agents alimentaires dataset. Résultats provenant de repo de démonstration:

Benchmark Résultats : Performances multimodales

Polars et DataFusion sont absents car ce test vise des opérations sur images natives plutôt que des expressions tabulaires. Lors de cette exécution sur 500 images, le chemin d’accès aux images géré par Daft a été 3,6 fois plus rapide que la solution de référence basée sur Pandas + Pillow, tandis que l’implémentation en Rust était 4,2 fois plus rapide. Cet échantillon est utile pour mettre en évidence les coûts liés à l’orchestration, mais il est trop restreint pour permettre une estimation fiable du débit en environnement de production.


Partie 3 : traitement distribué

Lorsqu’une seule machine devient insuffisante, il faut déterminer comment répartir les tâches. C’est à ce stade que les différences architecturales entre les moteurs commencent réellement à jouer un rôle important.

Apache Spark

Spark constitue une solution éprouvée pour les opérations ETL à grande échelle impliquant des tableaux de données, les jointures gourmandes en opérations de shuffle, ainsi que pour les organisations déjà intégrées à son écosystème. Sa capacité de récupération basée sur l’arborescence des données et son large soutien technique deviennent cruciales lorsque la fiabilité et la familiarité opérationnelle priment sur la vitesse locale. Les architectures mixtes CPU/GPU pipelines exigent une configuration des ressources et une conception pipeline plus méticuleuses que la voie SQL, domaine où Ray Data et Daft proposent une abstraction plus ciblée.

Depuis distributed_spark.ipynb:

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"))
)

Données Ray pour le calcul hétérogène

Données Ray a été conçu autour des charges de travail AI. Au lieu des barrières de phase de Spark, il utilise un modèle à flux continu cela permet de maintenir GPUs alimenté. La fonctionnalité qui justifie son existence est la planification mixte des ressources : vous pouvez indiquer qu’un acteur souhaite « 1 GPU, 4 CPUs » tandis qu’un autre ne demande que CPUs, et Ray s’occupe du reste.

Amazon a publié ses résultats financiers Des économies annuelles dépassant 120 millions de dollars. Après avoir transféré les charges de travail de traitement de données internes sélectionnées de Spark vers Ray, le rapport indique une efficacité coûteuse 91 % meilleure lors du test de concept, ainsi que 82 % de gain en environnement de production par GiB de données d’entrée S3. L’échelle est remarquable, mais il s’agit ici d’une étude de cas de migration portant sur des charges de travail spécifiques, et non d’économies attendues dans le cadre d’une adoption typique de Ray.

De distributed_ray.ipynb:

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/")

Un détail important à connaître : à mesure que l’on ajoute davantage de CPUs par GPU, Ray Data montre une croissance continue, tandis que d’autres moteurs atteignent leur plateau de performance. Dans Le benchmarks d’Anyscale, Le passage d’un rapport de 4:1 à un rapport de 32:1 CPU-à-GPU a permis d’obtenir une accélération de 3 fois en termes de déduction d’images, car les alimentateurs CPU ont finalement pu suivre le rythme des GPU.

Daft distribué

Daft s’élargit grâce à son Moteur de flottille: Un travailleur Swordfish par nœud, avec Flotilla qui gère le planification à l’échelle du cluster en amont. Swordfish s’occupe de l’exécution locale en Rust ainsi que des opérations d’entrée/sortie pipelines via un flux par petits lots, permettant ainsi à chaque nœud de rester actif sans avoir à attendre la phase suivante.

In benchmarks publié par Daft, Flotilla a fonctionné de 2 à 18 fois plus rapidement que les implémentations Spark testées sur quatre charges de travail multimodales. L’écart le plus important a été observé dans un cas de détection d’objets vidéo. Considérez ce résultat comme une preuve que la conception de l’exécution joue un rôle crucial pour ces pipelines, puis comparez-le aux résultats d’Anyscale présentés ci-dessous ainsi qu’à ceux obtenus lors de tests locaux.

Anyscale, l’entreprise derrière Ray, a publié un concurrent benchmark dans lequel Ray Data a comblé ou inversé l’écart concernant les instances à forte charge de CPU après ajustement. Ensemble, ces deux études menées par des fournisseurs justifient l’essai d’opérateurs représentatifs ainsi que des ratios CPU/GPU, sans pour autant permettre d’établir un classement définitif.

Depuis distributed_daft.ipynb:

import daft
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").url.download())

# Define a GPU UDF for embedding generation
@daft.udf(return_dtype=daft.DataType.list(daft.DataType.float32()))
class TextEmbedder:
    def __init__(self):
        from sentence_transformers import SentenceTransformer
        self.model = SentenceTransformer("all-MiniLM-L6-v2", device="cuda")

    def __call__(self, text_col):
        texts = text_col.to_pylist()
        embeddings = self.model.encode(texts, batch_size=32)
        return [emb.tolist() for emb in embeddings]

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

Comparaison distribuée

FonctionnalitéSparkDonnées Ray
Modèle d’exécutionTâche par cœur, basé sur des partitionsTâches de streaming et acteurs
Points fortsTraitement massif de données via SQL/ETL, tolérance aux pannesCalcul hétérogène, saturation de GPUTuyauterie multimodale, mémoire bornée
GPU de contrôlePlanification des ressources associée aux plugins de l’écosystèmeRessources explicites de tâche et d’acteurPlanification intégrée CPU/GPU pipeline
Ajustement typiqueExécutants, mémoire, partitions, accélérateursLot de traitement, acteur, stockage d’objetsTraitements par lots, ressources, concurrence E/S
Bon cas d’évaluationLes tâches volumineuses de type SQL/ETL et celles nécessitant beaucoup de shufflingEntraînement ou inférence à l’aide de ressources mixtesIngestion et transformation multimodales

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

Sous l’angle de la concurrence, c’est la convergence qui constitue l’aspect le plus intéressant. Polars, Daft et DataFusion utilisent largement Rust, tandis que Ray intègre des composants natifs en plus de son interface Python APIs. Tous ces outils peuvent échanger des données via des parties de leurs architectures respectives. Apache Arrow écosystème.

Convergence de Rust et Arrow

Le Interface PyCapsule flèche (__arrow_c_stream__ Ce protocole offre aux bibliothèques compatibles un moyen standard d’échanger des flux Arrow. Les transferts peuvent éviter la sérialisation ligne par ligne et réutiliser des buffers lorsque les schémas ainsi que les arrangements mémoire sont compatibles. La matérialisation, le rechunking, la conversion de type ou le transfert vers un dispositif peuvent néanmoins entraîner une copie des données ; il convient donc de vérifier ce processus par du profilage plutôt que de supposer qu’il s’effectue sans coût.

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 nouveaux moteurs utilisent-ils Rust pour l’exécution native ?

  1. Pas de collecteur de déchets à suivi dans le code Rust. Le mécanisme d’ownership confère aux opérateurs natifs un contrôle plus direct sur l’allouement et le libération de la mémoire. Cela peut réduire les latences liées au GC, mais il ne permet pas d’éviter les pressions mémoire ou les pannes dues à un manque de mémoire.

  2. Représentations natives compactes. Les structures Rust ne comportent pas d’en-têtes d’objet Java. L’avantage pratique dépend de la conception du moteur : les systèmes JVM colonnaires évitent également de représenter chaque valeur en tant qu’objet distinct.

  3. Vérifications de concurrence plus rigoureuses. Le langage Rust sécurisé permet d’éliminer de nombreuses courses de données au moment de la compilation. Le code de l’engine peut toutefois contenir des blocs non sécurisés ainsi que des erreurs de concurrence logique, mais le langage réduit considérablement les surfaces de défaillance possibles.

En combinaison avec le schéma de mémoire colonne de Arrow, ces choix permettent de réduire la charge liée à l’allocation et à la sérialisation des données. Ils ne parviennent cependant pas à l’éliminer à chaque frontière ; les objets Python, le transport réseau, les schémas incompatibles ainsi que les déplacements de dispositifs conservent une importance significative.


Partie 5 : adoption de plusieurs moteurs

Vous n’êtes pas obligé de faire passer l’ensemble de la plateforme par un seul moteur. La composition s’avère utile lorsque le coût lié à sa définition est inférieur à celui d’un système unique chargé de gérer toutes les charges de travail.

Le modèle de coexistence : multi-moteur Pipeline

Une solution de relais possible consiste à utiliser Spark pour les opérations de jointure typiques des lakehouses, à écrire les résultats dans un format de tableau ou de fichier ouvert, puis à transmettre ces données à Ray Data ou Daft afin d’y effectuer des inférences GPU ou des transformations multimodales. Polars ou DuckDB peuvent quant à eux servir à réaliser des analyses locales sur les mêmes fichiers. Une équipe plus petite n’a parfois besoin que d’un seul de ces moteurs ; on en ajoute un autre uniquement lorsque des goulots d’étranglement mesurables justifient l’extension des capacités opérationnelles.

Ce qui rend cette composition possible, ce sont les formats ouverts : Parquet, Delta, Iceberg, Arrow. Chaque moteur de la pile les lit et les écrit nativement, de sorte que le transfert entre les étapes se résume à un simple chemin de fichier.

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

Cas d’évaluationListe de sélectionValider
DataFrames analytiques locauxPolars / DuckDBAPI combinaison de capacité, de mémoire et de requêtes
SQL/ETL distribué existantApache SparkComportement de mélange, opérations, coût
Travail parallèle natif en PythonDask / RaySurcoût lié au planificateur, modèle de défaillance
Entraînement ou inférence mixte CPU/GPURay Data / DaftUtilisation de l’accélérateur, contrainte de contrepression
Ingestion multimodaleDaft / Ray DataOpérateurs natifs, comportement de tentative
Moteur de requête intégréDataFusionExtension APIs, couverture SQL

Échelle diagonale

Pendant longtemps, la solution envisagée se résumait soit à l’augmentation de l’échelle (un appareil plus puissant), soit à l’élargissement de l’échelle (plus d’appareils). Polars Cloud Il s’agit d’une approche intermédiaire qu’ils désignent sous le nom de diagonal scaling : on effectue une mise à l’échelle horizontale en lisant les données depuis le stockage cloud afin d’exploiter au maximum les opérations d’entrée/sortie, puis on réduit tout cela à un seul nœud majeur une fois que les filtres et les agrégations ont réduit la taille des données, évitant ainsi complètement le tri distribué.

La leçon principale à retenir est de comparer le coût par tâche réussie, et non le prix horaire de l’instance. Un nœud plus puissant peut s’avérer moins cher lorsque la réduction apportée par runtime compense son tarif plus élevé, mais le résultat dépend des opérations d’entrée/sortie, de la pression mémoire ainsi que de la quantité de travail qui peut être traitée avant un mélange des données.


Réaliser la comparaison

Le répertoire de démonstration accompagnateur Il contient les scripts, les notebooks ainsi que l’environnement utilisés pour les comparaisons. La commande rapide ci-dessous fait appel à un mois de données de taxis ; suivez les instructions dataset du répertoire afin de reproduire l’analyse sur l’ensemble de l’année présentée précédemment.

# Install dependencies
uv sync

# Run the tabular benchmark (~2.9M NYC taxi trips)
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 ..

Ce dépôt contient également des notebooks permettant des comparaisons côte à côte API, ainsi que des configurations Docker Compose pour exécuter localement de manière distribuée Spark, Ray et Daft.


Principaux enseignements

  1. Démontrer qu’une seule machine est insuffisante. Les résultats obtenus lors de l’escalade montrent que certaines analyses volumineuses peuvent être traitées sur une seule machine nodale. Il convient de mesurer les besoins en mémoire, en E/S, ainsi que les exigences liées aux fuites de données et à la récupération, avant d’accepter les coûts supplémentaires induits par le cluster.

  2. Les modèles d’exécution sont plus importants que la syntaxe familière. La planification différée, le pushdown, le traitement en flux et les opérateurs parallèles expliquent en grande partie l’écart entre un script DataFrame à exécution immédiate et un moteur analytique.

  3. Mesurer l’accélérateur dans son intégralité pipeline. La décodage, les opérations d’entrée/sortie réseau, le regroupement des tâches, la sérialisation ainsi que la contrainte de charge déterminent si un GPU reste actif. Ray Data et Daft permettent de traiter cette interdépendance comme une abstraction centrale ; Spark peut également y parvenir grâce à des architectures et des outils supplémentaires.

  4. Faites une première sélection en fonction de la charge de travail, puis benchmark. Utilisez l’adéquation avec l’écosystème pour réduire le champ des candidats, ainsi qu’un cas d’usage représentatif du cycle de vie complet, afin de choisir le fournisseur idéal. Les résultats obtenus grâce aux benchmark ne constituent que des hypothèses, et non des garanties.

  5. Les formats ouverts permettent de revenir en arrière facilement. Les interfaces compatibles avec les flèches, Parquet ainsi que les formats de table ouverts peuvent réduire les coûts de transfert. Vérifiez si une conversion donnée réutilise des buffers ou les copie.

Si vous avez besoin d’un point de départ, essayez Polars pour des données tabulaires locales pipeline, ainsi que Daft ou Ray Data pour des données de type CPU/GPU pipeline. Optez pour Spark lorsque son mode d’exécution distribuée, son écosystème ou sa base opérationnelle existante répondent à un besoin spécifique — et non seulement lorsque la taille des données dataset dépasse une certaine limite arbitraire.


Références

Benchmarks et données de performance

Moteurs et frameworks

Architecture et écosystème

Répertoire de démonstration

Datasets