Press "Enter" to skip to content

通过Spark在Python上实现并行化:Pandas并发选项

在使用Pandas时充分发挥Spark的优势

由Florian Steciuk在Unsplash上的照片

在我的上一个岗位上,我花了一些时间开展了一个内部项目,以预测我们托管服务客户数千个磁盘的未来磁盘存储空间使用情况。每个磁盘都有自己的使用模式,这意味着我们需要为每个磁盘创建一个单独的机器学习模型,该模型利用历史数据来预测磁盘的未来使用情况。尽管进行这种预测并选择正确的算法本身就充满挑战,但在大规模进行此操作也有其自身的问题。

为了利用更复杂的基础架构,我们可以考虑摆脱按顺序进行预测的方式,通过将工作负载并行化来加速预测的操作。本篇博文旨在比较Pandas UDF和’concurrent.futures’模块这两种并发处理方法,并确定各自的使用案例。

挑战

Pandas是Python中用于处理分析领域数据集的门户包。通过使用DataFrame,我们可以对数据进行分析和评估数据质量,进行探索性数据分析,构建对数据的描述性可视化以及预测未来趋势。

虽然这无疑是一个很好的工具,但Python的单线程性质意味着在处理较大数据集时,它的伸缩性可能较差,或者当您需要对多个数据子集执行相同分析时。

在大数据领域,我们期望在我们的方法中具备更多复杂性,因为我们还需要关注可伸缩性以保持良好的性能。Spark等语言使我们能够利用分布式处理来处理更大更复杂的数据结构。

在深入研究此特定示例之前,我们可以总结出一些需要并发处理数据的使用案例:

  • 对多个数据文件应用统一转换
  • 预测几个数据子集的未来值
  • 调整机器学习模型的超参数,选择最有效的配置

当我们的需求升级到像上述建议的工作负载时,在我们的案例中,Python和Pandas中最简单的方法是按顺序处理这些数据。对于我们的示例,我们将逐个磁盘运行上述流程。

数据

在我们的示例中,我们有数千个磁盘的数据,显示了随时间记录的可用空间,并且我们希望预测每个磁盘的未来可用空间值。

为了更清楚地描述,我提供了一个包含1,000个磁盘的csv文件,每个磁盘有一个月的历史空闲空间数据,以GB为单位。这足够大,让我们看到不同预测方法在大规模预测中的影响。

作者:示例DataFrame图像

对于这样的时间序列问题,我们希望使用历史数据来预测未来趋势,并且我们想知道哪种机器学习(ML)算法对于每个磁盘来说最合适。在确定适合一个数据集的适当模型时,AutoML等工具非常好用,但我们这里有1,000个数据集,所以这对我们的比较来说是过度的。

在这种情况下,我们将把要比较的算法数量限制为两个,并使用均方根误差(RMSE)作为验证指标,看哪个是最适合使用的模型。关于RMSE的更多信息可以在这里找到。这些算法包括:

  • 线性回归
  • Fbprophet(将数据拟合到更复杂的线条上)
  • Facebook的时间序列预测模型。
  • 用于带有季节性超参数的更复杂预测的模型。

如果我们想预测单个磁盘的未来空闲空间,我们现在已经准备好了所有组件。操作的步骤如下:

Image by Author: Data Lifecycle

现在我们想要进行扩展,对多个盘进行此流程,例如我们的示例中有 1,000 个盘。

作为我们的评估的一部分,我们将比较不同规模下使用不同算法计算 RMSE 值的性能。因此,我创建了一个包含前 100 个盘的子集来模拟。

这将为我们提供一些关于不同大小数据集的性能以及执行不同复杂度操作的有趣见解。

引入并发性

Python 以单线程而闻名,因此无法同时利用所有可用的计算资源。

因此,我看到了三个选择:

  1. 使用 for 循环按顺序计算预测值,采用单线程方法。
  2. 使用 Python 的futures 模块同时运行多个进程。
  3. 使用 Pandas UDF(用户自定义函数)在 PySpark 中利用分布式计算,同时保持我们的 Pandas 语法和兼容包。

我想在不同环境条件下进行一次相当深入的比较,所以我使用了一个单节点的 Databricks 集群以及另一个带有 4 个工作节点的 Databricks 集群,以利用 Spark 来实现我们的 Pandas UDF 方法。

我们将按照以下方法评估线性回归和 fbprophet 模型在每个盘上的适用情况:

  • 将数据拆分为训练集和测试集
  • 使用训练集作为输入,在测试集日期上进行预测
  • 将预测值与测试集中的实际值进行比较,得到均方根误差 (RMSE) 分数

