Pandas 该被淘汰了
/
本文内容来自我在 Latency Conference 上的演讲。如果你只想:
你没看错,Pandas 应该被淘汰。不是那种用于熊猫外交的可爱毛绒动物,而是那个 Python DataFrame 库。
为什么?因为 Pandas 的低效会让你在数据量还没到那个程度时,就不得不用上分布式查询系统,白白增加复杂度。我认为绝大多数工作负载根本不需要这些系统,它们不过是被营销包装出来的“银弹”。
要理解我的意思,得先看看 Pandas 的典型使用路径。
为什么我们都在用 Pandas?
下面的图大致展示了根据数据规模,什么时候适合选用哪个 DataFrame 库。从左往右看,你也能看到数据分析工具的典型升级路径,以及 Pandas 用户在数据量超过一定规模后遭遇的“断崖”。
人们通常从 Excel 起步,数据到了 GB 量级后转向 Pandas。Pandas 能一直撑到几十 GB,之后就会遇到内存不足、计算缓慢,或者被它那晦涩的 API 折磨得心力交瘁。这时候的传统答案是升级到 Spark、DataBricks、Snowflake 或 Dask 这类专为 大数据 ™️ 设计的“专业”(也就是昂贵)工具。

关键在于:在“Pandas 断崖”和真正需要分布式系统之间,存在一段越来越大的空白区。这段空白大约在 100GB 左右,而现代高性能单机工具完全可以填补它。我指的主要是 Polars 和 DuckDB。
我们为何如此在意这个约 100GB 的阈值?答案在于理解现实中存在多少“大数据”。
我确实有大数据,对吧?
2024 年,Amazon 发布了一篇题为 《为何 TPC 不够:分析 Amazon Redshift 集群》 的论文。该论文旨在将 Amazon 自家分布式分析型数据库 Redshift 的遥测数据,与行业标准数据库基准测试中使用的查询模式进行对比。作为分析的一部分,Amazon 发布了关于查询运行时间和表大小的集群统计数据。


