Motores de procesamiento de datos: Polars, DataFusion, Ray y Spark

Traducción automática Este artículo se tradujo automáticamente a partir de la versión original en inglés.

Durante años, Pandas se encargó del procesamiento tabular en memoria y Apache Spark, de la versión distribuida. Esta división funcionaba bien cuando los datos estaban estructurados.

Las pipelines modernas también procesan imágenes, audio y vídeo. En estas cargas de trabajo, la decodificación en CPU puede dejar sin trabajo a la inferencia en GPU, mientras que la recolección de basura de la JVM y el Global Interpreter Lock de Python limitan el throughput. Los motores más recientes utilizan Rust y Apache Arrow para reducir estos costes.

Para comparar las opciones, hice un benchmark de Pandas, Polars, DataFusion, Daft y Rust nativo con dos datasets reales. Spark y Ray se cubren en notebooks distribuidos independientes. Los viajes en taxi de Nueva York representan el procesamiento tabular; las imágenes de Food-101, una pipeline multimodal.

El repositorio engine-comparison-demo contiene el código. Ejecútalo en tu propio hardware, porque el rendimiento del motor depende de la máquina, el dataset y la carga de trabajo.


Tipos de cargas de trabajo de procesamiento de datos

Los motores están especializados, así que la primera pregunta que debes responder es qué tipo de carga de trabajo tienes realmente. La división principal es entre datos estructurados y multimodales.

Dos mundos de los datos: estructurados frente a multimodalesDos mundos de los datos: estructurados frente a multimodales

Los datos estructurados o tabulares implican filtrado, agregaciones y joins. Normalmente están limitados por la CPU y suelen caber en memoria. Muchas consultas de este tipo pueden ejecutarse en una sola máquina moderna, pero la memoria, la E/S, la concurrencia, la recuperación ante fallos y la forma de la carga de trabajo siguen determinando si resulta apropiado usar un clúster.

La AI multimodal procesa imágenes, audio o vídeo. La inferencia puede ejecutarse en una GPU mientras la decodificación en CPU, las lecturas remotas o el batching limitan la pipeline. La etapa que limita el throughput cambia según el tipo de instancia y los operadores, así que mide la utilización en todo el recorrido en lugar de dimensionar la GPU de forma aislada.

Los motores más recientes se dirigen a distintas partes de este espacio. La ejecución nativa puede reducir el overhead de Python, las interfaces compatibles con Arrow pueden abaratar el intercambio de datos y la ejecución en streaming puede limitar el uso de memoria en datasets mayores que la RAM. Ninguno de estos beneficios es automático para todos los operadores o conversiones.


Parte 1: procesamiento en un solo nodo

Pandas favorece un modelo de programación eager y en memoria. Muchas operaciones habituales crean resultados intermedios, y la ejecución paralela no se coordina mediante un optimizador de consultas. Esta decisión mantiene directo el código exploratorio, pero puede resultar costosa en exploraciones analíticas de mayor tamaño. En el benchmark PDS-H de Polars, con factor de escala 10, Pandas tardó aproximadamente 365 segundos en completar la suite, mientras que el motor de streaming de Polars tardó 3,89 segundos. El resultado de aproximadamente 94x describe ese benchmark y esa configuración; no es un factor de conversión general de Pandas a Polars.

Ejecución eager frente a lazyEjecución eager frente a lazy

Polars para pipelines tabulares locales

Polars es una opción por defecto práctica cuando una pipeline local con DataFrames ha superado la ejecución eager en un único proceso. Su API lazy construye un plan de consulta antes de ejecutar, lo que permite optimizaciones como predicate pushdown y projection pruning. Después, el motor puede ejecutar los operadores en paralelo y, cuando está disponible, en batches de streaming.

La diferencia aumenta a medida que crecen los datos. Con un factor de escala 100 (~100 GB), el motor de streaming de Polars terminó en 23,94 segundos frente a los 152,27 segundos de su propio motor en memoria: aproximadamente 6x más rápido, con datos que superan la RAM.