我们将在输出中返回两个内容:一个修改后的 DataFrame,其中包含预测结果,使我们能够绘制和比较预测与实际值,以及一个包含每个盘和算法的 RMSE 分数的 DataFrame。

下面是执行此操作的函数:

我们将比较上面提到的三种方法。我们有一些不同的情景,因此我们可以填写一个根据我们收集结果的表格:

使用以下组合:

方法

  • 顺序执行
  • futures
  • Pandas UDFs

算法

  • 线性回归
  • Fbprophet
  • 结合(每个盘使用两种算法) – 收集比较的最有效方法。

集群模式

  • 单节点集群
  • 具有 4 个工作节点的标准集群

盘的数量

  • 100
  • 1,000

结果以附录的格式呈现在此博客中,如果您希望进一步查看。

方法

方法 1:顺序执行

方法 2:concurrent.futures

使用此模块有两个选项:并行化内存密集型操作(使用 ThreadPoolExecutor)或 CPU 密集型操作(ProcessPoolExecutor)。如下博客中有关于此的描述。由于我们将处理 CPU 密集型问题,因此 ProcessPoolExecutor 是适合我们目标的选择。

方法 3:Pandas UDFs

现在我们将转向使用 Spark,并利用分布式计算来提高效率。由于我们使用的是 Databricks,所以大部分 Spark 配置已经为我们完成,但是我们对数据的一般处理还有一些调整。

首先,将数据导入到一个 PySpark DataFrame 中:

我们将使用 Pandas 分组映射 UDF(PandasUDFType.GROUPED_MAP),因为我们想要传入一个 DataFrame 并返回一个 DataFrame。由于Apache Spark 3.0之后我们不再需要显式声明此装饰器!

由于在PySpark中的DataFrame结构化,我们需要将fbprophet、回归和RMSE函数分离为Pandas UDFs,但不需要进行大量代码改动即可实现。

然后我们可以使用applyInPandas来生成我们的结果。

注意:上面的示例仅演示了使用线性回归的过程,以提高可读性。请参阅完整的notebook以了解完整的演示。

解释结果

通过Spark在Python上实现并行化:Pandas并发选项 四海 第4张

通过Spark在Python上实现并行化:Pandas并发选项 四海 第5张

通过Spark在Python上实现并行化:Pandas并发选项 四海 第6张

我们已经为不同的方法和不同的环境设置创建了图表,然后按算法和磁盘数量进行了分组,以便进行简单比较。

请注意,表格结果在本文的附录中。

下面是这些发现的要点:

  • 预测1000个磁盘与预测100个磁盘相比,通常需要更长的时间。
  • 顺序方法通常最慢,无法有效利用底层资源。
  • Pandas UDFs在较小、较简单的任务上效率不高。数据转换的开销更大,通过并行化来弥补这一点。
  • 顺序方法和concurrent.futures方法都没有利用Databricks中可用的集群计算资源,错过了额外的计算能力。

结束语

在确定最成功的方法时,上下文肯定起着重要作用,但考虑到Databricks和Spark通常用于大数据问题,我们可以看到在这里使用Pandas UDFs与这些更大、更复杂的数据集的好处。

在较小的数据集上使用Spark环境时,可以采用较小(且更便宜!)的计算配置来实现同样的高效率,正如使用concurrent.futures模块所示,因此在架构解决方案时要记住这一点。

如果您熟悉Python和Pandas,那么从顺序for循环方法转向并行方法应该不是一项艰巨的学习曲线,这在初学者教程中已经有所展示。

我们在这篇文章中没有进行详细探索,因为我发现当前版本存在差异和不兼容性,但最近的pyspark.pandas模块在未来肯定会更常见,也是一个值得关注的方法。这个API(连同由Databricks团队开发,但现已被弃用的Koalas)将Pandas的熟悉性与Spark的潜在好处结合起来。

为了演示我们想要实现的效果,我们只关注了每个磁盘生成的RMSE值,而没有真正预测未来时间序列的一组值。我们在这里设置的框架可以以同样的方式应用于此,通过逻辑来确定是否适用评估指标(以及磁盘的其他物理限制等)并使用所确定的算法尽可能地预测未来值。

如常,该notebook可以在我的GitHub中找到。

附录

最初发表于https://blog.coeo.com,适应本次转载。

Leave a Reply

Your email address will not be published. Required fields are marked *