•6 min read

Fast Data Science: DuckDB and Polars for High-Performance Analytics

Fast Data Science: DuckDB and Polars for High-Performance Analytics

If you are still using Pandas for multi-gigabyte data transformations in Python, you are paying a massive "Pandas tax": single-threaded CPU execution, eager in-memory copies, and memory bloat requiring 5x to 10x the raw dataset size in RAM before throwing an OutOfMemoryError.

In modern data pipelines, Polars and DuckDB have replaced Pandas as the standard toolkit for high-performance single-node analytics. Both leverage columnar storage, vectorized SIMD execution, and Apache Arrow, delivering 10x to 50x speedups while processing datasets larger than physical RAM.

Here is how they work, how they compare, and how to use them together in production.


Audio Briefing
0:00 / 0:00

The Core Problem with Pandas

Dataset on Disk (Parquet/CSV): 2.0 GB
Pandas in RAM:                12.0 GB - 18.0 GB (Peak during merge/groupby)
Polars in RAM:                 2.3 GB (Zero-copy Arrow + streaming)
DuckDB in RAM:                 0.8 GB (Chunked vectorized execution engine)

Pandas was designed in 2008 around NumPy 1D arrays. It has three fundamental bottlenecks:

  1. The GIL and Single-Threading: Operations run on a single CPU core unless using external wrappers.
  2. Eager Evaluation: Every intermediate step creates an entire new DataFrame copy in memory.
  3. Missing Data Overhead: Pandas historically cast integer columns with NaN to float64, doubling memory usage.

Advertisement

1. Polars: The Rust-Powered DataFrame Engine

Polars is written from scratch in Rust and built directly on the Apache Arrow columnar memory format. It processes data using all available CPU cores without Python GIL restrictions.

Eager vs Lazy Execution

In production, you should almost always use Polars' Lazy API (LazyFrame). Instead of running computations immediately, Polars constructs a Logical Plan, optimizes it (predicate pushdown, projection pushdown, slice pushdown), and executes it in parallel:

import polars as pl

# Construct lazy query plan — zero disk I/O occurs here
lazy_query = (
    pl.scan_parquet("s3://analytics-bucket/events/*.parquet")
    .filter(pl.col("timestamp") >= pl.date(2026, 1, 1))
    .filter(pl.col("event_type").is_in(["purchase", "subscription"]))
    .with_columns([
        (pl.col("amount_cents") / 100.0).alias("amount_usd"),
        pl.col("user_id").n_unique().over("country").alias("unique_users_per_country")
    ])
    .group_by(["country", "event_type"])
    .agg([
        pl.col("amount_usd").sum().alias("total_revenue"),
        pl.col("amount_usd").mean().alias("avg_order_value"),
        pl.len().alias("transaction_count")
    ])
    .sort("total_revenue", descending=True)
)

# Inspect the optimized query plan
print(lazy_query.explain())

# Execute optimized plan in parallel across all CPU cores
result_df = lazy_query.collect(streaming=True)
print(result_df)

Why Polars Query Optimization Matters

When you call .scan_parquet() with filters:

  • Predicate Pushdown: Polars inspects the Parquet metadata footer and skips entire row groups that don't match timestamp >= 2026-01-01 without reading the data from disk.
  • Projection Pushdown: Polars only reads the 4 columns referenced in the query (timestamp, event_type, amount_cents, country), ignoring the remaining 50 columns in the file.
  • Streaming Engine (streaming=True): Processes data in streaming micro-batches, allowing transformations on datasets that exceed your machine's physical RAM.

2. DuckDB: The "SQLite for Columnar Analytics"

While Polars provides a DataFrame API, DuckDB is an embedded in-process SQL OLAP database. It runs inside your Python process with zero external server dependencies, zero network latency, and native SQL dialect support.

Querying Remote Parquet and S3 Directly in SQL

DuckDB can execute SQL queries directly against compressed Parquet, CSV, or JSON files on disk or remote S3 without loading them into database tables first:

import duckdb

# Connect to in-process DuckDB instance (or persist to 'analytics.duckdb')
con = duckdb.connect()