Lo siguiente es un esquema ilustrativo de la API, no un extracto literal de la versión actual de engine_comparison_examples.ipynb ni de 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 diferencia importante está en el modelo de ejecución. En la consulta lazy anterior, Polars puede introducir el filtro y la selección de columnas en la lectura de Parquet. Esto puede permitir omitir row groups y evitar la lectura de columnas no utilizadas. Una pipeline eager típica de Pandas no dispone de un plan de consulta global que optimizar, aunque el uso cuidadoso de filtros de Parquet, columnas seleccionadas y backends alternativos puede recuperar parte de la diferencia.

DataFusion para motores de consultas integrados

Mientras que Polars es una biblioteca que utilizas directamente, DataFusion es el motor sobre el que construyes otros motores. Da soporte a InfluxDB 3.0, GreptimeDB y al acelerador Comet Spark de Apple.

En una ejecución de ClickBench publicada por el proyecto DataFusion en noviembre de 2024, DataFusion lideró las configuraciones probadas de Parquet en un solo nodo. Posteriormente, el equipo de Embucket publicó un caso práctico del proveedor para TPC-H con factor de escala 1000 utilizando datafusion-cli sobre r6gd.metal, con 64 vCPUs, 512 GB de RAM, NVMe local y mediciones con la caché caliente. Q18 y Q21 fallaron o tardaron demasiado en la ejecución original, y después se reescribieron estructuralmente para obtener los resultados completados. Esto muestra una ampliación vertical en un nodo bajo esa configuración, no un resultado general de TPC-H.

Lo siguiente es un esquema ilustrativo de la API, no un extracto literal de la versión actual de engine_comparison_examples.ipynb ni de 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()

El punto fuerte de DataFusion es su modularidad. Sus APIs de extensión cubren catálogos personalizados, proveedores de tablas, reglas del optimizador y planes de ejecución. Por eso aparece siempre que alguien está construyendo una plataforma de datos personalizada o integrando un motor de consultas en su propio producto.

Daft para datos multimodales

Daft integra imágenes, audio, vídeo y embeddings en el flujo de trabajo con DataFrames. Proporciona expresiones nativas para operaciones como la decodificación de imágenes y las descargas desde URL, por lo que las rutas habituales de preprocesamiento no requieren un bucle de Python fila a fila. Esta integración es el principal motivo para elegirlo frente a un motor tabular generalista.

Lo siguiente es un esquema ilustrativo de la API, no un extracto literal de la versión actual de engine_comparison_examples.ipynb ni de 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()
)

La API resultará familiar si conoces Pandas, pero el motor utiliza la misma pila de Rust + Arrow que encontrarás en Polars y DataFusion. Daft resulta especialmente útil cuando las «filas» de tu DataFrame son imágenes, PDF o tensores.

Benchmark en un solo nodo

Ejecuté Pandas, Polars, DataFusion y Daft, además de Rust nativo mediante Polars-rs, sobre aproximadamente 41 millones de viajes en taxis amarillos de Nueva York correspondientes a todo 2024. Spark y Ray aparecen en los ejemplos distribuidos, no en este benchmark de un solo nodo. Los resultados completos están en el repositorio de demostración:

Resultados del benchmark: rendimiento en un solo nodo (alcances temporales mixtos)Resultados del benchmark: rendimiento en un solo nodo (alcances temporales mixtos)

El gráfico sirve como contexto para el benchmark complementario, no como una clasificación justa y comparable de motores: Pandas carga los datos de viajes y zonas antes de su función cronometrada, mientras que Polars, DataFusion, Daft y Rust incluyen dentro de la sección cronometrada distintos trabajos de lectura de archivos, registro o configuración de la consulta. El script repite cada operación tres veces e informa de la mediana después de cargar las importaciones; el estado de la caché del sistema de archivos, el calentamiento, la CPU y otras condiciones de ejecución pueden cambiar los valores. Alinea el límite temporal y vuelve a ejecutar las pruebas en condiciones controladas antes de interpretar el orden como un resultado de rendimiento.

