使用 Dask 进行并行计算
Dask 是一个灵活的开源 Python 库,用于并行计算。在本文中,我们将了解并行计算以及为什么应该选择 Dask 来实现并行计算。
我们将把它与 Spark、Ray 和 Modin 等其他库进行比较。我们还讨论了 Dask 的用例。
并行计算
并行计算是一种同时执行多个计算或进程的计算类型。大型问题通常被分解成可管理的部分,以便分别解决。
并行计算分为四类:
位级
指令级
数据级
作业并行。
尽管并行性在高性能计算中应用已久,但由于频率扩展的物理限制,它直到最近才变得更加流行。
Dask 的必要性
我想到的一个问题是,我们为什么需要 Dask?
借助 Numpy、Sklearn、Seaborn 等 Python 库,数据操作和机器学习任务变得简单。对于大多数数据分析任务,Python [pandas] 模块就足够了。数据可以通过多种不同的方式进行操作,并且可以使用这些数据创建机器学习模型。
但是,如果您的数据量超过可用的 RAM,Pandas 将变得不够用。这是一个相当常见的问题。您可以使用 Spark 或 Hadoop 来解决这个问题。然而,这些都不是 Python 环境。因此,您将无法使用 NumPy、Pandas、Sklearn、TensorFlow 和其他常用的 Python 机器学习工具。有办法解决这个问题吗?有!这时 Dask 就派上用场了。
Dask 简介
Dask 是一个并行计算框架,可与 Jupyter Notebook 无缝集成。最初,它是为了扩展 NumPy、Pandas 和 Scit-kit 的计算能力而创建的,以突破单机的存储限制。DASK 的类似物可用于学习,但很快它就被用作通用分布式系统。
Dask 有两个主要优势 -
可扩展性
Dask 可以原生扩展 Pandas、NumPy 和 Scikit-Learn 的 Python 版本,并在多核集群上弹性运行。它也可以缩减规模以在单个系统上运行。
规划
与 Airflow 和 Luigi 类似,Dask 任务调度器针对计算进行了优化。它提供快速反馈,使用任务图管理任务,并支持本地和分布式诊断,使其动态且响应迅速。
此外,Dask 还提供实时动态仪表板,每 100 毫秒更新一次,并显示进度、内存利用率等各种信息。
您可以根据自己的喜好,克隆 git 仓库或使用 Conda/pip 安装 Dask。
conda install dask
仅安装核心组件 -
conda install dask-core
Dask-core 是 Dask 的受限版本,仅安装必要组件。 pip 也是如此。如果您只想使用 dask 数据框和 dask 数组来扩展 pandas、numpy 或两者,那么您也可以只安装 dask 数据框或 dask 数组。
python -m pip install dask
安装数据框的要求
python -m pip install "dask[dataframe]" #
安装阵列的要求
python -m pip install "dask[list]"
让我们看一下该库用于并行计算的几个实例。我们的代码使用 dask.delayed 来实现并行性。
注意 - 以下两个代码片段应在 Jupyter Notebook 的两个不同单元中运行
import time import random def calcprofit(a, b): time.sleep(random.random()) return a + b def calcloss(a, b): time.sleep(random.random()) return a - b def calctotal(a, b): time.sleep(random.random()) return a + b
现在运行下面的代码片段 -
%%time profit = calcprofit(10, 22) loss = calcloss(18, 3) total = calctotal(profit, loss) print(total)
输出
47 CPU times: user 4.13 ms, sys: 1.23 ms, total: 5.36 ms Wall time: 1.35 s
虽然这些函数彼此独立,但它们将按顺序依次执行。因此,我们可以并发执行它们以节省时间。
import dask calcprofit = dask.delayed(calcprofit) calcloss = dask.delayed(calcloss) calctotal = dask.delayed(calctotal)
现在运行下面的代码片段 -
%%time profit = calcprofit(10, 22) loss = calcloss(18, 3) total = calctotal(profit, loss) print(total)
输出
Delayed('calctotal-9e3e896e-b4de-400c-aeb8-9e4c0961fe11')
CPU times: user 3.3 ms, sys: 0 ns, total: 3.3 ms
Wall time: 10.2 ms
即使在这个简单的示例中,运行时间也得到了改善。我们还可以看到如下任务图:
total.visualize(rankdir='LR')
Spark vs. Dask
Spark 是一个强大的集群计算框架工具,它将数据和处理划分为可管理的部分,并将它们分布到任意规模的集群上,并并发执行。
尽管 Spark 是事实上的大数据分析标准技术,但 Dask 似乎前景光明。Dask 轻量级且作为 Python 组件开发,而 Spark 则拥有额外的功能,主要使用 Scala 开发,并且还支持 Python/R。如果您想要一个切实可行的解决方案,甚至拥有 JVM 基础架构,Spark 可能是您的首选。但是,如果您想要快速、轻量级的并行处理,Dask 也是一个可行的选择。只需快速安装 pip 即可使用。
Dask、Ray 和 Modin
Ray 和 Dask 有不同的调度策略。集群的所有作业都由 Dask 使用的中央调度程序管理。由于 Ray 是去中心化的,每台计算机都有自己的调度程序,这使得计划任务的问题可以在特定计算机级别而不是整个集群级别解决。Ray 缺乏 Dask 提供的丰富的高级集合 API(例如数据帧、分布式数组等)。
相反,Modin 则超越了 Dask 或 Ray。只需添加一行代码 import modin.pandas as pd,我们就可以使用 Modin 快速扩展我们的 Pandas 进程。尽管 Modin 尽力将 Pandas API 的大部分功能并行化,但 Dask DataFrame 有时无法扩展完整的 Pandas API。
Dask 用例示例
Dask 应用案例分为两类 -
我们可以利用动态任务调度来优化计算。
大型数据集可以使用"大数据"集合(例如并行数组和数据帧)来处理。
任务图是我们数据处理作业组织的可视化描述,它是使用 Dask 集合制作的。
任务图使用 Dask 调度程序执行。
Dask 使用并行编程完成作业。
"并行编程"一词指的是同时执行多个任务。
这样,我们可以高效地利用资源,同时完成多项任务。
让我们来看看 Dask 提供的一些数据集。
Dask.array - 使用 Numpy 接口,dask.array 将巨大的数组划分为更小的数组,使我们能够对大于系统内存的数组进行计算。
Dask.bag - 它提供对标准 Python 对象集合的操作,例如 filter、map、group by 和 fold。
Dask.dataframe - 像 Pandas 一样分发数据帧。它是一个由多个微小数据帧构成的巨大的并行数据帧。
结论
在本文中,我们了解了 Dask 和并行计算。我们帮助您增强了对 Dask 的了解,包括它的需求以及它与其他库的比较。

