Dask 分布式 DataFrame

FreeGuideOnline 最新 2026-07-12

bash pip install dask[dataframe] distributed


#### 创建 Dask DataFrame

从 Pandas DataFrame 转换:

```python
import pandas as pd
import dask.dataframe as dd

# 创建大型 Pandas DataFrame
df_pd = pd.DataFrame({'x': range(100000), 'y': range(100000, 200000)})

# 转换为 Dask DataFrame,分区数 npartitions=4
ddf = dd.from_pandas(df_pd, npartitions=4)
ddf

从 CSV 文件读取(自动分区):

ddf = dd.read_csv('large_file.csv', blocksize='64MB')

blocksize 参数控制每个分区的大小,Dask 会根据文件大小自动创建多个分区。

惰性计算与触发执行

# 构建计算图——此时不会执行
result = ddf.groupby('category').amount.mean()

# 触发计算并获得 Pandas DataFrame
computed = result.compute()

你还可以使用 .head() 先预览前几行,验证逻辑:

ddf.head()

分区机制详解

Dask DataFrame 的每个分区都是一个独立的 Pandas DataFrame。了解分区数是性能调优的基础。

查看分区数和每个分区的界限:

ddf.npartitions          # 分区数量
ddf.divisions            # 每个分区的最小索引值(如果索引已排序)

常见设置分区的方式:

  • 从文件读取时dd.read_csv('*.csv', blocksize=None) 或使用通配符。
  • repartition 方法:增加或减少分区。
# 合并为 2 个分区
ddf_small = ddf.repartition(npartitions=2)

# 按列值重新分区(隐式排序并建立已知分区)
ddf_sorted = ddf.set_index('timestamp')

常用数据操作

Dask DataFrame 支持绝大多数 Pandas 操作,以下是一些示例:

选择与过滤

# 选择列
ddf[['x', 'y']]

# 布尔过滤
ddf[ddf.x > 50000]

分组聚合

ddf.groupby('category').amount.agg(['sum', 'mean']).compute()

注意:分组操作会触发大量数据移动,如果分组键基数很高,性能会下降。

连接与合并

merged = dd.merge(left_ddf, right_ddf, on='key')

当两个 DataFrame 都按照连接键排序且分区对齐时,连接非常高效。否则会触发昂贵的 shuffle 操作。

应用自定义函数

ddf['new_col'] = ddf['x'].map(lambda v: v * 2, meta=('x', 'int64'))

meta 参数用于告知 Dask 输出列的名称和类型,对性能至关重要。

分布式计算与集群

当本地资源不足时,你可以将 Dask 连接到分布式集群。

创建本地分布式集群(多核)

from dask.distributed import Client, LocalCluster

cluster = LocalCluster(n_workers=4, threads_per_worker=2)
client = Client(cluster)
client

现在所有 Dask 操作都会自动使用集群资源。不要忘记在Client创建之后再读取数据或定义计算。

查看任务图

result.visualize(filename='task_graph.png')

监控仪表盘

启动客户端后,浏览器访问 http://localhost:8787/status 查看任务进度、内存使用和工作者状态。

实践:大规模数据清洗与分析

假设你需要处理一个 50GB 的 CSV 日志文件,统计每小时请求量。

import dask.dataframe as dd
from dask.distributed import Client

# 启动集群
client = Client(n_workers=4)

# 读取数据,只选择需要的列以减少内存
ddf = dd.read_csv('server_logs_*.csv',
                  usecols=['timestamp', 'status'],
                  parse_dates=['timestamp'],
                  blocksize='200MB')

# 添加小时列
ddf['hour'] = ddf.timestamp.dt.floor('H')

# 过滤成功请求并分组计数
hourly_requests = ddf[ddf.status == 200].groupby('hour').size()

# 触发计算并保存结果
hourly_requests.compute().to_csv('hourly_requests.csv', index=True)