Con este tamaño (~660 MB de Parquet), los tres motores Python más recientes terminan rápidamente. La cifra de 94x de Polars frente a Pandas procede de PDS-H; mi carga de trabajo de taxis con 41 millones de filas produjo una diferencia menor. En esta ejecución con alcances temporales mixtos, Polars, DataFusion, Daft y Polars-rs quedaron en el mismo rango general, y DataFusion registró el tiempo más bajo. Este orden corresponde a esta consulta y a este límite de medición, no a los motores en general. Cuando los motores están tan próximos, la adecuación de la API y las mediciones repetidas con datos representativos deberían guiar la decisión.


Parte 2: datos multimodales

El ETL tabular normalmente reduce los datos: filtra, agrega y escribe menos de lo que lee. Las pipelines de AI multimodal hacen lo contrario. Una única ruta de documento puede expandirse en docenas de fragmentos de texto y vectores de embedding.

En una pipeline PySpark convencional, las imágenes y el audio suelen entrar como datos binarios y pasar a bibliotecas Python para su decodificación o transformación. Cruzar el límite JVM/Python y serializar los resultados puede convertirse en una parte importante del tiempo de reloj. Spark puede utilizar rutas basadas en Arrow, UDF vectorizadas e integraciones con aceleradores, pero estas opciones requieren un diseño y unas mediciones explícitos.

Ejecución en pipeline y utilización de la GPU

Spark planifica el trabajo en etapas separadas por límites de shuffle. Una implementación directa puede colocar la descarga, la decodificación y la inferencia en una secuencia que deja esperando parte de la capacidad de la CPU o de la GPU. Spark admite la planificación de recursos GPU y plugins, por lo que el tiempo de inactividad no es inevitable. Sin embargo, el solapamiento y el backpressure no son propiedades automáticas de un trabajo ordinario de PySpark con DataFrames.

Modelo de ejecución en pipelineModelo de ejecución en pipeline

La ejecución en pipeline es la alternativa. En lugar de ejecutar las etapas una detrás de otra, el motor solapa la E/S, el trabajo de CPU y la inferencia en GPU. Cuando las etapas y los buffers están equilibrados, este solapamiento puede reducir el tiempo de inactividad y mejorar el throughput; el arranque y el vaciado, el backpressure o una etapa más lenta aún pueden hacer que los recursos esperen.

Benchmark de procesamiento de imágenes

Hice un benchmark de procesamiento de imágenes con 500 fotografías reales del dataset Food-101. Resultados del repositorio de demostración:

Resultados del benchmark: rendimiento multimodalResultados del benchmark: rendimiento multimodal

Polars y DataFusion no aparecen porque esta prueba se centra en operaciones nativas con imágenes, no en expresiones tabulares. En esta ejecución con 500 imágenes, la ruta de imágenes de Daft fue aproximadamente 3,7x más rápida que la línea base de Pandas + Pillow. El valor exacto de multimodal_results.json, 1.0148517909983639s, dividido por el valor de rust_multimodal_results.json, 0.235865209s, es 4.30267692, por lo que la implementación en Rust fue 4,3x más rápida. El README complementario muestra los valores redondeados 1.01s y 0.24s, cuya relación es 4,2x; se debe al redondeo mostrado, mientras que los valores JSON son la fuente autorizada. La muestra resulta útil para poner de manifiesto el overhead de orquestación, pero es demasiado pequeña para predecir por sí sola el throughput en producción.


Parte 3: procesamiento distribuido

Cuando una sola máquina ya no basta, tienes que decidir cómo repartir el trabajo.

Apache Spark

Spark es una opción madura para ETL tabular a gran escala, joins con mucho shuffle y organizaciones que ya operan su ecosistema. Su recuperación basada en lineage y su amplio soporte de plataformas son importantes cuando la fiabilidad y la familiaridad operativa pesan más que la velocidad local. Las pipelines mixtas de CPU/GPU requieren una configuración de recursos y un diseño más deliberados que la ruta SQL; ahí es donde Ray Data y Daft ofrecen una abstracción más centrada.

Lo siguiente es un esquema ilustrativo de la API de Spark, no un extracto literal de la versión actual de distributed_spark.ipynb ni 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 para computación heterogénea

Ray Data se diseñó pensando en las cargas de trabajo de AI. En lugar de las barreras entre etapas de Spark, utiliza un modelo de streaming que mantiene alimentadas las GPUs. La característica más importante para las pipelines de AI es la planificación de recursos mixtos: puedes declarar que un actor quiere «1 GPU y 4 CPUs» y que otro solo quiere CPUs, y Ray se encarga del resto.

Amazon informó de más de $120 millones de dólares de ahorro anual tras migrar determinadas cargas internas de procesamiento de datos de Spark a Ray. El informe cita una mejora de la eficiencia de costes del 91 % en la prueba de concepto y del 82 % en producción por GiB de entrada de S3. La escala es destacable, pero se trata de un caso práctico de migración para cargas concretas, no del ahorro esperado en una adopción típica de Ray.

Lo siguiente es un esquema ilustrativo de la API de Ray Data, no un extracto literal de la versión actual de distributed_ray.ipynb ni 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/")

La relación CPU-GPU merece atención. En los benchmarks de Anyscale, Ray Data siguió escalando al aumentar las CPUs por GPU, mientras que los demás motores probados llegaron a una meseta. La fuente informa de una aceleración de hasta 3x en inferencia de imágenes al reducirse la falta de CPU, porque los procesos de alimentación de CPU pueden seguir el ritmo de la GPU.

Daft distribuido

Daft escala horizontalmente mediante su motor Flotilla: un worker Swordfish por nodo, con Flotilla encima para realizar la planificación en todo el clúster. Swordfish gestiona la ejecución local en Rust y crea pipelines de E/S y computación mediante streaming en batches pequeños, de modo que cada nodo permanece ocupado sin esperar a la siguiente etapa.

En los benchmarks publicados por Daft, Flotilla fue entre 4 y 18 veces más rápido que las implementaciones de Spark probadas en cuatro cargas multimodales. La mayor diferencia se produjo en un caso de detección de objetos en vídeo. Interpreta el resultado como una evidencia de que el diseño de ejecución importa en estas pipelines; después compáralo con los resultados de Anyscale que aparecen a continuación y con una prueba local.

Anyscale, la empresa responsable de Ray, publicó un benchmark comparativo en el que Ray Data cerró o invirtió la diferencia en instancias con muchas CPUs después de aplicar tuning. En conjunto, ambos estudios de proveedores son un motivo para probar operadores representativos y relaciones CPU-GPU, no para inferir una clasificación estable.

Lo siguiente es un esquema ilustrativo de la API de Daft, no un extracto literal de la versión actual de distributed_daft.ipynb ni 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/")

Comparativa distribuida

CaracterísticaSparkRay DataDaft (Flotilla)
Modelo de ejecuciónUna tarea por core, basada en particionesTareas y actores en streamingUn Swordfish por nodo, batches en streaming
Puntos fuertesSQL/ETL masivo, tolerancia a fallosComputación heterogénea, saturación de GPUPipelines multimodales, memoria acotada
Controles de GPUPlanificación de recursos y plugins del ecosistemaRecursos explícitos para tareas y actoresPlanificación integrada de pipelines CPU/GPU
Tuning habitualExecutors, memoria, particiones, aceleradoresBatches, actores, object storeBatches, recursos, concurrencia de E/S
Caso de evaluación recomendadoTrabajos SQL/ETL grandes y con mucho shuffleEntrenamiento o inferencia con recursos mixtosIngesta y transformación multimodal

