8.7 并行计算 (例如:Dask 与 Pandas 集成)


文档摘要

8.7 并行计算 (例如:Dask 与 Pandas 集成) 8.7 并行计算:Dask 与 Pandas 集成 在处理大型数据集时,Pandas 的性能可能会成为瓶颈。单核 CPU 运算限制了处理速度,使得某些操作耗时过长。为了解决这个问题,我们可以利用并行计算技术,将任务分解成多个子任务,分配给多个 CPU 核心同时执行,从而显著提高数据处理速度。Dask 是一个灵活的并行计算库,它与 Pandas 紧密集成,可以无缝地扩展 Pandas 的处理能力,使其能够处理超出内存的数据集。 8.7.1 Dask 简介 Dask 是一个用于并行计算的 Python 库。它提供了两种核心组件: Dask Delayed: 允许你延迟函数的执行,并创建一个任务图,描述了各个任务之间的依赖关系。

8.7 并行计算 (例如:Dask 与 Pandas 集成)

8.7 并行计算:Dask 与 Pandas 集成

在处理大型数据集时,Pandas 的性能可能会成为瓶颈。单核 CPU 运算限制了处理速度,使得某些操作耗时过长。为了解决这个问题,我们可以利用并行计算技术,将任务分解成多个子任务,分配给多个 CPU 核心同时执行,从而显著提高数据处理速度。Dask 是一个灵活的并行计算库,它与 Pandas 紧密集成,可以无缝地扩展 Pandas 的处理能力,使其能够处理超出内存的数据集。

8.7.1 Dask 简介

Dask 是一个用于并行计算的 Python 库。它提供了两种核心组件:

  • Dask Delayed: 允许你延迟函数的执行,并创建一个任务图,描述了各个任务之间的依赖关系。

  • Dask DataFrames: 模仿 Pandas DataFrame 的 API,但可以处理大于内存的数据集。Dask DataFrames 将数据分割成多个小的 Pandas DataFrames (称为 partitions),并并行地在这些 partitions 上执行操作。

Dask 的核心优势在于:

  • 并行性: 充分利用多核 CPU 和分布式计算资源。

  • 延迟执行: 只有在需要结果时才执行计算,避免不必要的计算开销。

  • 内存管理: 可以处理超出内存的数据集,通过磁盘溢出等机制进行内存管理。

  • 与 Pandas 集成: 与 Pandas DataFrame API 高度兼容,学习曲线平缓。

8.7.2 Dask DataFrame 的创建

Dask DataFrame 可以从多种数据源创建,包括 CSV 文件、Parquet 文件、数据库等。

1. 从 CSV 文件创建:

import dask.dataframe as dd import pandas as pd # 创建一个示例 CSV 文件 data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'B', 'C', 'D', 'E', 'F', 'G', 'H', 'I', 'J']} df = pd.DataFrame(data) df.to_csv('example.csv', index=False) # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv('example.csv') # 查看 Dask DataFrame 的信息 print(ddf.head()) print(ddf.npartitions) # 查看分区数量

这段代码首先创建了一个小的 Pandas DataFrame,并将其保存为 CSV 文件。然后,使用 dd.read_csv() 函数从 CSV 文件创建 Dask DataFrame。ddf.head() 显示 Dask DataFrame 的前几行,ddf.npartitions 显示 DataFrame 被分割成了多少个分区。默认情况下,Dask 会根据文件大小和可用内存自动确定分区数量。

2. 从 Pandas DataFrame 创建:

import dask.dataframe as dd import pandas as pd # 创建一个 Pandas DataFrame data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'B', 'C', 'D', 'E', 'F', 'G', 'H', 'I', 'J']} pdf = pd.DataFrame(data) # 使用 Dask 从 Pandas DataFrame 创建 Dask DataFrame ddf = dd.from_pandas(pdf, npartitions=2) # 指定分区数量 # 查看 Dask DataFrame 的信息 print(ddf.head()) print(ddf.npartitions)

这段代码展示了如何从现有的 Pandas DataFrame 创建 Dask DataFrame。dd.from_pandas() 函数将 Pandas DataFrame 分割成指定数量的分区,并创建一个 Dask DataFrame。

8.7.3 Dask DataFrame 的操作

Dask DataFrame 提供了与 Pandas DataFrame 相似的 API,可以执行各种数据操作,例如:

  • 数据选择: 使用 lociloc 进行数据选择。

  • 数据过滤: 使用布尔索引进行数据过滤。

  • 数据聚合: 使用 groupby() 进行数据聚合。

  • 数据转换: 使用 map_partitions()apply() 进行数据转换。

1. 数据过滤和选择:

import dask.dataframe as dd import pandas as pd # 创建一个示例 CSV 文件 data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'B', 'C', 'D', 'E', 'F', 'G', 'H', 'I', 'J']} df = pd.DataFrame(data) df.to_csv('example.csv', index=False) # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv('example.csv') # 过滤数据 ddf_filtered = ddf[ddf['col1'] > 5] # 选择列 ddf_selected = ddf_filtered[['col1']] # 查看结果 print(ddf_selected.head())

2. 数据聚合:

