Skip to content
HN On Hacker News ↗

Pandas Should Go Extinct

▲ 194 points • 101 comments • by __eddie__ • 4w ago • HN discussion ↗

Pangram verdict · v3.3

We believe that this entire text is human-written.

1 %

AI likelihood · overall

Human
100% human-written 0% AI-generated
SEGMENTS · HUMAN 1 of 1
SEGMENTS · AI 0 of 1
WORD COUNT 1,472
PEAK AI % 1% · §1
Analyzed
Sep 12
backend: pangram/v3.3
Segments scanned
1 windows
avg 1472 words each
Distribution
100 / 0%
human / AI fraction
Verdict
Human
Pangram v3.3

Article text · 1,472 words · 1 segments analyzed

Human AI-generated
§1 Human · 1%

You read that correctly, Pandas should go extinct. Not the cute fluffy things used for international diplomacy, but the Python DataFrame library. Why? Because Pandas’ inefficiencies force you to adopt distributed querying systems before your workloads justify the added complexity. I posit that most workloads will never justify those systems, they are just well marketed “silver bullets”. To understand what I’m talking about we first must understand the typical adoption pathway for Pandas. Why do we use Pandas? The diagram below shows a rough guide of when you typically would consider adopting a given DataFrame library based on the data size you are working with. Following it from left to right, you also see the typical adoption pathway for data analysis tools, and the cliff that Pandas’ users experience beyond a certain data size. People typically start with Excel and graduate to Pandas somewhere in the GB range. Pandas serves them well into the 10s of GBs range, and then they start hitting memory issues, slow computation, or become frustrated with Pandas’ baroque API. The traditional answer at this point is to graduate to a “real” (read: expensive) tool like Spark, DataBricks, Snowflake, or Dask designed for Big Data ™ Rough guide to dataframe libary adoption by data size Here’s the thing: there’s a growing gap between the “Pandas cliff” and the scale where distributed systems are genuinely necessary. This gap, sits somewhere around the 100GB mark, and can be effectively filled by modern, high-performance, single-machine tools. I’m primarily talking about Polars and DuckDB. Why do we care so much about this ~100GB threshold? The answer lies in understanding how much “Big Data” exists in the wild. I have Big Data, right? In 2024, Amazon published a paper entitled “Why TPC is not enough: An analysis of the Amazon Redshift fleet”. The aim of this paper was to compare telemetry data from Amazon’s own distributed analytics database, Redshift, with the query patterns used in industry standard database benchmarks. As part of their analysis Amazon published fleet statistics on query run times and table sizes. Bucketed query runtimes for the Amazon Redshift fleet Bucketed table sizes for the Amazon Redshift fleet If we’re willing to make a couple of assumptions we draw some interesting conclusions about how Amazon’s customers are using analytics databases. Let’s assume that: The average size of a row in a Redshift table is 1KB Every RedShift cluster is comprised of 10 machines that are each capable of guzzling data at 8GB/s from S3, and do nothing but this We find that: 94.68% of tables in the Redshift fleet contain fewer than 100GB of data 86.9% of queries operate on 80GB of data or less If you are interested in another, deeper look at this dataset, Jordan Tigani of MotherDuck did a deep dive here. Note: MotherDuck is a SaaS business selling DuckDB hosting, so some scepticism is perhaps warranted. Explanation of calculationsSumming up the first 3 rows of the runtime table we calculate that 86.9% of queries run in less than a second.Using our assumptions that we have 10 machines in the cluster guzzling data at 8GB/s we compute: 10 machines * 8GB * 1 second = 80GB of dataThe assumption of 8GB/s is based on this admittedly outdated benchmark.To arrive at the claim of “94.68% of tables contain less than 100GB” we sum the rows up to the 10^8 limit, giving us 94.68% of rows. We then take our assumption of 1KB per row and compute: 10^8 rows in a table * 1KB = 100GBPerhaps the assumption of 1KB/row is too optimistic, but even assuming 10KB, you still arrive at a table size of 1TB. But what does this all mean? You likely do not have Big Data, and probably never will. You have Medium Data problems, and need Medium Data solutions. Meet the alternatives The alternatives I propose, as alluded to earlier are DuckDB and Polars. In broad strokes, Polars is a Rust-based DataFrame library that feels familiar to Pandas, but differs in several important ways we will explore. DuckDB is an in-memory analytics DB - essentially SQLite for analytics. To get a feel for these tools and how they differ from Pandas let’s look at an example. The 1 Billion Row Challenge was a challenge to write the fastest Java program which could compute the min, mean and max of a 1 billion row CSV containing weather station data. The fastest implementation accepted for the competition ran in 1.5 seconds. The original challenge used a bare metal Hetzner AX161 server with 32 cores and 128GB of RAM running Debian 12. Because the author is a serial procrastinator, a skinflint and Hetzner requires you to develop a “reputation” in order to rent large boxes, a m7a.8xlarge from AWS was instead used for these tests, also running Debian 12. This fundamental configuration is the same as the original challenge: 32 cores and 128GB of RAM on an AMD CPU. However not using bare metal dedicated hardware may affect reproducibility somewhat(sorry). Shut up and show me the code Without further ado, let’s look at some implementations. Pandas This should look very familiar to anyone who has touched Pandas before. We read the data in from the CSV, group by the weather station and then compute the aggregate min, mean and max figures. Where’s the output serialisation?For performance tests the output serialisation specified in the original challenge is skipped. The implementations all include the ability to serialise the output, which was used to unit test the implementations (e.g. the Pandas code). Given that the output format for the 1 Billion Row challenge is non-standard, it didn’t feel like a relevant test of the various libraries to test serialisation. def do_1brc_pandas(file_path: str): df = ( pd.read_csv(file_path, sep=";", names=["station", "measurement"]) .groupby("station") .agg({"measurement": ["min", "mean", "max"]}) .round(2) ) The key part of this example to remember is that Pandas executes each step of this computation sequentially and eagerly. It reads in the entire dataset, groups it and then performs aggregation. Polars The Polars code looks similar to Pandas, but it works very differently at runtime as we will see. def do_1brc_polars(file_path: str): df = ( pl.scan_csv( file_path, separator=";", new_columns=["station", "measurement"], has_header=False, ) .group_by("station") .agg( pl.col("measurement").min().round(2).alias("min"), pl.col("measurement").mean().round(2).alias("mean"), pl.col("measurement").max().round(2).alias("max"), ) .collect(new_streaming=True) # Stream the input data and perform computations in chunks ) The data is scanned in chunks, grouped and aggregated. The key detail here is that scan_csv is lazily evaluated and the call to .collect executes the query pipeline. If this sounds like database terminology it should. This lazy evaluation allows Polars to construct an optimised query graph, similar to a database, and leverage 40 years worth of database optimisations to read the data in a chunk-wise fashion and parallelise the work across threads as necessary. Much like a database, we can visualise the optimised and unoptimised query plan by replacing our call to .collect with a call to .explain(streaming=True) and .explain(streaming=True, optimized=False) respectively. The optimised query planThis query plan isn’t hugely exciting, it scans the CSV, does a 2 column projection, and aggregates. For queries involving filtering we’d expect to see predicate push down applied, where rows are filtered before aggregation occurs. This is unlike Pandas, where all rows are loaded into memory and then filtered.AGGREGATE [col("measurement").min().round().alias("min"), col("measurement").mean().round().alias("mean"), col("measurement").max().round().alias("max")] BY [col("station")] FROM STREAMING: simple π 2/2 ["measurement", "station"] Csv SCAN [/Users/eddie/Documents/code/pandas-should-go-extinct/data/measurements.csv] PROJECT 2/2 COLUMNS DuckDB The DuckDB code reads like vanilla SQL - columns are selected with an aggregation function applied and a group by criteria. def do_1brc_duckdb(file_path: str): df = duckdb.read_csv(file_path, names=["station", "measurements"]) src = duckdb.sql(""" create table src as select station, min(measurements) min, max(measurements) max, cast(avg(measurements) as decimal(8, 1)) avg from df group by station """ ) The key takeaways from this code sample is that DuckDB provides an SQL interface over your data, however and wherever it is stored. It also has the ability to query Python objects in memory, in the listing above the object df is created by reading the CSV and is queried using SQL. DuckDB, much like Polars, constructs and executes query plans which are applied in a lazy, multi-threaded, chunk-wise manner depending on if the query engine deems it appropriate. By tacking on a call to .explain() we can also view the query plan that DuckDB generates for the query. The optimised query planAgain, the query plan isn’t hugely exciting, it scans the CSV, does a 2 column projection, and aggregates.Note: this is the output from running the plan on my M1 MBA, which was not used for performance profiling results below.┌─────────────────────────────────────┐ │┌───────────────────────────────────┐│ ││ Query Profiling Information ││ │└───────────────────────────────────┘│ └─────────────────────────────────────┘ explain analyze create or replace table src as select station, min(measurements) min, max(measurements) max, cast(avg(measurements) as decimal(8, 1)) avg from df group by station ┌────────────────────────────────────────────────┐ │┌──────────────────────────────────────────────┐│ ││ Total Time: 56.08s ││ │└──────────────────────────────────────────────┘│ └────────────────────────────────────────────────┘ ┌───────────────────────────┐ │ QUERY │ └─────────────┬─────────────┘ ┌─────────────┴─────────────┐ │ EXPLAIN_ANALYZE │ │ ──────────────────── │