Parte 4: el papel de Rust y Arrow

Bajo la competencia entre motores, la convergencia es la historia más interesante. Polars, Daft y DataFusion utilizan Rust de forma intensiva, mientras que Ray incluye componentes nativos junto con sus APIs de Python. Todos pueden intercambiar datos mediante partes del ecosistema de Apache Arrow.

Convergencia de Rust y ArrowConvergencia de Rust y Arrow

La Arrow PyCapsule Interface (protocolo __arrow_c_stream__) proporciona a las bibliotecas compatibles una forma estándar de intercambiar streams de Arrow. Las transferencias pueden evitar la serialización fila a fila y reutilizar buffers cuando los esquemas y las distribuciones de memoria son compatibles. La materialización, el rechunking, la conversión de tipos o la transferencia al dispositivo aún pueden copiar datos, así que verifica el traspaso mediante profiling en lugar de asumir que es gratuito.

Lo siguiente es un esquema ilustrativo de la API, no un script independiente: presupone que events.parquet ya existe. El engine_comparison_examples.ipynb complementario crea .data/events.parquet en sus celdas de configuración antes del ejemplo de 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")

¿Por qué utilizan Rust estos motores más recientes para la ejecución nativa?

  1. No hay un garbage collector de tracing en el código Rust. La propiedad permite a los operadores nativos controlar de forma más directa la asignación y liberación de memoria. Esto puede reducir la latencia relacionada con el GC, pero no evita la presión de memoria ni los fallos por falta de memoria.

  2. Representaciones nativas compactas. Las estructuras de Rust no llevan cabeceras de objetos Java. El beneficio práctico depende de la distribución del motor: los sistemas JVM columnares también evitan representar cada valor como un objeto individual.

  3. Comprobaciones de concurrencia más estrictas. Las reglas de Rust seguro descartan muchas data races en tiempo de compilación. El código del motor aún puede contener bloques unsafe y bugs lógicos de concurrencia, pero el lenguaje reduce la superficie de fallos.

Combinadas con la distribución de memoria columnar de Arrow, estas decisiones pueden reducir el overhead de asignación y serialización. No lo eliminan en todos los límites; los objetos Python, el transporte de red, los esquemas incompatibles y el movimiento de datos entre dispositivos siguen siendo importantes.


Parte 5: adoptar varios motores

No tienes por qué hacer pasar toda la plataforma por un único motor. La composición resulta útil cuando el coste del traspaso entre motores es inferior al de hacer que un solo sistema gestione todas las cargas de trabajo.

Modelo de coexistencia: pipeline con varios motoresModelo de coexistencia: pipeline con varios motores

Una posible combinación utiliza Spark para los joins consolidados del lakehouse, escribe una tabla o un formato de archivos abierto y entrega el resultado a Ray Data o Daft para la inferencia en GPU o las transformaciones multimodales. Polars o DuckDB pueden encargarse del análisis local sobre los mismos archivos. Un equipo pequeño quizá solo necesite uno de estos motores; añade otro cuando un cuello de botella medido justifique la frontera operativa.

Lo que hace posible esta composición es la interoperabilidad. Parquet es el formato de archivos común para el traspaso en estos ejemplos, mientras que Arrow puede proporcionar el intercambio en memoria. La compatibilidad con Delta e Iceberg depende del motor y del conector, así que verifica cada frontera. Un traspaso mediante archivos también tiene costes de E/S y conversión; no es automáticamente gratuito.

Selección del motor según el caso de evaluación

Caso de evaluaciónLista cortaValidar
DataFrames analíticos localesPolars / DuckDBAdecuación de la API, memoria, mezcla de consultas
SQL/ETL distribuido existenteApache SparkComportamiento del shuffle, operaciones, coste
Trabajo paralelo nativo de PythonDask / RayOverhead del scheduler, modelo de fallos
Entrenamiento o inferencia mixta CPU/GPURay Data / DaftUtilización del acelerador, backpressure
Ingesta multimodalDaft / Ray DataOperadores nativos, comportamiento de los reintentos
Motor de consultas integradoDataFusionAPIs de extensión, cobertura SQL