import dask.dataframe as dd import pandas as pd # 创建一个示例 CSV 文件 data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'A', 'B', 'B', 'C', 'C', 'A', 'A', 'B', 'B']} df = pd.DataFrame(data) df.to_csv('example.csv', index=False) # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv('example.csv') # 分组聚合 ddf_grouped = ddf.groupby('col2')['col1'].sum() # 查看结果 print(ddf_grouped.compute()) # 使用 compute() 触发计算

重要提示: Dask DataFrame 的操作是延迟执行的。这意味着当你执行一个操作时,Dask 并不会立即执行计算,而是将该操作添加到任务图中。只有当你调用 compute() 方法时,Dask 才会真正执行计算并返回结果。

3. 自定义函数与 map_partitions

map_partitions 允许你将自定义函数应用到 Dask DataFrame 的每个分区。 这对于执行复杂的数据转换非常有用。

import dask.dataframe as dd import pandas as pd # 创建一个示例 CSV 文件 data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'A', 'B', 'B', 'C', 'C', 'A', 'A', 'B', 'B']} df = pd.DataFrame(data) df.to_csv('example.csv', index=False) # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv('example.csv') def my_function(partition): # 自定义函数:将 col1 的值乘以 2 partition['col1_multiplied'] = partition['col1'] * 2 return partition ddf_transformed = ddf.map_partitions(my_function) print(ddf_transformed.head())

8.7.4 Dask 的调度器

Dask 提供了多种调度器,用于执行任务图。常见的调度器包括:

  • 单线程调度器: 在单个线程中执行任务,适用于调试和小型数据集。

  • 线程池调度器: 使用线程池并行执行任务,适用于 CPU 密集型任务。

  • 进程池调度器: 使用进程池并行执行任务,适用于 CPU 密集型任务和避免 GIL 限制。

  • 分布式调度器: 在分布式集群上执行任务,适用于大规模数据集和复杂计算。

可以通过 dask.config.set() 函数来配置 Dask 的调度器。例如,使用线程池调度器:

import dask dask.config.set(scheduler='threads')

选择合适的调度器取决于你的计算需求和可用资源。对于 CPU 密集型任务,进程池调度器通常比线程池调度器更有效,因为它可以避免 GIL (Global Interpreter Lock) 的限制。对于大规模数据集和复杂计算,分布式调度器是最佳选择。

8.7.5 Dask 任务图可视化

Dask 允许你可视化任务图,这有助于理解计算流程和诊断性能问题。可以使用 visualize() 方法来生成任务图。

import dask.dataframe as dd import pandas as pd # 创建一个示例 CSV 文件 data = {'col1': [1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 'col2': ['A', 'A', 'B', 'B', 'C', 'C', 'A', 'A', 'B', 'B']} df = pd.DataFrame(data) df.to_csv('example.csv', index=False) # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv('example.csv') # 分组聚合 ddf_grouped = ddf.groupby('col2')['col1'].sum() # 可视化任务图 ddf_grouped.visualize(filename='dask_graph.png')

这段代码会生成一个名为 dask_graph.png 的图像文件,其中包含了任务图的可视化表示。

任务图示例 (mermaid graph TD):

8.7.6 代码实践:处理大型 CSV 文件

假设有一个大型 CSV 文件 (超过内存大小),需要计算每个类别的平均值。

import dask.dataframe as dd import pandas as pd import numpy as np import os # 模拟生成一个大型 CSV 文件 file_path = 'large_data.csv' chunk_size = 1000000 num_chunks = 10 if not os.path.exists(file_path): for i in range(num_chunks): data = {'category': np.random.choice(['A', 'B', 'C'], size=chunk_size), 'value': np.random.rand(chunk_size)} df = pd.DataFrame(data) if i == 0: df.to_csv(file_path, mode='w', header=True, index=False) else: df.to_csv(file_path, mode='a', header=False, index=False) print(f"Large CSV file '{file_path}' generated.") else: print(f"Large CSV file '{file_path}' already exists.") # 使用 Dask 读取 CSV 文件 ddf = dd.read_csv(file_path) # 计算每个类别的平均值 ddf_grouped = ddf.groupby('category')['value'].mean() # 执行计算 result = ddf_grouped.compute() # 打印结果 print(result)

这段代码首先模拟生成一个大型 CSV 文件。然后,使用 Dask 读取 CSV 文件,并使用 groupby()mean() 函数计算每个类别的平均值。最后,调用 compute() 方法执行计算并打印结果。Dask 会自动将数据分割成多个分区,并并行地在这些分区上执行计算,从而有效地处理大型数据集。

8.7.7 总结

Dask 是一个强大的并行计算库,它可以与 Pandas 无缝集成,扩展 Pandas 的处理能力,使其能够处理超出内存的数据集。通过使用 Dask DataFrame,你可以利用多核 CPU 和分布式计算资源,加速数据处理任务。掌握 Dask 的基本概念和使用方法,对于处理大型数据集和提高数据分析效率至关重要。Dask 提供的延迟执行、内存管理和任务图可视化等功能,可以帮助你更好地理解和优化你的数据处理流程。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U