dask,一个并行计算的 Python 库!
据IDC预测,2025年全球数据量将达到175 ZB。然而大部分人的电脑内存只有 16GB,一个 50GB 的 CSV 文件根本无法被 Pandas 加载。这时候,你是否只能选择升级服务器、买更大内存的云主机,或者望洋兴叹?
2026 年第一季度,Python 生态中有一个低调而强大的库——Dask,其 PyPI 月下载量已突破 800 万,全球使用率在过去一年增长了 47%。如果你每天和数据打交道,那么 P****s 你多半用过,但 Dask 你可能只是听说过。在解决“单机内存装不下、计算太慢跑不动”这两个核心问题上,Dask 展现出了惊人的能力。
Dask 是一个开源并行计算库,从单机多核扩展到大型集群只需极少的代码改动,Dask 的技术方案跟 Pandas 内存不兼容时的多页模式、NumPy 内存不够时的跨核心处理模式有关。Dask 在底层会将数据按块切分,为每个块单独构造一个 Python 内置的数据结构(Pandas DataFrame 或 NumPy Array),再把这些块拼接成完整的 Dask DataFrame 或 Dask Array。你写代码的时候仍然在用 Pandas 或 NumPy 那套习惯的语法,但实际的运算会并发到多个 CPU 核心上,数据不足时还能按需加载不会撑爆内存。这才是 Dask 真正改变游戏规则的地方。
一、Dask 简介:它在日常生活中有何实际用处?
想象一下你有一个电商工作室,每天从淘宝、京东、抖音小店导出的订单 Excel 文件加起来有 300 多个。如果你用 Pandas,第一步就会卡住——30 多张 Excel 每张 200MB,Pandas 一条一条读下来耗时超 40 分钟,中途还频繁内存溢出。
换成 Dask 的话,你只需要把 pd.read_csv 改成 dd.read_csv,然后像写普通 Pandas 代码一样继续分组聚合和统计,Dask 会在底层自动把 300 多个文件并行加载到多个核心。整个过程通常几十秒就能跑完。
面对规模飙升到 TB 级的海量数据,Dask 还支持部署一个能够跨多台机器并行计算的分布式集群。这意味着你可以用一批低价机器跑出单台服务器无法达到的计算吞吐量,工程层面的成本大幅缩减。
当今数据科学的应用场景中,无论是财务分析的亿级用户行为记录、气象研究的上百年卫星数据,还是医疗领域的四五十GB超大表格,Dask 都在为原本“单机跑不动”的问题提供务实的工程解决方案。
二、安装 Dask
安装 Dask 主库及相关组件很简单:
bash
pip install dask pip install dask[complete] # 安装全部依赖,包括分布式调度、DataFrame扩展等
若只是使用 dask.dataframe 和 dask.array 等核心组件,安装最小包即可。安装完成后验证:
python
import dask import dask.dataframe as dd print(dask.__version__) # 显示最新版,如 2026.3.0
三、基本用法——一份典型的数据分析工作流
这里以处理一个真实的电商销售记录为例,分四个小步骤来演示 Dask 的基本用法。
步骤 1:创建 Dask DataFrame
假设我们有一个包含 1 亿行订单记录的 CSV 大文件,Pandas 直接读取会因内存不足而崩溃,而 Dask 通过分块读取的方式优雅地解决这个问题:
python
import dask.dataframe as dd
# 分块读取大型 CSV,Dask 会根据机器内存自动估算分区数
df = dd.read_csv('sales_2026_*.csv', assume_missing=True)
# 查看分区结构
print(df.npartitions)
步骤 2:初步数据探查
Dask 操作在调用 .compute() 前只构建逻辑执行计划,不会执行真正的计算,因此可以使用相同的 Pandas 语法放心编写探查代码:
python
# 查看前几行(仅读取相关分区的前几行) print(df.head()) print(df.shape.compute()) # 查看 schema print(df.dtypes)
步骤 3:数据清洗与过滤
python
# 过滤掉无效销售额
clean_df = df[df['sales_amount'] > 0]
# 处理缺失值——用每种商品的均值来填充缺失的销售量
clean_df = clean_df.assign(
quantity=clean_df.groupby('product_id')['quantity'].transform(
lambda x: x.fillna(x.mean())
).astype('int')
)
步骤 4:分组聚合与计算触发
python
# 按区域和商品类别进行聚合
result = (
clean_df.groupby(['region', 'category'])
.agg({
'sales_amount': ['sum', 'mean'],
'quantity': 'sum'
})
)
# 触发真实计算——Dask 在此处才会将延迟执行的任务分发到各个核心
final_result = result.compute()
print(final_result.head())
注意 final_result 返回的是 Pandas DataFrame,这意味着计算结果可以被无缝地倒入其他 Python 数据工具中进行后续分析和可视化。
四、高级用法:dask.delayed 精细化并行与本地分布式集群
4.1 dask.delayed:解耦复杂任务流
当遇到数据比较零散、不适合用 DataFrame 进行统一批处理时,dask.delayed 是一个极其灵活的低级接口。dask.delayed 不是简单地把函数扔进线程池,而是构建一个延迟执行的有向无环图(DAG),后续能做任务调度、重试、内存感知和跨节点分发。
python
from dask import delayed
import time
def process_file(filepath, threshold):
"""模拟读取一份日志并做统计"""
# 模拟 IO 操作
time.sleep(0.5)
return {"file": filepath, "count": 100, "sum": 2000}
def aggregate(results):
"""汇总多个文件的统计"""
return {"total_count": sum(r["count"] for r in results),
"total_sum": sum(r["sum"] for r in results)}
# 构建 DAG——所有过程先不执行
file_list = [f"log_{i}.txt" for i in range(12)]
lazy_tasks = [delayed(process_file)(f, 1000) for f in file_list]
lazy_aggregate = delayed(aggregate)(lazy_tasks)
# 真正执行并计算结果
result = lazy_aggregate.compute()
print(result)
有了 dask.delayed,你就可以像搭乐高积木一样把函数调用按依赖关系拼起来,等 .compute() 时才真正执行,非常适合处理非标准 I/O、API 分页请求和多阶段 ETL 等场景。
4.2 启动本地分布式集群(LocalCluster)
当计算量显著增大时,可以启动一个本地分布式集群来调度任务——它会在单机上模拟分布式的完整机制,跑通了再无缝切换到真实集群,节约了大量联调成本:
python
from dask.distributed import Client, LocalCluster
cluster = LocalCluster(n_workers=4, threads_per_worker=2, memory_limit='8GB')
client = Client(cluster)
print(client.scheduler_info())
# 执行一个稍微复杂的并行任务
def costly_computation(x):
# 模拟耗时运算
return sum(i**2 for i in range(x))
futures = client.map(costly_computation, range(10_000))
results = client.gather(futures)
print(f"计算结果总数: {len(results)}")
# 查看仪表盘——访问 http://localhost:8787 可以看到实时的 Worker 信息
client.close()
五、实际应用场景:40GB 医疗数据高效处理
我们来写一个可以直接复现的 Python 脚本,演示用 Dask 快速处理 40GB MIMIC-IV(重症监护医学信息集市)医疗数据案例。
python
import dask.dataframe as dd
from dask.distributed import Client
# 1. 启动本地集群——在 8GB 笔记本上也能跑
client = Client(n_workers=2, threads_per_worker=2)
print(client.dashboard_link)
# 2. 分块读取超大型数据集
# 假设我们的 Mimic_Patients 文件夹下有 patients.parquet / stays.parquet / admissions.csv
df_patients = dd.read_parquet("/data/mimic_iv/patients.parquet")
df_stays = dd.read_parquet("/data/mimic_iv/icustays.parquet")
df_admissions = dd.read_csv("/data/mimic_iv/admissions.csv",
blocksize="100MB",
assume_missing=True)
# 3. 数据合并操作
merged = (
df_patients.merge(df_stays, on="subject_id", how="inner")
.merge(df_admissions, on="hadm_id", how="left")
)
# 4. 计算院内死亡率初步指标——使用 groupby 统计各性别入院次数
result = (
merged.filter(merged["deathtime"].notnull())
.groupby(["gender", "anchor_age_group"])
.agg({
"subject_id": "count",
"los": "mean"
})
)
# 5. 并行计算最终结果
final = result.compute()
print("死亡率统计初步结果:")
print(final.head())
# 6. 把结果存为 Parquet(列式存储,后续快速加载)
merged.to_parquet("/output/mimic_cleaned/", write_index=False)
client.close()
为什么这段代码有意义:
-
一台普通的 16GB 内存笔记本就能承载原本需要 64GB 服务器才能跑的医疗分析任务
-
blocksize="100MB"确保每个分区的大小不会过于夸张,内存平稳可控 -
Parquet 格式与 Dask 天然适配,后续可以直接加载
dd.read_parquet("/output/mimic_cleaned/"),速度比 CSV 快 5~10 倍 -
整个过程代码风格与 Pandas 几乎一致,降低了学习曲线和重构成本
在金融行业,有量化研究团队基于 Dask 构建了遗传规划高频因子挖掘框架,借助 Dask 的分布式计算图与惰性求值实现了分钟频数据分块存储和按需加载,在突破内存限制的同时保持了高性能计算。在气象领域,科研人员用 Dask 整合了卫星遥感和地面观测站的 TB 级数据以挖掘降水时空规律。这些现实案例都说明了 Dask 如何在大型数据量面前帮助用户落地。
六、总结
Dask 通过构建声明式计算图(Task Graph)和延迟执行调度层,将原本串行的操作转化为可分布调度的任务流。这种“最小侵入性、最大扩展性”的设计理念,让开发者可以在几乎不重写代码的前提下实现从笔记本电脑到千节点集群的平滑过渡。当 Pandas 遇上内存不足,当 NumPy 处理不了超大数组,当你的分析流程需要真正的横向并行时,Dask 是一座坚实的工程桥梁。
但 Dask 也遵循“如果 Pandas 能在几秒内跑完,不必为用而用”的原则——它引入的任务调度、序列化和跨核心协调都有额外开销,不适合微数据量或高度串行的算法。
今天的分享就到这里,如果你有过“单机内存装不下真实数据”的苦恼,或者用 Dask 处理过类似的项目,欢迎在评论区留言交流你遇到的瓶颈和解决方式!
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐


所有评论(0)