Escalado diagonal

Durante mucho tiempo, la elección consistía en escalar verticalmente —una máquina más grande— o escalar horizontalmente —más máquinas—. Polars Cloud está haciendo algo intermedio, que denomina escalado diagonal. Escala horizontalmente mientras lee del almacenamiento en la nube para maximizar la E/S y, después, concentra el trabajo en un único nodo grande una vez que los filtros y las agregaciones han reducido los datos. El shuffle distribuido nunca llega a producirse.

La lección general es comparar el coste por trabajo completado correctamente, no el precio horario de la instancia. Un nodo más grande puede resultar más barato cuando la reducción del tiempo de ejecución compensa su tarifa superior, pero el resultado depende de la E/S, la presión de memoria y la cantidad de trabajo eliminada antes del shuffle.


Reproducir la comparación

El repositorio de demostración complementario contiene los scripts, notebooks y el entorno utilizados para las comparaciones. El comando rápido siguiente utiliza el dataset de taxis de todo 2024 que viene por defecto en el repositorio, en línea con la ejecución de mayor tamaño mostrada anteriormente.

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

El repositorio también incluye notebooks para comparar APIs en paralelo y configuraciones de Docker Compose para ejecuciones distribuidas locales de Spark, Ray y Daft.


Conclusiones principales

  1. Demuestra que una sola máquina no basta. Los resultados de escalado vertical muestran que algunas exploraciones analíticas grandes caben en un único nodo. Mide los requisitos de memoria, E/S, spill y recuperación antes de asumir el overhead de un clúster.

  2. Los modelos de ejecución importan más que la sintaxis conocida. La planificación lazy, el pushdown, el streaming y los operadores paralelos explican gran parte de la diferencia entre un script eager con DataFrames y un motor analítico.

  3. Mide toda la pipeline del acelerador. La decodificación, la E/S de red, el batching, la serialización y el backpressure determinan si una GPU permanece ocupada. Ray Data y Daft convierten este solapamiento en una abstracción central; Spark puede admitirlo mediante diseño y herramientas adicionales.

  4. Haz una lista corta por carga de trabajo y después ejecuta un benchmark. Utiliza la adecuación al ecosistema para reducir las opciones y un trabajo representativo de extremo a extremo para elegir. Los resultados de los benchmarks de proveedores son hipótesis, no garantías.

  5. Los formatos abiertos mantienen reversible la elección. Las interfaces compatibles con Arrow, Parquet y los formatos de tablas abiertos pueden reducir los costes de traspaso. Confirma si una conversión concreta reutiliza los buffers o los copia.

Si necesitas un punto de partida, prueba Polars para una pipeline tabular local y Daft o Ray Data para una pipeline mixta de CPU/GPU. Elige Spark cuando su ejecución distribuida, su ecosistema o tu base operativa existente resuelvan un requisito. El tamaño del dataset, por sí solo, es un motivo débil.


Referencias

Benchmarks y datos de rendimiento

Motores y frameworks

  • Polars - Biblioteca de DataFrames en Rust con ejecución lazy
  • Apache DataFusion - Motor de consultas integrable en Rust (GitHub)
  • Daft - DataFrame distribuido nativo para datos multimodales
  • Ray - Framework de computación distribuida para AI (Data Internals)
  • Apache Spark - Motor distribuido de ETL y SQL
  • DuckDB - Motor SQL analítico en proceso
  • Dask - Computación paralela nativa de Python

Arquitectura y ecosistema

Repositorio de demostración

  • engine-comparison-demo - Código complementario de este artículo: benchmarks, notebooks y Docker Compose para Spark/Ray/Daft

Datasets