如果我们愿意接受几个假设,就能得出关于 Amazon 客户如何使用分析型数据库的有趣结论。让我们假设:
- Redshift 表中每行的平均大小为 1KB
- 每个 Redshift 集群由 10 台机器组成,每台机器能以 8GB/s 的速度从 S3 吸入数据,且仅执行此项任务
我们得出的发现是:
- Redshift 集群中 94.68% 的表包含的数据量少于 100GB
- 86.9% 的查询操作的数据量在 80GB 或更少
如果你对这份数据的更深入解读感兴趣,MotherDuck 的 Jordan Tigani 在此 做了一次深入挖掘。注意:MotherDuck 是一家提供 DuckDB 托管服务的 SaaS 公司,因此对其结论持些许保留态度或许更为恰当。
计算过程说明汇总运行时间表格的前 3 行数据,我们计算出 86.9% 的查询运行时间不到一秒。
基于集群中有 10 台机器、每台以 8GB/s 速度吸入数据的假设,我们计算得出:
10 台机器 * 8GB * 1 秒 = 80GB 数据8GB/s 的假设基于这个公认已有些过时的 基准测试。
关于“94.68% 的表格大小小于 100GB”这一结论,我们是这样推导的:首先累加行数直至 $10^8$ 的上限,得到占总行数 94.68% 的比例;接着假设每行数据大小为 1KB,计算如下:
10^8 rows in a table * 1KB = 100GB也许 1KB/行的假设过于乐观,但即便按 10KB 估算,表格大小也仅相当于 1TB。
这一切意味着什么?
你可能并没有“大数据”问题,将来恐怕也不会有。你面临的其实是“中数据”难题,因此需要的是“中数据”解决方案。
看看这些替代方案
我推荐的替代方案如前文所述,分别是 DuckDB 和 Polars。简而言之,Polars 是一个基于 Rust 的 DataFrame 库,使用感受与 Pandas 相似,但在若干关键方面存在差异,后文将详细探讨。DuckDB 则是一个内存分析数据库,本质上可视为用于分析场景的 SQLite。为了让大家感受这些工具的特质及其与 Pandas 的区别,我们来看一个具体示例。
10 亿行挑战要求编写一个尽可能快的 Java 程序,用于计算包含气象站数据的 10 亿行 CSV 文件中的最小值、平均值和最大值。该竞赛中入围的最快实现仅需 1.5 秒。
原挑战使用的是搭载 32 核 CPU 和 128GB 内存、运行 Debian 12 的裸机 Hetzner AX161 服务器。由于作者是个习惯性拖延且极度抠门的人,而 Hetzner 要求用户具备良好“信誉”才能租用大型服务器,因此这些测试转而使用了 AWS 的 m7a.8xlarge 实例,同样运行 Debian 12。
该基础配置与原挑战保持一致:在 AMD CPU 上运行 32 核 128GB 内存。不过,未使用裸机专用硬件可能会对结果的可复现性产生一定影响(抱歉)。
别废话,直接上代码
废话少说,我们直接看几种实现方式。
Pandas
对于任何用过 Pandas 的人来说,这段代码应该非常眼熟。我们从 CSV 文件读取数据,按气象站分组,然后计算汇总统计指标:最小值、平均值和最大值。
输出的序列化去哪了?性能测试跳过了原挑战中要求的输出序列化。各个实现都保留了序列化输出的能力,并以此对实现做了单元测试(例如 Pandas 的代码)。考虑到 1 Billion Row 挑战的输出格式并非标准格式,对各个库做序列化测试意义不大。
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)
)
这个例子需要记住的关键点是:Pandas 会逐步顺序执行每个计算步骤,且是即时执行的。它会读入全部数据集,先分组,再做聚合。
Polars
Polars 的代码看起来和 Pandas 差不多,但运行机制截然不同,下面就能看到。
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
)
数据以分块方式扫描,然后分组、聚合。关键细节在于 scan_csv 是惰性求值的,直到调用 .collect 才会执行整个查询流水线。这听起来像数据库术语,也确实如此——惰性求值让 Polars 能像数据库一样构建优化过的查询图,并借助数据库领域四十年的优化成果,以分块方式读取数据,并按需在多线程间并行处理。
和数据库类似,通过分别将 `.collect` 调用替换为 `.explain(streaming=True)` 和 `.explain(streaming=True, optimized=False)`,我们可以查看优化前后的查询执行计划。
优化后的查询执行计划这个查询计划并不复杂:它扫描 CSV,执行 2 列投影,然后进行聚合。在涉及过滤的查询中,我们通常期待看到谓词下推被应用,即在聚合前就过滤掉不相关的行。这与 Pandas 不同,Pandas 会将所有行加载到内存中再进行过滤。
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
DuckDB 的读法与普通 SQL 很像:选择列,应用聚合函数,并指定 group by 条件。
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
"""
)
这段代码的核心要点是:无论数据存储在哪里,DuckDB 都提供 SQL 接口来查询数据。它还支持查询内存中的 Python 对象,在上面的示例中,通过读取 CSV 创建的对象 `df` 就可以直接用 SQL 查询。
与 Polars 一样,DuckDB 构建并执行查询计划,根据查询引擎判断是否合适,以惰性、多线程、分块方式执行。通过追加调用 `.explain()`,我们也能查看 DuckDB 为该查询生成的执行计划。
优化后的查询执行计划再次强调,该执行计划依然很简单:扫描 CSV,进行 2 列投影,然后聚合。
注意:以下输出来自 M1 MacBook Air,该机器未用于下方的性能剖析结果。
┌─────────────────────────────────────┐
│┌───────────────────────────────────┐│
││ 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 │
│ ──────────────────── │
│ 0 Rows │
│ (0.00s) │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│ CREATE_TABLE_AS │
│ ──────────────────── │
│ 1 Rows │
│ (0.00s) │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│ PROJECTION │
│ ──────────────────── │
│ station │
│ min │
│ max │
│ avg │
│ │
│ 8888 Rows │
│ (0.00s) │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│ HASH_GROUP_BY │
│ ──────────────────── │
│ Groups: #0 │
│ │
│ Aggregates: │
│ min(#1) │
│ max(#2) │
│ avg(#3) │
│ │
│ 8888 Rows │
│ (121.01s) │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│ PROJECTION │
│ ──────────────────── │
│ station │
│ measurements │
│ measurements │
│ measurements │
│ │
│ 1000000000 Rows │
│ (0.29s) │
└─────────────┬─────────────┘
┌─────────────┴─────────────┐
│ TABLE_SCAN │
│ ──────────────────── │
│ Function: │
│ READ_CSV_AUTO │
│ │
│ Projections: │
│ station │
│ measurements │
│ │
│ 1000000000 Rows │
│ (316.38s) │
└───────────────────────────┘
性能测试结果
性能测试的方式是:为每个库编写一个独立脚本,分别处理磁盘上一个 10 亿行的 CSV 文件。
脚本由一个轻量的自制基准测试工具执行。该工具会启动一个全新的 Python 解释器来运行脚本,并使用 psutil 以 50 毫秒为间隔持续采集内存和 CPU 指标,直到子进程退出。每个基准测试先预热运行两次,再正式重复执行三十次。
这种方法谈不上完美,但就本文的对比目的而言,在准确性和开销之间取得了不错的平衡。
| 库 | 耗时中位数 | CPU 峰值中位数 | USS 峰值中位数 | Swap 峰值中位数 |
|---|---|---|---|---|
| Pandas | 4 分 28 秒 | 113.0% | 38.12 GB | 0 MB |
| Polars | 5.04 秒 | 3202.60% | 18.02 GB | 0 MB |
| DuckDB | 5.19 秒 | 3174.64% | 1.93 GB | 0 MB |
结果一目了然:Polars 和 DuckDB 比 Pandas 快得多,内存占用分别只有后者的约 1/2 和 1/19,性能已接近手工调优的 Java 实现。
Polars 的内存占用看起来还是偏高,我怀疑并非所有计算都以流式方式执行。DuckDB 则无可挑剔——只用了极少的代码、无需任何调优,就交出了惊人的性能。
本地开发环境的性能
不过,生产环境的性能提升只是一方面。DuckDB 和 Polars 在加快本地开发迭代速度上同样表现出色。
下面在一台性能不错但稍显过时的笔记本上重复同样的测试。我用的是 Framework 13,搭载 Intel i5-1135G7(8 核)和 16GB 内存。
| 库 | 耗时中位数 | CPU 峰值中位数 | USS 峰值中位数 | Swap 峰值中位数 |
|---|---|---|---|---|
| Pandas | 12 分 15 秒 | 110.35% | 15.67 GB | 21.02 GB |
| Polars | 39 秒 | 765.75% | 15.22 GB | 35.85 MB |
| DuckDB | 47 秒 | 807.0% | 546.87 MB | 0 MB |
Polars 和 DuckDB 再次交出了亮眼的性能答卷,不过 Polars 内存占用较高,还轻微触发了交换空间。相比之下,Pandas 跑起来慢如嚼蜡,内存开销也极大。
我能免费获得什么优势?
这项基准测试凸显了相比传统 Pandas 工作流的几个关键优势:
- 无痛多线程:Polars 和 DuckDB 会自动调用所有 CPU 核心,无需手动管理线程或进程。既然你花钱买了这些核心,就利用起来!
- 高效的内存使用与流式处理:DuckDB 和 Polars 均支持分块处理数据,有助于降低内存占用。
- 惰性求值:通过预先定义整个计算流程,这两个库都能优化执行计划,应用谓词下推(提前过滤数据)等技术,正如“正宗”数据库那样。
- 潜在的磁盘溢出:当优化后的操作超出内存限制时,这些工具内置了机制,可将中间结果智能地写入磁盘,这通常比依赖操作系统通用的交换机制更高效。
如果你对更标准的 TPC-H 基准测试结果感兴趣,可以点击此处查看。
关于 TPC-H 基准测试结果的说明 这些 TPC-H 基准测试由 Coiled 运行,该公司提供托管式 Dask 服务,并在这一领域与 Polars / DuckDB 争夺市场份额。这并不意味着他们的测试结果有误,但保持一定的怀疑态度(包括对我自己的怀疑)是必要的。
借助 Apache Arrow 无缝采用
我们都被那些光鲜亮丽的工具坑过,有些人甚至现在正被它们折磨着。在不重写所有代码的前提下测试这些新工具的关键,在于Apache Arrow。Arrow 正在成为列式数据内存表示的事实标准,其创始人正是 Pandas 的作者 Wes McKinney。Pandas 自 2023 年 4 月的2.0 版本起便支持 Arrow。
Polars 和 DuckDB 都原生支持 Arrow。这意味着你可以在 Polars、Pandas 和 DuckDB 之间移动 DataFrame 而无需复制内存,使得在不同框架间切换几乎“零成本”。这里唯一需要注意的坑是,Pandas 默认不会创建基于 Arrow 的 DataFrame,你需要在创建 DataFrame 时显式指定 dtype_backend 为 pyarrow。
示例:纽约出租车数据
让我们分析一下纽约出租车数据集(包含2009年至今的行程数据,按月存储为 Parquet 文件),看看在疫情期间(2019-2022年)现金支付是否变得不那么普遍。这需要处理大约3GB的 Parquet 数据。
在这个示例中,我们将实现用于读取 Parquet 文件和执行计算的函数,分别在 DuckDB 和 Pandas 中实现。这些函数的 Polars 版本留给读者作为练习 😉。
Pandas 代码:
COLUMNS = ["tpep_pickup_datetime", "payment_type"]
MIN_DATETIME = "2019-01-01"
MAX_DATETIME = "2023-01-01"
# 现金支付类型的枚举值
CASH = 2
def read_data_pandas(folder: Path) -> pd.DataFrame:
df = None
for entry in folder.iterdir():
if not entry.name.endswith("parquet"):
continue
temp_df = pd.read_parquet(entry, columns=COLUMNS, dtype_backend="pyarrow")
if df is not None:
df = pd.concat([df, temp_df])
else:
df = temp_df
# 将数据限制在我们感兴趣的时间范围内
df = df[df["tpep_pickup_datetime"] > pd.to_datetime(MIN_DATETIME)]
df = df[df["tpep_pickup_datetime"] < pd.to_datetime(MAX_DATETIME)]
df["month"] = df["tpep_pickup_datetime"].dt.month
df["year"] = df["tpep_pickup_datetime"].dt.year
df = df.drop(["tpep_pickup_datetime"], axis="columns")
return df
def calculate_cash_pandas(df: pd.DataFrame):
df = (
df.groupby(["year", "month", "payment_type"])
.agg({"payment_type": "count"})
.unstack(fill_value=0, level=2)["payment_type"]
.reset_index()
)
df["total_payments"] = df.iloc[:, 2:8].sum(axis=1)
df["cash_pct"] = (df[CASH] / df["total_payments"]) * 100
DuckDB 代码:
def read_data_duck(folder: str):
return duckdb.sql(f"""
select
datepart('year', tpep_pickup_datetime) year,
datepart('month', tpep_pickup_datetime) month,
payment_type
from '{folder}/*.parquet'
where
tpep_pickup_datetime > '{MIN_DATE}'
and tpep_pickup_datetime < '{MAX_DATE}'"""
)
def calculate_cash_duck(data):
return duckdb.sql(f"""
with total as (
select year, month, count(payment_type) payments from df
group by year, month
),
total_cash as (
select year, month, count(payment_type) cash from df
where payment_type={CASH}
group by year, month
)
select total.*, cash, (cash / total) * 100 cash_pct from total
join total_cash
on total.year=total_cash.year and total.month=total_cash.month
order by total.year, total.month
""").df() # 强制执行:物化为 df,否则执行是惰性的
接下来我们可以把这几种读写组合放在一起跑,看看在读数据、写数据或两者兼用时,选哪种工具更划算。
def pure_pandas(folder: Path):
data = read_data_pandas(folder)
df = calculate_cash_pandas(data)
def duck_reads_panda_thinks(folder: Path):
data = read_data_duck(folder).df()
df = calculate_cash_pandas(data)
def panda_reads_duck_thinks(folder: Path):
data = read_data_pandas(folder)
df = calculate_cash_duck(data)
def pure_duck(folder: Path):
data = read_data_duck(folder)
df = do_taxi_duck_compute(data)
用和之前相同的基准测试脚本、相同的笔记本配置,结果如下:
| 方案 | 耗时中位数 | 最大 CPU 占用中位数 | 最大 USS 中位数 | 最大 Swap 中位数 |
|---|---|---|---|---|
| 纯 Pandas | 41.88s | 146.10% | 14.52 GB | 1.92 GB |
| DuckDB 读、Pandas 算 | 28.39s | 793.7% | 14.79 GB | 1.22 GB |
| Pandas 读、DuckDB 算 | 29.25s | 765.4% | 12.39 GB | 0 MB |
| 纯 DuckDB | 21.70s | 814.95% | 216.76 MB | 0 MB |
同样的规律再次浮现。DuckDB 能充分利用机器的 CPU 核心,同时内存占用极低。一旦涉及 Pandas,速度立刻骤降,内存占用则急剧攀升。
如果你想了解疫情期间的现金使用是否下降,答案当然是肯定的。不过,相关不等于因果,请不要据此得出任何有意义的结论。

为什么不该听我的?
在科技行业,保持怀疑态度是好事,所以我整理了一份“别听我”的理由清单:
- 我只是一个跑基准测试的人。所有代码都开源,欢迎自行阅读,判断其中是否有缺陷。我强烈建议你这么做。
- 如果你深度绑定在 Pandas 生态系统中,迁移成本可能过高。不过自从我那次演讲以来,“迁移成本过高”的标准其实已经发生了巨大变化。
- 让 Pandas 再发展一阵子。Pandas 正在改进,虽然因其关键地位而进展缓慢。然而,DuckDB 和 Polars 的优势并不止于性能。据我经验,两者的 API 设计都更不易令人困惑;尤其是 DuckDB,SQL 是一项高可迁移的技能。
Polars vs DuckDB,谁更好?
这完全取决于你的工作负载、经验和个人偏好。数据工程师通常偏爱 SQL,软件工程师则更青睐 Polars。不妨两个都试试,看看哪个更顺手。
真正不该做的,是仅因 Pandas 性能不佳,就盲目引入分布式查询系统及其附带的所有复杂性。从长期来看,你真正需要这类系统的概率其实非常低。
高性能从未如此亲民,货比三家吧!