# Enable S3 / HTTP filesystem extension
con.execute("INSTALL httpfs; LOAD httpfs;")
con.execute("""
    SET s3_region='us-east-1';
    SET s3_access_key_id='YOUR_KEY';
    SET s3_secret_access_key='YOUR_SECRET';
""")

# Query 100GB of remote Parquet files using Vectorized SQL
query = """
    SELECT 
        country,
        event_type,
        COUNT(DISTINCT user_id) AS unique_users,
        ROUND(SUM(amount_cents) / 100.0, 2) AS total_revenue_usd,
        ROUND(AVG(amount_cents) / 100.0, 2) AS aov_usd
    FROM read_parquet('s3://analytics-bucket/events/year=2026/*/*.parquet')
    WHERE event_type IN ('purchase', 'subscription')
    GROUP BY country, event_type
    HAVING total_revenue_usd > 10000
    ORDER BY total_revenue_usd DESC
    LIMIT 20;
"""

# Execute and fetch directly to Arrow, Polars, or Python dictionaries
results = con.execute(query).pl()  # Returns native Polars DataFrame
print(results)

3. Zero-Copy Interop: Polars + DuckDB + Apache Arrow

Because both Polars and DuckDB use Apache Arrow for memory representation, you can pass data between them with zero memory serialization or copy overhead:

import polars as pl
import duckdb

# 1. Load and clean data with Polars
df = pl.DataFrame({
    "user_id": [101, 102, 103, 104],
    "scores": [88.5, 92.0, 79.5, 95.0],
    "tier": ["gold", "platinum", "gold", "platinum"]
})

# 2. Run complex analytical window SQL in DuckDB directly on the Polars DataFrame
con = duckdb.connect()

# DuckDB can reference the 'df' Python variable directly in the FROM clause!
sql_result = con.execute("""
    SELECT 
        user_id,
        scores,
        tier,
        RANK() OVER (PARTITION BY tier ORDER BY scores DESC) as rank_in_tier,
        AVG(scores) OVER (PARTITION BY tier) as tier_avg
    FROM df
""").arrow()  # Zero-copy Arrow Table

# 3. Convert back to Polars instantaneously
final_df = pl.from_arrow(sql_result)
print(final_df)

Advertisement

4. Benchmark: 10 Million Rows (1.2 GB Parquet)

Aggregating and filtering a 10M-row dataset on an 8-core, 16GB RAM developer laptop:

OperationPandas 2.2Polars (Eager)Polars (Lazy)DuckDB (SQL)
Read Parquet + Filter4.82s0.82s0.29s0.31s
Group By + 4 Aggs3.15s0.41s0.28s0.24s
Window Function2.40s0.35s0.26s0.21s
Peak RAM Usage~4.6 GB~1.4 GB~0.7 GB~0.4 GB

5. Architectural Decision Matrix: When to Use What

                                Dataset Scale & Use Case
                                           │
         ┌─────────────────────────────────┴─────────────────────────────────┐
         ▼                                                                   ▼
  Single Machine (< 500GB)                                          Distributed Cluster (> 1TB)
         │                                                                   │
    ┌────┴──────────────────────────┐                              ┌─────────┴─────────┐
    ▼                               ▼                              ▼                   ▼
SQL-Heavy / S3 Parquet       DataFrame Transformations      Batch Pipeline       Real-time OLAP
    ▼                               ▼                              ▼                   ▼
 DuckDB                          Polars                       Apache Spark         ClickHouse
 (Embedded SQL)              (Rust Multi-thread)               / Ray               / StarRocks
  • Choose Polars when: You are writing feature engineering pipelines, ML data preprocessing, ETL scripts, or complex procedural dataframe logic where type safety and Pythonic expressions shine.
  • Choose DuckDB when: You want standard SQL, ad-hoc analytics on local/S3 Parquet files, app-embedded analytical dashboards, or integration with BI tools via standard ODBC/JDBC/Python connectors.
  • Avoid Spark for data under 100GB: A modern multi-core machine running Polars or DuckDB is frequently faster, cheaper, and 10x easier to maintain than spinning up a distributed Spark cluster with JVM and network serialization overhead.

You Might Also Like

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement