原文:zh.annas-archive.org/md5/7b21344f8b9a93e8b4b323b6ffeec37b

译者:飞龙

协议:CC BY-NC-SA 4.0

前言

计算机行业的特点是寻求不断提高和高效性能,从网络、电信、航空电子领域的高端应用,到桌面计算机、笔记本电脑和视频游戏中的低功耗嵌入式系统。这种发展路径导致了多核系统,其中双核、四核和八核处理器只是未来向不断增加的计算核心数量扩展的起点。

然而,这种扩展不仅对半导体行业,也对能够进行并行计算的应用程序的开发构成了挑战。

并行计算实际上代表了同时使用多个计算资源来解决处理问题,这样它就可以在多个 CPU 上执行,将问题分解成可以同时处理的离散部分,其中每个部分进一步分解成可以在不同的 CPU 上串行执行的指令序列。

计算资源可以包括具有多个处理器的单个计算机、通过网络连接的任意数量的计算机,或者这两种方法的组合。并行计算一直被认为是计算技术的极致或未来,直到几年前,它还受到复杂系统和各种领域相关情况的数值模拟的推动:天气和气候预测、化学反应和核反应、人类基因组图谱、地震和地质活动、机械装置的行为(从假肢到航天飞机)、电子电路和制造过程。

现在,然而,越来越多的商业应用正日益要求开发速度更快的计算机来以复杂的方式处理大量数据。这类应用包括数据挖掘、并行数据库、石油勘探、网络搜索引擎、网络化商业服务、计算机辅助医疗诊断、跨国公司的管理、高级图形和虚拟现实(尤其是在视频游戏行业)、多媒体和视频网络技术,以及协作工作环境。

最后但同样重要的是,并行计算代表了尝试最大化那种无限但同时也越来越宝贵和稀缺的资源——时间。这就是为什么并行计算正从只为少数人保留的昂贵超级计算机的世界转向基于多个处理器、图形处理单元GPUs)或少量互联计算机的经济解决方案,这些计算机可以克服串行计算的约束和单个 CPU 的限制。

为了介绍并行编程的概念,选择了最受欢迎的编程语言之一——Python。Python 的流行部分归因于其灵活性,因为它是一种被网页和桌面开发者、系统管理员和代码开发者、以及最近的数据科学家和机器学习工程师经常使用的语言。

从技术角度来看,在 Python 中,没有像 C 语言那样的独立编译阶段(例如,从源代码生成可执行文件)。由于其伪解释特性,Python 成为了一种可移植的语言。一旦编写了源代码,它就可以在大多数当前使用的平台上进行解释和执行,无论是来自苹果(macOS X)还是 PC(Microsoft Windows 和 GNU/Linux)。

Python 的另一个优点是其易于学习。任何人都可以在几天内学会使用它并编写他们的第一个应用程序。在这种情况下,语言的开放结构起着基本的作用,没有冗余的声明,因此与口语非常相似。最后,Python 是免费软件:不仅 Python 解释器和 Python 在我们应用程序中的使用是免费的,而且 Python 也可以根据完全开源许可证的规则自由修改和重新分发。

《Python 并行编程食谱,第二版》包含大量示例,为读者提供了解决实际问题的机会。它考察了并行架构的软件设计原则,强调程序清晰性的重要性,并避免使用复杂术语,而是使用清晰直接示例。

每个主题都作为完整的、可工作的 Python 程序的一部分进行介绍,始终跟随相关程序的输出。各个章节的模块化组织提供了一个经过验证的路径,沿着这个路径可以从最简单的论点过渡到最先进的论点,但它也适合只想学习几个特定问题的人。

本书面向对象

《Python 并行编程食谱,第二版》旨在为希望使用并行编程技术编写强大和高效代码的软件开发人员编写。阅读本书将使您掌握并行计算的基本和高级方面。

Python 编程语言易于使用,并允许非专业人士轻松应对并理解本书概述的主题。

本书涵盖内容

第一章,Python 并行计算入门,概述了并行编程架构和编程模型。本章介绍了 Python 编程语言,讨论了该语言的特点、易学易用性、可扩展性以及丰富的软件库和应用程序,所有这些都使 Python 成为任何应用程序,尤其是并行计算的有价值工具。

第二章,基于线程的并行性,讨论了使用threading Python 模块的线程并行性。读者将通过完整的编程示例学习如何同步和操作线程以实现他们的多线程应用程序。

第三章,基于过程的并行性,引导读者了解如何通过基于过程的途径并行化一个程序。一系列完整的示例将展示读者如何使用multiprocessing Python 模块。

第四章,消息传递,专注于消息传递交换通信系统。特别是,将用大量的应用示例来描述mpi4py库。

第五章,异步编程,解释了用于并发编程的异步模型。在某些方面,它比线程模型更简单,因为有一个单一的指令流,任务明确地放弃控制权而不是任意挂起。本章展示了读者如何使用asyncyio模块将每个任务组织成一系列必须以异步方式执行的小步骤。

第六章,分布式 Python,向读者介绍了分布式计算,这是将多个计算单元聚合起来以透明和一致的方式协同运行单个计算任务的过程。特别是,本章提供的示例应用程序描述了使用socket和 Celery 模块来管理分布式任务。

第七章,云计算,提供了与 Python 编程语言相关的主要云计算技术的概述。PythonAnywhere平台对于在云上部署 Python 应用程序非常有用,本章将对其进行考察。本章还包含示例应用程序,展示了容器无服务器技术的使用。

第八章,异构计算,探讨了现代 GPU 在增加编程复杂性的代价下为数值计算提供突破性性能。实际上,GPU 的编程模型要求程序员手动管理 CPU 和 GPU 之间的数据传输。本章将通过编程示例和用例,教读者如何利用PyCUDANumbaPyOpenCL强大的 Python 模块来发挥 GPU 卡的计算能力。

第九章,Python 调试和测试,是介绍软件工程中两个重要主题的最后一章:调试和测试。特别是,以下 Python 框架将被描述:winpdb-reborn用于调试,以及unittestnose用于软件测试。

要充分利用这本书

本书是自包含的:在开始阅读之前,唯一的基本要求是对编程的热情和对本书涵盖主题的好奇心。

下载示例代码文件

您可以从 www.packt.com 的账户下载本书的示例代码文件。如果您在其他地方购买了这本书,您可以访问 www.packtpub.com/support 并注册,以便将文件直接通过电子邮件发送给您。

您可以通过以下步骤下载代码文件:

  1. www.packtpub.com 登录或注册。

  2. 选择“支持”选项卡。

  3. 点击“代码下载”。

  4. 在搜索框中输入书籍名称,并遵循屏幕上的说明。

文件下载完成后,请确保您使用最新版本的以下软件解压或提取文件夹:

  • WinRAR/7-Zip for Windows

  • Zipeg/iZip/UnRarX for Mac

  • 7-Zip/PeaZip for Linux

书籍的代码包也托管在 GitHub 上,网址为 github.com/PacktPublishing/Python-Parallel-Programming-Cookbook-Second-Edition。我们还有其他来自我们丰富的图书和视频目录的代码包可供使用,网址为**github.com/PacktPublishing/**。查看它们吧!

下载彩色图像

我们还提供了一份包含本书中使用的截图/图表彩色图像的 PDF 文件。您可以从这里下载:static.packt-cdn.com/downloads/9781789533736_ColorImages.pdf

使用的约定

本书使用了多种文本约定。

CodeInText:表示文本中的代码单词、数据库表名、文件夹名、文件名、文件扩展名、路径名、虚拟 URL、用户输入和 Twitter 昵称。以下是一个示例:“可以通过使用terminate方法立即终止进程。”

代码块设置如下:

import socket
port=60000
s =socket.socket()
host=socket.gethostname()

当我们希望您注意代码块中的特定部分时,相关的行或项目将以粗体显示:

 p = multiprocessing.Process(target=foo)
 print ('Process before execution:', p, p.is_alive())
 p.start()

任何命令行输入或输出都应如下所示:

> python server.py

粗体: 表示新术语、重要单词或屏幕上看到的单词。例如,菜单或对话框中的单词在文本中显示如下。以下是一个示例:“转到系统属性 | 环境变量 | 用户或系统变量 | 新建。”

警告或重要注意事项看起来像这样。

小贴士和技巧看起来像这样。

部分

在这本书中,您将找到几个频繁出现的标题(准备工作如何做…它是如何工作的…更多内容…参考以下内容)。

要清楚地说明如何完成食谱,请按以下方式使用这些部分:

准备工作

本节告诉您在食谱中可以期待什么,并描述如何设置任何软件或任何为食谱所需的初步设置。

如何做…

本节包含遵循食谱所需的步骤。

它是如何工作的…

本节通常包含对上一节发生情况的详细解释。

更多内容…

本节包含有关食谱的附加信息,以便您对食谱有更深入的了解。

参考以下内容

本节提供了对食谱有用的其他信息的链接。

联系我们

我们欢迎读者的反馈。

一般反馈: 发送电子邮件至feedback@packtpub.com,并在邮件主题中提及书籍标题。如果您对本书的任何方面有疑问,请通过questions@packtpub.com发送电子邮件给我们。

勘误: 尽管我们已经尽最大努力确保内容的准确性,但错误仍然可能发生。如果您在这本书中发现了错误,如果您能向我们报告,我们将不胜感激。请访问www.packtpub.com/support/errata,选择您的书籍,点击勘误提交表单链接,并输入详细信息。

盗版: 如果您在互联网上以任何形式发现我们作品的非法副本,如果您能向我们提供位置地址或网站名称,我们将不胜感激。请通过copyright@packtpub.com与我们联系,并提供材料的链接。

如果您有兴趣成为作者: 如果您在某个主题上具有专业知识,并且您有兴趣撰写或为书籍做出贡献,请访问authors.packtpub.com

评论

请留下评论。一旦您阅读并使用过这本书,为何不在您购买它的网站上留下评论呢?潜在读者可以查看并使用您的客观意见来做出购买决定,Packt 公司可以了解您对我们产品的看法,我们的作者也可以看到他们对书籍的反馈。谢谢!

如需了解 Packt 的更多信息,请访问packtpub.com.

第一章:开始使用并行计算和 Python

并行分布式计算模型基于同时使用不同的处理单元来执行程序。尽管并行计算和分布式计算之间的区别非常微妙,其中一种可能的定义将并行计算模型与共享内存计算模型相关联,将分布式计算模型与消息传递模型相关联。

从现在开始,我们将使用术语并行计算来指代并行和分布式计算模型。

以下几节提供了并行编程架构和编程模型的概述。这些概念对于第一次接触并行编程技术的经验不足的程序员来说很有用。此外,它也可以作为经验丰富的程序员的基本参考。并行系统的双重特征也被介绍。第一种特征基于系统架构,而第二种特征基于并行编程范式。

本章以对 Python 编程语言的简要介绍结束。该语言的特点、易用性和学习性,以及软件库和应用的扩展性和丰富性使 Python 成为任何应用的宝贵工具,也是并行计算的宝贵工具。线程和进程的概念在它们在语言中的应用中被介绍。

在本章中,我们将涵盖以下食谱:

  • 我们为什么需要并行计算?

  • 飞利浦分类法

  • 内存组织

  • 并行编程模型

  • 评估性能

  • 介绍 Python

  • Python 和并行编程

  • 介绍进程和线程

我们为什么需要并行计算?

现代计算机提供的计算能力的增长导致我们在相对较短的时间内面临越来越复杂的计算问题。直到 2000 年代初,复杂性是通过增加晶体管的数量以及单处理器系统的时钟频率来处理的,这些系统的峰值达到了 3.5-4 GHz。然而,晶体管数量的增加导致了处理器本身功耗的指数级增长。本质上,因此存在一个物理限制,阻止了单处理器系统性能的进一步改进。

因此,在近年来,微处理器制造商将注意力集中在多核系统上。这些系统基于几个物理处理器共享相同内存的核心,从而绕过了之前描述的功耗问题。近年来,四核八核系统也已成为普通桌面和笔记本电脑配置的标准。

另一方面,这种重大的硬件变化也导致了软件结构的演变,软件结构一直是设计为在单个处理器上顺序执行的。为了利用增加处理器数量所提供的更多计算资源,现有的软件必须重新设计成适合 CPU 并行结构的适当形式,以便通过同时执行同一程序的多个部分的单个单元来获得更高的效率。

飞利浦的分类法

飞利浦的分类法是一个用于分类计算机架构的系统。它基于两个主要概念:

  • 指令流:一个拥有 n 个 CPU 的系统有 n 个程序计数器,因此有 n 个指令流。这对应于一个程序计数器。

  • 数据流:一个计算数据列表上函数的程序有一个数据流。计算多个不同数据列表上相同函数的程序有更多的数据流。这由一组操作数组成。

由于指令和数据流是独立的,因此有四种并行机类别:单指令单数据SISD)、单指令多数据SIMD)、多指令单数据MISD)和多指令多数据MIMD):

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/315c390f-4c31-4a69-811e-5696dff064d1.png

飞利浦的分类法

单指令单数据(SISD)

SISD 计算系统类似于冯·诺伊曼机器,这是一种单处理器机器。正如你在 Flynn’s taxonomy 图中可以看到的,它执行一个操作单个数据流的指令。在 SISD 中,机器指令是顺序处理的。

在一个时钟周期内,CPU 执行以下操作:

  • 取指令:CPU 从一个内存区域取数据和指令,这个区域称为寄存器

  • 解码:CPU 解码指令。

  • 执行:指令在数据上执行。操作的结果存储在另一个寄存器中。

一旦执行阶段完成,CPU 将自己设置为开始另一个 CPU 周期:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/8d131bb3-cd52-4f51-969c-8c836a394f89.png

取指令、解码和执行周期

在这类计算机上运行的算法是顺序的(或串行的),因为它们不包含任何并行性。SISD 计算机的一个例子是具有单个 CPU 的硬件系统。

这些架构的主要元素(即冯·诺伊曼架构)如下:

  • 中央存储单元:用于存储指令和程序数据。

  • CPU:这是用来从内存单元获取指令和/或数据,它解码指令并依次执行它们。

  • I/O 系统:这指的是程序的输入和输出数据。

传统的单处理器计算机被归类为 SISD 系统:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/1a6b6929-c4c4-41bb-96ae-2ec6b2aea55d.png

SISD 架构方案

下图具体显示了 CPU 在取指、解码和执行阶段所使用的区域:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/9ceec645-0aa1-4c90-97a7-cba9f8eb5031.png

取指-解码-执行阶段的 CPU 组件

多指令单数据(MISD)

在这个模型中,n个处理器,每个都有自己的控制单元,共享一个单一的内存单元。在每个时钟周期,从内存接收到的数据由所有处理器同时处理,每个处理器根据其控制单元接收到的指令进行处理。

在这种情况下,通过在相同的数据上执行多个操作来获得并行性(指令级并行性)。在这些架构中可以有效地解决的问题类型相当特殊,例如数据加密。因此,MISD 计算机在商业领域没有找到空间。MISD 计算机更多的是一种智力练习,而不是实际配置。

单指令多数据(SIMD)

一个 SIMD 计算机由n个相同的处理器组成,每个处理器都有自己的局部内存,可以存储数据。所有处理器在单个指令流的控制下工作。此外,还有n个数据流,每个处理器一个。处理器在每一步同时工作,执行相同的指令,但针对不同的数据元素。这是一个数据级并行的例子。

SIMD 架构比 MISD 架构更灵活。SIMD 计算机上的并行算法可以解决广泛应用的众多问题。另一个有趣的特点是,这些计算机的算法相对容易设计、分析和实现。局限性在于,只有那些可以分解为多个子问题(这些子问题都是相同的,每个子问题将通过同一组指令同时解决)的问题才能用 SIMD 计算机解决。

根据这种范例开发的超级计算机,我们必须提到连接机(思考机器,1985)和MPP(NASA,1983)。

正如我们将在第六章“分布式 Python”和第七章“云计算”中看到的,现代图形卡(GPU)的出现,这些图形卡由许多 SIMD 嵌入式单元构建,导致了这种计算范例的更广泛使用。

多指令多数据(MIMD)

根据弗林分类,这类并行计算机是最通用和最强大的,它包含n个处理器、n个指令流和n个数据流。每个处理器都有自己的控制单元和局部内存,这使得 MIMD 架构比 SIMD 架构在计算上更强大。

每个处理器在其控制单元发出的指令流的控制下运行。因此,处理器可以潜在地运行不同的程序,使用不同的数据,这使得它们可以解决不同且可以是单个更大问题一部分的子问题。在 MIMD 中,通过线程和/或进程的并行级别来实现架构。这也意味着处理器通常以异步方式运行。

现在,这种架构已应用于许多个人电脑、超级计算机和计算机网络。然而,你需要考虑的一个反例是:异步算法难以设计、分析和实现:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/f3fa5a98-1a89-4d76-a8af-c25a376c1ac8.png

SIMD 架构(A)和 MIMD 架构(B)

通过考虑 SIMD 机器可以分为两个子组,Flynn 的分类法可以扩展:

  • 数值超级计算机

  • 向量机

另一方面,MIMD 可以分为具有共享内存的机器和具有分布式内存的机器。

事实上,下一节将重点介绍 MIMD 机器内存组织的最后一个方面。

内存组织

为了评估并行架构,我们需要考虑的另一个方面是内存组织,或者说数据访问的方式。无论处理单元有多快,如果内存不能以足够的速度维护和提供指令和数据,那么性能就不会有所提高。

我们需要克服的主要问题,使内存的响应时间与处理器的速度相匹配,是内存周期时间,它定义为两次连续操作之间经过的时间。处理器的周期时间通常比内存的周期时间短得多。

当处理器启动对内存的传输时,处理器的资源将在整个内存周期内保持占用;此外,在此期间,由于正在进行的传输,没有任何其他设备(例如,I/O 控制器、处理器或甚至请求该资源的处理器)能够使用内存:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/d56ae362-cff1-4d47-97aa-78419437b2b7.png

MIMD 架构中的内存组织

解决内存访问问题的解决方案导致了 MIMD 架构的二分法。第一种系统,称为共享内存系统,具有高虚拟内存,所有处理器都可以平等地访问该内存中的数据和指令。另一种系统是分布式内存模型,其中每个处理器都有本地内存,其他处理器无法访问。

分布式内存共享内存的特点是内存访问的管理,这由处理单元执行;这种区别对于程序员来说非常重要,因为它决定了并行程序的不同部分必须如何通信。

特别是,分布式内存机器必须在每个本地内存中复制共享数据。这些副本是通过从一个处理器向另一个处理器发送包含要共享的数据的消息来创建的。这种内存组织的缺点是,有时这些消息可以非常大,并且需要相对较长的时间来传输,而在共享内存系统中,没有消息交换,主要问题在于同步访问共享资源。

共享内存

下图显示了共享内存多处理器系统的架构。这里的物理连接相当简单:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/0f53e868-04c0-4493-9b33-f9be28089ca2.png

共享内存架构图

在这里,总线结构允许任意数量的设备(如图中所示的前一个图中CPU + Cache)共享相同的通道(主内存,如图中所示的前一个图)。总线协议最初是为了允许单个处理器和一台或多台磁盘或磁带控制器通过这里的共享内存进行通信而设计的。

每个处理器都关联了缓存内存,因为假设处理器需要将数据或指令保留在本地内存中的概率非常高。

当处理器修改其他处理器同时使用的内存系统中存储的数据时,问题就出现了。新的值将从已更改的处理器缓存传递到共享内存。然而,稍后它还必须传递给所有其他处理器,以便它们不与过时的值一起工作。这个问题被称为缓存一致性问题——这是内存一致性问题的特例,它需要能够处理并发问题和同步的硬件实现,类似于线程编程。

共享内存系统的主要特点如下:

  • 所有处理器的内存都是相同的。例如,所有与相同数据结构相关的处理器都将使用相同的逻辑内存地址,从而访问相同的内存位置。

  • 通过读取各种处理器的任务并允许共享内存来获得同步。实际上,处理器一次只能访问一个内存。

  • 在另一个任务访问它时,共享内存位置不得从任务中更改。

  • 在任务之间共享数据非常快。所需通信的时间是其中一个任务读取单个位置所需的时间(取决于内存访问的速度)。

共享内存系统中的内存访问如下:

  • 统一内存访问UMA):该系统的基本特征是每个处理器和任何内存区域的访问时间都是恒定的。因此,这些系统也被称为对称多处理器SMPs)。它们相对容易实现,但可扩展性不强。程序员负责通过在管理资源的程序中插入适当的控制、信号量、锁等来管理同步。

  • 非统一内存访问NUMA):这些架构将内存划分为分配给每个处理器的高速访问区域,以及用于数据交换的公共区域,访问速度较慢。这些系统也被称为分布式共享内存DSM)系统。它们可扩展性很强,但开发复杂。

  • 无远程内存访问NoRMA):内存物理上分布在处理器之间(本地内存)。所有本地内存都是私有的,只能访问本地处理器。处理器之间的通信是通过用于交换消息的通信协议进行的,这被称为消息传递协议

  • 仅缓存内存架构COMA):这些系统只配备了缓存内存。在分析 NUMA 架构时,注意到这种架构将数据的本地副本存储在缓存中,并且这些数据在主内存中以副本的形式存储。这种架构消除了副本,只保留缓存内存;内存物理上分布在处理器之间(本地内存)。所有本地内存都是私有的,只能访问本地处理器。处理器之间的通信也是通过消息传递协议进行的。

分布式内存

在具有分布式内存的系统中,内存与每个处理器相关联,处理器只能访问其自己的内存。一些作者将这种类型的系统称为多计算机,反映了系统元素本身是小型且完整的处理器和内存系统的事实,如下面的图所示:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/edbcb72a-7807-4dba-97da-a1e69af3c29d.png

分布式内存架构方案

这种组织方式有几个优点:

  • 在通信总线或交换机层面没有冲突。每个处理器都可以使用其本地内存的全部带宽,而不会受到其他处理器的干扰。

  • 没有公共总线意味着处理器数量的内在限制。系统的规模仅受连接处理器的网络大小限制。

  • 没有缓存一致性方面的问题。每个处理器负责其自己的数据,不必担心升级任何副本。

主要缺点是处理器之间的通信更难实现。如果一个处理器需要另一个处理器的内存中的数据,那么这两个处理器不一定要通过消息传递协议交换消息。这引入了两个减速源:从一个处理器向另一个处理器构建和发送消息需要时间,而且,任何处理器都应该停止以管理从其他处理器接收到的消息。设计用于在分布式内存机器上工作的程序必须组织成一组独立任务,这些任务通过消息进行通信:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/6174b7b7-c606-48e1-a585-3412e6d542c7.png

基本消息传递

分布式内存系统的主要特点如下:

  • 内存在处理器之间物理分布;每个本地内存只能由其处理器直接访问。

  • 通过在处理器之间移动数据(即使只是消息本身)来实现同步。

  • 数据在本地内存中的细分会影响机器的性能——确保细分准确是至关重要的,以便最小化 CPU 之间的通信。此外,协调这些分解和组合操作的处理器必须有效地与操作数据结构各个部分的处理器进行通信。

  • 使用消息传递协议,以便 CPU 可以通过交换数据包相互通信。消息是离散的信息单元,从它们有明确的身份这一意义上说,它们总是可以相互区分。

大规模并行处理(MPP)

MPP 机器由数百个处理器(在某些机器中,处理器数量可以高达数十万个)组成,这些处理器通过通信网络连接。世界上最快的计算机基于这些架构;这些架构系统的例子包括地球模拟器、蓝基因、ASCI White、ASCI Red、ASCI Purple 和 Red Storm。

工作站集群

这些处理系统基于通过通信网络连接的经典计算机。计算集群属于这一分类。

在集群架构中,我们将节点定义为参与集群的单个计算单元。对于用户来说,集群是完全透明的——所有硬件和软件的复杂性都被掩盖,数据和应用程序都可以像来自单个节点一样访问。

在这里,我们已经确定了三种类型的集群:

  • 故障转移集群:在这种情况下,节点的活动会持续监控,当某个节点停止工作时,另一台机器将接管这些活动的责任。目的是通过架构的冗余确保连续的服务。

  • 负载均衡集群:在这个系统中,任务请求被发送到活动较少的节点。这确保了处理任务所需的时间更少。

  • 高性能计算集群:在这个集群中,每个节点都配置为提供极高的性能。过程也被划分为多个节点上的多个任务。任务被并行化,并将被分配到不同的机器上。

异构架构

在同构的超级计算世界中引入 GPU 加速器已经改变了超级计算机的使用和编程的本质。尽管 GPU 提供了高性能,但它们不能被视为一个自主的处理单元,因为它们应该始终伴随着 CPU 的组合。因此,编程范式非常简单:CPU 接管控制并以串行方式计算,将计算成本高且具有高度并行性的任务分配给图形加速器。

CPU 和 GPU 之间的通信不仅可以通过使用高速总线进行,还可以通过共享物理或虚拟内存的单个区域进行。实际上,在两种设备都没有配备自己的内存区域的情况下,可以使用由各种编程模型(如CUDAOpenCL)提供的软件库来引用一个公共内存区域。

这些架构被称为异构架构,其中应用程序可以在单个地址空间中创建数据结构,并将任务发送到设备硬件,这对于任务的解决是合适的。由于原子操作,几个处理任务可以在同一区域安全地运行,以避免数据一致性问题的发生。

因此,尽管 CPU 和 GPU 看起来似乎没有高效地协同工作,但使用这种新的架构,我们可以优化它们与并行应用程序的交互以及性能:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/46dac58e-d70e-4177-bc35-1018279093a9.png

异构架构方案

在下一节中,我们将介绍主要的并行编程模型。

并行编程模型

并行编程模型作为硬件和内存架构的抽象而存在。实际上,这些模型并不特定,也不指向任何特定的机器或内存架构。它们可以在任何类型的机器上实现(至少在理论上)。与之前的细分相比,这些编程模型在更高的层面上进行,代表了软件必须以何种方式实现以执行并行计算。每个模型都有其与其它处理器共享信息的方式,以便访问内存并分配工作。

从绝对意义上讲,没有一种模型比另一种更好。因此,最佳解决方案将非常依赖于程序员应该解决和解决的问题。最广泛使用的并行编程模型如下:

  • 共享内存模型

  • 多线程模型

  • 分布式内存/消息传递模型

  • 数据并行模型

在这个菜谱中,我们将为您概述这些模型。

共享内存模型

在这个模型中,任务共享一个单一的内存区域,我们可以异步地读取和写入。有机制允许编码者控制对共享内存的访问;例如,锁或信号量。这种模型的优势在于编码者不必明确任务间的通信。从性能的角度来看,一个重要的缺点是它变得难以理解和管理数据局部性。这指的是将数据保留在处理器的本地,以节省内存访问、缓存刷新和当多个处理器使用相同数据时发生的总线流量。

多线程模型

在这个模型中,一个进程可以有多个执行流。例如,创建一个顺序部分,随后创建一系列可以并行执行的任务。通常,这种类型的模型用于共享内存架构。因此,对于我们来说,管理线程之间的同步非常重要,因为它们在共享内存上操作,程序员必须防止多个线程同时更新同一位置。

当代 CPU 在软件和硬件层面都支持多线程。POSIX(即可移植操作系统接口)线程是软件层面实现多线程的典型例子。英特尔 Hyper-Threading 技术通过在某个线程停滞或等待 I/O 时切换到另一个线程,在硬件层面实现多线程。即使数据对齐是非线性的,从这个模型中也可以实现并行性。

消息传递模型

消息传递模型通常应用于每个处理器都有自己的内存(分布式内存系统)的情况。更多的任务可以驻留在同一台物理机器上或任意数量的机器上。编码者负责确定通过消息发生的并行性和数据交换,并且需要在代码中请求和调用函数库。

一些例子自 1980 年代以来就存在,但直到 1990 年代中期才创建了一个标准化的模型,从而产生了被称为消息传递接口MPI)的事实上的标准。

MPI 模型显然是为分布式内存设计的,但作为并行编程的模型,多平台模型也可以与共享内存机器一起使用:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/bd5deb5a-ea45-42d8-ba4e-b7852b6e0fcd.png

消息传递范式模型

数据并行模型

在此模型中,我们拥有更多在相同数据结构上操作的任务,但每个任务操作的是数据的不同部分。在共享内存架构中,所有任务都通过共享内存和分布式内存架构访问数据,其中数据结构被分割并驻留在每个任务的本地内存中。

为了实现此模型,编码者必须开发一个程序,该程序指定了数据的分布和对齐;例如,当前一代 GPU 只有在数据(任务 1任务 2任务 3)对齐时才能高效运行,如下面的图所示:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/93cf041f-f65b-46e5-b36d-59d4ed537910.png

数据并行范式模型

设计并行程序

利用并行性的算法设计基于一系列操作,这些操作必须执行,以便程序能够正确执行任务而不会产生部分或错误的结果。为了正确并行化一个算法,必须执行以下宏观操作:

  • 任务分解

  • 任务分配

  • 聚合

  • 映射

任务分解

在这个第一阶段,软件程序被分割成任务或一组指令,然后可以在不同的处理器上执行以实现并行化。为了执行这种细分,使用了两种方法:

  • 域分解:在此,问题的数据被分解。该应用对所有在数据的不同部分上工作的处理器都是通用的。当我们必须处理大量数据时,我们使用这种方法。

  • 功能分解:在这种情况下,问题被分割成任务,每个任务将对所有可用数据进行特定的操作。

任务分配

在此步骤中,指定了任务将在各个进程之间如何分配的机制。这一阶段非常重要,因为它建立了不同处理器之间的工作负载分配。负载均衡在这里至关重要;实际上,所有处理器都必须连续工作,避免长时间处于空闲状态。

为了执行此操作,编码者会考虑到系统可能的异构性,并尝试将更多任务分配给性能更好的处理器。最后,为了提高并行化的效率,有必要尽可能减少处理器之间的通信,因为它们往往是减速和资源消耗的来源。

聚合

聚合是将较小的任务与较大的任务结合起来的过程,以提高性能。如果设计过程的先前两个阶段将问题分割成远超过可用处理器数量的任务,并且如果计算机没有专门设计来处理大量的小任务(一些架构,如 GPU,处理这些任务很好,并且确实从运行数百万甚至数十亿的任务中受益),那么设计可能会非常低效。

通常,这是因为任务需要传达给处理器或线程,以便它们计算所述任务。大多数通信的成本与传输的数据量不成比例,但每次通信操作(如设置 TCP 连接时固有的延迟)都会产生固定成本。如果任务太小,那么这种固定成本很容易使设计变得低效。

映射

在并行算法设计过程的映射阶段,我们指定每个任务将在何处执行。目标是使总执行时间最小化。在这里,你通常必须做出权衡,因为两种主要策略往往相互冲突:

  • 应该将频繁通信的任务放置在同一处理器中,以增加局部性。

  • 应该将可以并发执行的任务放置在不同的处理器中,以增强并发性。

这被称为映射问题,并且已知它是NP 完全的。因此,在一般情况下,不存在该问题的多项式时间解。对于大小相等且具有易于识别的通信模式的任务,映射是直接的(我们也可以在这里进行聚簇,以将映射到同一处理器的任务组合在一起)。然而,如果任务的通信模式难以预测或每个任务的工作量不同,那么设计一个有效的映射和聚簇方案就很难。

对于这类问题,可以在运行时使用负载均衡算法来识别聚簇和映射策略。最困难的问题是那些在程序执行过程中通信量或任务数量发生变化的问题。对于这类问题,可以使用动态负载均衡算法,这些算法在执行期间定期运行。

动态映射

对于各种问题,存在许多负载均衡算法:

  • 全局算法:这些算法需要全局了解正在进行的计算,这通常会增加很多开销。

  • 本地算法:这些算法仅依赖于与所讨论任务局部相关的信息,与全局算法相比,这降低了开销,但它们通常在寻找最佳聚簇和映射方面表现较差。

然而,降低开销可能会减少执行时间,尽管映射本身可能更差。如果任务很少在执行的开始和结束时进行通信,那么通常会使用任务调度算法,该算法简单地将任务映射到空闲的处理器。在任务调度算法中,维护一个任务池。任务被放置在这个池中,并由工作者从池中取出。

在这个模型中有三种常见的方法:

  • 管理/工作员:这是所有工作员连接到一个集中式管理员的基动态映射方案。管理员反复向工作员发送任务并收集结果。这种策略可能适用于相对较少的处理器。通过提前获取任务,可以改进基本策略,以便通信和计算重叠。

  • 分层管理/工作员:这是具有半分布式布局的管理员/工作员变体。工作员被分成组,每组都有自己的管理员。这些组管理员与中央管理员(以及可能彼此之间)通信,而工作员从组管理员那里请求任务。这样可以在几个管理员之间分散负载,并且可以处理更多的处理器,如果所有工作员都从同一个管理员那里请求任务。

  • 去中心化:在这个方案中,一切都是去中心化的。每个处理器维护自己的任务池,并与其他处理器通信以请求任务。处理器如何选择其他处理器来请求任务各不相同,并且基于问题来确定。

评估并行程序的性能

并行编程的发展产生了对性能指标的需求,以便决定其使用是否方便。事实上,并行计算的重点是在相对较短的时间内解决大型问题。有助于实现这一目标的因素包括,例如,使用的硬件类型、问题的并行程度以及采用的并行编程模型。为了便于此,引入了基本概念的分析,它比较了从原始序列获得的并行算法。

通过分析和量化使用的线程数和/或进程数来实现性能。为了分析这一点,让我们引入一些性能指标:

  • 加速

  • 效率

  • 扩展性

并行计算的局限性是由阿姆达尔定律引入的。为了评估顺序算法并行化的效率程度,我们有古斯塔夫森定律

加速

加速是显示并行解决问题益处的度量。它定义为在单个处理元素(Ts)上解决问题所需的时间与在 p 个相同处理元素(Tp)上解决问题所需时间的比率。

我们如下表示加速:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/d06f79cc-0130-4c88-9669-be45342198b8.png

我们有一个线性加速,如果 S=p,这意味着执行速度会随着处理器数量的增加而增加。当然,这是一个理想的情况。当 Ts 是最佳顺序算法的执行时间时,加速是绝对的;当 Ts 是单个处理器的并行算法的执行时间时,加速是相对的。

让我们回顾这些条件:

  • S = p 是线性或理想加速。

  • S < p 是一个真实加速。

  • S > p 是一个超线性加速。

效率

在一个理想的世界里,具有 p 个处理单元的并行系统可以给我们一个等于 p 的加速。然而,这很少实现。通常,一些时间浪费在空闲或通信上。效率是衡量处理单元将多少执行时间用于有用工作的度量,表示为所花费时间的分数。

我们用 E 来表示它,并可以如下定义:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/a51cb8a0-063e-4177-a8e3-caa123753d53.png

具有线性加速的算法其值为 E = 1。在其他情况下,它们的值小于 1。以下是对三种情况进行的识别:

  • E = 1 时,这是一个线性情况。

  • E < 1 时,这是一个真实情况。

  • E << 1 时,这是一个低效率可并行化的问题。

规模化

规模化定义为在并行机器上保持效率的能力。它按处理器数量比例识别计算能力(执行速度)。通过增加问题的大小和同时增加处理器数量,在性能方面将不会有所损失。

可扩展的系统,根据不同因素的增量,可能保持相同的效率或提高效率。

阿姆达尔定律

阿姆达尔定律是一个广泛使用的定律,用于设计处理器和并行算法。它指出,可以达到的最大加速比受程序串行部分限制:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/d0ef21cb-ef7a-401d-990a-5a3b45325d05.png

1 – P 表示程序中(未并行化)的串行部分。

这意味着,例如,如果一个程序中 90%的代码可以并行化,但 10%必须保持串行,那么即使对于无限数量的处理器,最大可达到的加速比也是 9。

古斯塔夫森定律

古斯塔夫森定律表述如下:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/5b33302b-8561-4073-a6df-09072923e8f5.png

在这里,正如我们在方程中指出的,以下适用:

  • P处理器数量

  • S加速 因子。

  • α 是任何并行过程中 不可并行化部分

古斯塔夫森定律与阿姆达尔定律形成对比,正如我们描述的那样,阿姆达尔定律假设程序的整体工作量不会随着处理器数量的增加而改变。

事实上,古斯塔夫森定律建议程序员首先设定并行解决一个问题的 时间,然后基于这个(即时间)来调整问题的大小。因此,并行系统越快,在相同的时间内可以解决的问题就越多。

古斯塔夫森定律的影响是将计算机研究的目标指向选择或重新表述问题,以便在相同的时间内解决更大的问题仍然是可能的。此外,该定律重新定义了 效率 的概念,即需要 至少减少程序中的顺序部分,尽管 工作量增加

介绍 Python

Python 是一种强大、动态和解释型编程语言,广泛应用于各种应用。其一些特性如下:

  • 清晰且易于阅读的语法。

  • 一个非常广泛的标准库,通过额外的软件模块,我们可以添加数据类型、函数和对象。

  • 易于学习的快速开发和调试。在 Python 中开发 Python 代码可以比在 C/C++代码中快 10 倍。代码也可以作为原型,然后翻译成 C/C++。

  • 基于异常的错误处理。

  • 强大的内省功能。

  • 文档的丰富性和软件社区。

Python 可以被视为一种粘合语言。使用 Python,可以开发出更好的应用程序,因为不同类型的程序员可以一起在项目上工作。例如,在构建科学应用时,C/C++程序员可以实现高效的数值算法,而同一项目中的科学家可以编写 Python 程序来测试和使用这些算法。科学家不必学习低级编程语言,C/C++程序员也不需要理解涉及的科学。

您可以从www.python.org/doc/essays/omg-darpa-mcc-position了解更多相关信息。

让我们看看一些非常基本的代码示例,以了解 Python 的功能。

以下部分可以成为大多数人的复习资料。我们将在第二章基于线程的并行性和第三章基于进程的并行性中实际使用这些技术。

帮助函数

Python 解释器已经提供了一个有效的帮助系统。如果您想了解如何使用一个对象,只需键入help(object)

例如,让我们看看如何使用help函数在整数0上:

>>> help(0)
Help on int object:

class int(object)
 | int(x=0) -> integer
 | int(x, base=10) -> integer
 | 
 | Convert a number or string to an integer, or return 0 if no 
 | arguments are given. If x is a number, return x.__int__(). For 
 | floating point numbers, this truncates towards zero.
 | 
 | If x is not a number or if base is given, then x must be a string,
 | bytes, or bytearray instance representing an integer literal in the
 | given base. The literal can be preceded by '+' or '-' and be
 | surrounded by whitespace. The base defaults to 10\. Valid bases are 0 
 | and 2-36.
 | Base 0 means to interpret the base from the string as an integer 
 | literal.
>>> int('0b100', base=0)

int对象的描述后面跟着一个适用于它的方法列表。前五种方法如下:

 | Methods defined here:
 | 
 | __abs__(self, /)
 | abs(self)
 | 
 | __add__(self, value, /)
 | Return self+value.
 | 
 | __and__(self, value, /)
 | Return self&value.
 | 
 | __bool__(self, /)
 | self != 0
 | 
 | __ceil__(...)
 | Ceiling of an Integral returns itself.

同样有用的是dir(object),它列出了对象可用的方法:

>>> dir(float)
['__abs__', '__add__', '__and__', '__bool__', '__ceil__', '__class__', '__delattr__', '__dir__', '__divmod__', '__doc__', '__eq__', '__float__', '__floor__', '__floordiv__', '__format__', '__ge__', '__getattribute__', '__getnewargs__', '__gt__', '__hash__', '__index__', '__init__', '__int__', '__invert__', '__le__', '__lshift__', '__lt__', '__mod__', '__mul__', '__ne__', '__neg__', '__new__', '__or__', '__pos__', '__pow__', '__radd__', '__rand__', '__rdivmod__', '__reduce__', '__reduce_ex__', '__repr__', '__rfloordiv__', '__rlshift__', '__rmod__', '__rmul__', '__ror__', '__round__', '__rpow__', '__rrshift__', '__rshift__', '__rsub__', '__rtruediv__', '__rxor__', '__setattr__', '__sizeof__', '__str__', '__sub__', '__subclasshook__', '__truediv__', '__trunc__', '__xor__', 'bit_length', 'conjugate', 'denominator', 'from_bytes', 'imag', 'numerator', 'real', 'to_bytes']

最后,对象的相应文档由.__doc__函数提供,如下例所示:

>>> abs.__doc__
'Return the absolute value of the argument.'

语法

Python 不采用语句终止符,代码块通过缩进来指定。需要缩进级别的语句必须以冒号(:)结尾。这导致以下情况:

  • Python 代码更清晰、更易于阅读。

  • 程序结构始终与缩进一致。

  • 任何列表中的缩进风格都是统一的。

错误的缩进可能导致错误。

以下示例展示了如何使用if构造:

print("first print")
if condition:
    print(“second print)
print(“third print)

在这个例子中,我们可以看到以下内容:

  • 以下语句:print("first print")if condition:print("third print")具有相同的缩进级别,并且总是被执行。

  • if语句之后,有一个更高缩进级别的代码块,其中包含print ("second print")语句。

  • 如果if的条件为真,则执行print ("second print")语句。

  • 如果if的条件为假,则不执行print ("second print")语句。

因此,注意缩进非常重要,因为在程序解析过程中始终会评估缩进。

注释

注释以井号(#)开头,并且位于单行上:

# single line comment

多行字符串用于多行注释:

""" first line of a multi-line comment
second line of a multi-line comment."""

赋值

使用等号(=)进行赋值。对于相等性测试,使用相同的量(==)。你可以使用+=-=运算符增加和减少一个值,后面跟一个附加项。这适用于许多数据类型,包括字符串。你可以在同一行上赋值和使用多个变量。

一些例子如下:

>>> variable = 3
>>> variable += 2
>>> variable
5
>>> variable -= 1
>>> variable
4

>>> _string_ = "Hello"
>>> _string_ += " Parallel Programming CookBook Second Edition!"
>>> print (_string_) 
Hello Parallel Programming CookBook Second Edition!

数据类型

Python 中最显著的结构是列表元组字典。集合自 Python 2.5 版本以来已集成(旧版本可在sets库中找到):

  • 列表:这些类似于一维数组,但你可以创建包含其他列表的列表。

  • 字典:这些是包含键值对的数组(哈希表)。

  • 元组:这些是不可变的单维对象。

数组可以是任何类型,因此你可以将整数和字符串等变量混合到你的列表、字典和元组中。

任何类型数组的第一个对象的索引始终为零。允许使用负索引,并从数组的末尾开始计数;-1表示数组的最后一个元素:

#let's play with lists
list_1 = [1, ["item_1", "item_1"], ("a", "tuple")]
list_2 = ["item_1", -10000, 5.01]

>>> list_1
[1, ['item_1', 'item_1'], ('a', 'tuple')]

>>> list_2
['item_1', -10000, 5.01]

>>> list_1[2]
('a', 'tuple')

>>>list_1[1][0]
['item_1', 'item_1']

>>> list_2[0]
item_1

>>> list_2[-1]
5.01

#build a dictionary 
dictionary = {"Key 1": "item A", "Key 2": "item B", 3: 1000}
>>> dictionary 
{'Key 1': 'item A', 'Key 2': 'item B', 3: 1000} 

>>> dictionary["Key 1"] 
item A

>>> dictionary["Key 2"]
-1

>>> dictionary[3]
1000

你可以使用冒号(:)获取数组范围:

list_3 = ["Hello", "Ruvika", "how" , "are" , "you?"] 
>>> list_3[0:6] 
['Hello', 'Ruvika', 'how', 'are', 'you?'] 

>>> list_3[0:1]
['Hello']

>>> list_3[2:6]
['how', 'are', 'you?']

字符串

Python 字符串使用单引号(')或双引号(")表示,并且可以在由另一个分隔的字符串中使用一种表示法:

>>> example = "she loves ' giancarlo"
>>> example
"she loves ' giancarlo"

在多行中,它们被三重引号(或三个单引号)包围('''多行字符串'''):

>>> _string_='''I am a 
multi-line 
string'''
>>> _string_
'I am a \nmulti-line\nstring'

Python 也支持 Unicode;只需使用u "This is a unicode string"语法:

>>> ustring = u"I am unicode string"
>>> ustring
'I am unicode string'

要在字符串中输入值,请输入%运算符和一个元组。然后,每个%运算符被从左到右的元组元素替换:*

>>> print ("My name is %s !" % ('Mr. Wolf'))
My name is Mr. Wolf!

流控制

流控制指令是ifforwhile

在下一个例子中,我们检查数字是正数、负数还是零,并显示结果:

num = 1

if num > 0:
    print("Positive number")
elif num == 0:
    print("Zero")
else:
    print("Negative number")

以下代码块使用for循环找出存储在列表中的所有数字的总和:

numbers = [6, 6, 3, 8, -3, 2, 5, 44, 12]
sum = 0
for val in numbers:
    sum = sum+val
print("The sum is", sum)

我们将执行while循环,直到条件结果为真来迭代代码。由于我们不知道迭代次数,我们将使用这个循环而不是for循环。在这个例子中,我们使用while来计算自然数之和sum = 1+2+3+...+n

n = 10
# initialize sum and counter
sum = 0
i = 1
while i <= n:
    sum = sum + i
    i = i+1 # update counter

# print the sum
print("The sum is", sum)

前三个示例的输出如下:

Positive number
The sum is 83
The sum is 55
>>>

函数

Python 函数使用def关键字声明:

def my_function():
    print("this is a function")

要运行一个函数,使用函数名称,后跟括号,如下所示:

>>> my_function()
this is a function

参数必须在函数名之后、括号内指定:

def my_function(x):
    print(x * 1234)

>>> my_function(7)
8638

多个参数必须用逗号分隔:

def my_function(x,y):
    print(x*5+ 2*y)

>>> my_function(7,9)
53

使用等号来定义默认参数。如果你不传递参数调用函数,则将使用默认值:

def my_function(x,y=10):
    print(x*5+ 2*y)

>>> my_function(1)
25

>>> my_function(1,100)
205

函数的参数可以是任何类型的数据(例如字符串、数字、列表和字典)。在这里,以下列表 lcities 被用作 my_function 的参数:

def my_function(cities):
    for x in cities:
        print(x)

>>> lcities=["Napoli","Mumbai","Amsterdam"]
>>> my_function(lcities)
Napoli
Mumbai
Amsterdam

使用 return 语句从函数返回一个值:

def my_function(x,y):
    return x*y >>> my_function(6,29)
174 

Python 支持一种有趣的语法,允许你即时定义小型、单行的函数。这些 lambda 函数源自 Lisp 编程语言,可以在需要函数的地方使用。

一个 lambda 函数 functionvar 的示例如下所示:

# lambda definition equivalent to def f(x): return x + 1

functionvar = lambda x: x * 5
>>> print(functionvar(10))
50

Python 支持类的多重继承。传统上(不是语言规则),私有变量和方法通过在前面加上两个下划线(__)来声明。我们可以将任意属性(属性)分配给类的实例,如下例所示:

class FirstClass:
    common_value = 10
    def __init__ (self):
        self.my_value = 100
    def my_func (self, arg1, arg2):
        return self.my_value*arg1*arg2

# Build a first instance
>>> first_instance = FirstClass()
>>> first_instance.my_func(1, 2)
200

# Build a second instance of FirstClass
>>> second_instance = FirstClass()

#check the common values for both the instances
>>> first_instance.common_value
10

>>> second_instance.common_value
10

#Change common_value for the first_instance
>>> first_instance.common_value = 1500
>>> first_instance.common_value
1500

#As you can note the common_value for second_instance is not changed
>>> second_instance.common_value
10

# SecondClass inherits from FirstClass. 
# multiple inheritance is declared as follows:
# class SecondClass (FirstClass1, FirstClass2, FirstClassN)

class SecondClass (FirstClass):
    # The "self" argument is passed automatically
    # and refers to the class's instance
    def __init__ (self, arg1):
        self.my_value = 764
        print (arg1)

>>> first_instance = SecondClass ("hello PACKT!!!!")
hello PACKT!!!!

>>> first_instance.my_func (1, 2)
1528

异常

Python 中的异常通过 try-except 块(exception_name)来管理:

def one_function():
     try:
         # Division by zero causes one exception
         10/0
     except ZeroDivisionError:
         print("Oops, error.")
     else:
         # There was no exception, we can continue.
         pass
     finally:
         # This code is executed when the block
         # try..except is already executed and all exceptions
         # have been managed, even if a new one occurs
         # exception directly in the block.
         print("We finished.")

>>> one_function()
Oops, error.
We finished

导入库

使用 import [library name] 来导入外部库。或者,你可以使用 from [library name] import [function name] 语法来导入特定的函数。以下是一个示例:

import random
randomint = random.randint(1, 101)

>>> print(randomint)
65

from random import randint
randomint = random.randint(1, 102)

>>> print(randomint)
46

文件管理

为了让我们能够与文件系统交互,Python 提供了内置的 open 函数。这个函数可以被调用以打开一个文件并返回一个文件对象。后者允许我们对文件执行各种操作,如读取和写入。当我们完成与文件的交互后,我们必须最后记得使用 file.close 方法来关闭它:

>>> f = open ('test.txt', 'w') # open the file for writing
>>> f.write ('first line of file \ n') # write a line in file
>>> f.write ('second line of file \ n') # write another line in file
>>> f.close () # we close the file
>>> f = open ('test.txt') # reopen the file for reading
>>> content = f.read () # read all the contents of the file
>>> print (content)
first line of the file
second line of the file
>>> f.close () # close the file

列表推导式

列表推导式是创建和操作列表的有力工具。它们由一个表达式组成,后面跟着一个 for 子句,然后是一个或多个 if 子句。列表推导式的语法很简单,如下所示:

[expression for item in list]

然后,执行以下操作:

#list comprehensions using strings
>>> list_comprehension_1 = [ x for x in 'python parallel programming cookbook!' ]
>>> print( list_comprehension_1)

['p', 'y', 't', 'h', 'o', 'n', ' ', 'p', 'a', 'r', 'a', 'l', 'l', 'e', 'l', ' ', 'p', 'r', 'o', 'g', 'r', 'a', 'm', 'm', 'i', 'n', 'g', ' ', 'c', 'o', 'o', 'k', 'b', 'o', 'o', 'k', '!']

#list comprehensions using numbers
>>> l1 = [1,2,3,4,5,6,7,8,9,10]
>>> list_comprehension_2 = [ x*10 for x in l1 ]
>>> print( list_comprehension_2)

[10, 20, 30, 40, 50, 60, 70, 80, 90, 100]

运行 Python 脚本

要执行 Python 脚本,只需调用 Python 解释器后跟脚本名称,在这种情况下,my_pythonscript.py。或者,如果我们处于不同的工作目录,则使用其完整地址:

> python my_pythonscript.py 

从现在开始,对于每次调用 Python 脚本,我们将使用前面的表示法;即,python 后跟 script_name.py,假设启动 Python 解释器的目录是脚本要执行的目录。

使用 pip 安装 Python 包

pip 是一个工具,允许我们搜索、下载和安装在 Python 包索引上找到的 Python 包,该索引是一个包含成千上万用 Python 编写的包的存储库。这也允许我们管理已下载的包,使我们能够更新或删除它们。

安装 pip

pip 已经包含在 Python 版本 ≥ 3.4 和 ≥ 2.7.9 中。要检查此工具是否已安装,我们可以运行以下命令:

C:\>pip

如果 pip 已经安装,则此命令将显示已安装的版本。

更新 pip

建议检查您所使用的 pip 版本是否始终是最新的。要更新它,我们可以使用以下命令:

 C:\>pip install -U pip

使用 pip

pip 支持一系列命令,允许我们执行各种操作,包括 搜索、下载、安装、更新删除 包。

要安装 PACKAGE,只需运行以下命令:

C:\>pip install PACKAGE 

介绍 Python 并行编程

Python 提供了许多库和框架,这些库和框架有助于高性能计算。然而,由于 全局解释器锁GIL),使用 Python 进行并行编程可能会相当微妙。

事实上,最广泛和最常用的 Python 解释器 CPython 是用 C 编程语言开发的。CPython 解释器需要 GIL 以进行线程安全操作。使用 GIL 意味着当您尝试访问线程中包含的任何 Python 对象时,您将遇到全局锁。并且一次只有一个线程可以获取 Python 对象或 C API 的锁。

幸运的是,事情并没有那么严重,因为,在 GIL 的领域之外,我们可以自由地使用并行性。这一类别包括我们在下一章中将要讨论的所有主题,包括多进程、分布式计算和 GPU 计算。

因此,Python 并非真正的多线程。那么什么是线程?什么是进程?在接下来的章节中,我们将介绍这两个基本概念以及 Python 编程语言如何处理它们。

进程和线程

线程 可以与轻量级进程相比,从某种意义上说,它们提供了与进程类似的优势,但不需要进程的典型通信技术。线程允许您将程序的主要控制流程划分为多个并发运行的控制流。相比之下,进程有自己的 地址空间 和自己的资源。因此,在运行在不同进程上的代码部分之间的通信只能通过适当的管理机制进行,包括管道、代码 FIFO、邮箱、共享内存区域和消息传递。另一方面,线程允许创建程序的并发部分,其中每个部分都可以访问相同的地址空间、变量和常量。

以下表格总结了线程和进程之间的主要区别:

线程 进程
共享内存。 不共享内存。
开始/更改计算成本较低。 开始/更改计算成本较高。
需要较少的资源(轻量级进程)。 需要更多的计算资源。
需要同步机制来正确处理数据。 不需要内存同步。

在这简短的介绍之后,我们最终可以展示进程和线程是如何操作的。

尤其是我们想比较以下函数的串行、多线程和多进程执行时间,该函数do_something执行一些基本计算,包括构建一个随机选择的整数列表(do_something.py文件):

import random

def do_something(count, out_list):
  for i in range(count):
    out_list.append(random.random())

接下来,是串行(serial_test.py)实现。让我们从相关的导入开始:

from do_something import *
import time 

注意导入模块time,它将用于评估执行时间,在本例中,以及do_something函数的串行实现。要构建的列表size等于10000000,而do_something函数将被执行10次:

if __name__ == "__main__":
    start_time = time.time()
    size = 10000000 
    n_exec = 10
    for i in range(0, exec):
        out_list = list()
        do_something(size, out_list)

    print ("List processing complete.")
    end_time = time.time()
    print("serial time=", end_time - start_time)   

接下来,我们有多线程实现(multithreading_test.py)。

导入相关库:

from do_something import *
import time
import threading

注意导入threading模块的重要性,以便操作 Python 的多线程功能。

在这里,是do_something函数的多线程执行。我们不会深入评论以下代码中的指令,因为它们将在第二章中更详细地讨论,基于线程的并行性

然而,也应该注意,在这种情况下,列表的长度显然与串行情况相同,size = 10000000,而定义的线程数是 10,threads = 10,这也是do_something函数必须执行的次数:

if __name__ == "__main__":
    start_time = time.time()
    size = 10000000
    threads = 10 
    jobs = []
    for i in range(0, threads):

还要注意通过threading.Thread方法构建单个线程:

out_list = list()
thread = threading.Thread(target=list_append(size,out_list))
jobs.append(thread)

我们启动线程并立即停止它们的循环序列如下:

    for j in jobs:
        j.start()
    for j in jobs:
        j.join()

    print ("List processing complete.")
    end_time = time.time()
    print("multithreading time=", end_time - start_time)

最后,是多进程实现(multiprocessing_test.py)。

我们首先导入必要的模块,特别是multiprocessing库,其功能将在第三章中深入解释,基于进程的并行性

from do_something import *
import time
import multiprocessing

与前例一样,要构建的列表长度、do_something函数的大小和执行次数保持不变(procs = 10):

if __name__ == "__main__":
    start_time = time.time()
    size = 10000000 
    procs = 10 
    jobs = []
    for i in range(0, procs):
        out_list = list()

在这里,通过multiprocessing.Process方法调用实现单个进程受以下影响:

        process = multiprocessing.Process\
                  (target=do_something,args=(size,out_list))
        jobs.append(process)

接下来,启动进程并立即停止它们的循环序列执行如下:

    for j in jobs:
        j.start()

    for j in jobs:
        j.join()

    print ("List processing complete.")
    end_time = time.time()
    print("multiprocesses time=", end_time - start_time)

然后,我们打开命令行并运行前面描述的三个函数。

前往已复制函数的文件夹,然后输入以下内容:

> python serial_test.py

在具有以下特性的机器上获得的结果——CPU Intel i7/8 GB 的内存,如下所示:

List processing complete.
serial time= 25.428767204284668

多线程实现的情况下,我们有以下内容:

> python multithreading_test.py

输出如下:

List processing complete.
multithreading time= 26.168917179107666

最后,是多进程的实现:

> python multiprocessing_test.py

其结果如下:

List processing complete.
multiprocesses time= 18.929869890213013

如所示,串行实现(即使用serial_test.py)的结果与使用多线程实现(使用multithreading_test.py)获得的结果相似,其中线程实际上是依次启动的,优先考虑一个然后是另一个,直到结束,而我们在使用 Python 多进程能力(使用multiprocessing_test.py)方面获得了执行时间上的好处。

第二章:基于线程的并行

目前,在软件应用程序中管理并发最广泛使用的编程范式是基于多线程的。通常,一个应用程序由一个单一进程组成,该进程被划分为多个独立的线程,这些线程代表不同类型的活动,它们并行运行并相互竞争。

现在,使用多线程的现代应用程序已经被大规模采用。事实上,所有当前的处理器都是多核的,这样它们就可以执行并行操作并利用计算机的计算资源。

因此,多线程编程无疑是实现并发应用的好方法。然而,多线程编程往往隐藏一些非平凡的困难,这些困难必须得到适当的处理,以避免死锁或同步问题等错误。

我们将首先定义基于线程和多线程编程的概念,然后介绍multithreading库。我们将学习线程定义、管理和通信的主要指令。

通过multithreading库,我们将看到如何通过不同的技术解决问题,例如RLock信号量条件事件屏障队列

在本章中,我们将涵盖以下食谱:

  • 什么是线程?

  • 如何定义线程

  • 如何确定当前线程

  • 如何在子类中使用线程

  • 使用锁进行线程同步

  • 使用 RLock 进行线程同步

  • 使用信号量进行线程同步

  • 使用条件进行线程同步

  • 使用事件进行线程同步

  • 使用屏障进行线程同步

  • 使用队列进行线程通信

我们还将探讨 Python 提供的用于线程编程的主要选项。为此,我们将专注于使用threading模块。

什么是线程?

线程是一个可以与其他系统中的线程并行和并发执行的独立执行流程。

多个线程可以共享数据和资源,利用所谓的共享信息空间。线程和进程的具体实现取决于你计划在哪个操作系统上运行应用程序,但一般来说,可以这样说,线程包含在进程内部,并且同一进程中的不同线程条件共享一些资源。相比之下,不同的进程不会与其他进程共享它们自己的资源。

线程由三个元素组成:程序计数器、寄存器和栈。与同一进程中的其他线程共享的资源基本上包括数据操作系统资源。此外,线程有自己的执行状态,即线程状态,并且可以与其他线程同步

线程状态可以是就绪、运行或阻塞:

  • 当线程被创建时,它进入就绪状态。

  • 线程由操作系统(或运行时支持系统)调度执行,当轮到它时,它进入运行状态开始执行。

  • 线程可以等待一个条件发生,从运行状态转换为阻塞状态。一旦锁定条件终止,阻塞线程将返回到就绪状态:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/1c9d8391-719e-4277-a1ae-dd4155345659.png

线程生命周期

多线程编程的主要优势在于性能,因为进程之间的上下文切换比同一进程内的线程之间的上下文切换要重得多。

在下一节中,直到本章结束,我们将检查 Python 的threading模块,通过编程示例介绍其主要功能。

Python 线程模块

Python 使用 Python 标准库提供的threading模块来管理线程。此模块提供了一些非常有趣的功能,使得基于线程的方法变得非常简单;实际上,threading模块提供了几个非常简单的同步机制。

线程模块的主要组件如下:

  • thread对象

  • lock对象

  • RLock对象

  • semaphore对象

  • condition对象

  • event对象

在下面的菜谱中,我们通过不同的应用示例检查threading库提供的功能。对于下面的示例,我们将参考 Python 3.5.0 发行版(www.python.org/downloads/release/python-350/)。

定义线程

使用线程的最简单方法是使用目标函数实例化它,然后调用 start 方法让它开始工作。

准备工作

Python 的threading模块提供了一个Thread类,用于在不同的线程中运行进程和函数:

class threading.Thread(group=None, 
                       target=None, 
                       name=None, 
                       args=(), 
                       kwargs={})  

这里是Thread类的参数:

  • group:这是group值,应该是None;这是为未来的实现保留的。

  • target:这是启动线程活动时要执行的功能。

  • name:这是线程的名称;默认情况下,分配给它一个唯一的名称,形式为Thread-N

  • args:这是要传递给目标函数的参数元组。

  • kwargs:这是用于target函数的关键字参数字典。

在下一节中,我们将学习如何定义线程。

如何实现…

我们将通过传递一个数字来定义线程,这个数字代表线程号,最后将打印出结果:

  1. 使用以下 Python 命令导入threading模块:
import threading
  1. main程序中,使用名为my_func的目标函数实例化一个Thread对象。然后,将包含在输出消息中的函数参数传递:
t = threading.Thread(target=function , args=(i,))
  1. 线程只有在调用start方法后才开始运行,join方法使调用线程等待,直到线程完成执行,如下所示:
import threading

def my_func(thread_number):
    return print('my_func called by thread N°\
        {}'.format(thread_number))

def main():
    threads = []
    for i in range(10):
        t = threading.Thread(target=my_func, args=(i,))
        threads.append(t)
        t.start()
        t.join()

if __name__ == "__main__":
    main()

它是如何工作的…

main程序中,我们初始化线程列表,并将创建的每个线程的实例添加到该列表中。创建的线程总数为 10,而 i^(th)线程的i-索引作为参数传递给 i^(th)线程:

my_func called by thread N°0
my_func called by thread N°1
my_func called by thread N°2
my_func called by thread N°3
my_func called by thread N°4
my_func called by thread N°5
my_func called by thread N°6
my_func called by thread N°7
my_func called by thread N°8
my_func called by thread N°9

还有更多…

所有当前处理器都是多核的,因此提供了执行多个并行操作的可能性,并充分利用计算机的计算资源。尽管这是真的,但多线程编程隐藏了许多非平凡困难,这些困难必须得到适当的处理,以避免死锁或同步问题等错误。

确定当前线程

使用参数来识别或命名线程是繁琐且不必要的。每个Thread实例都有一个name,具有默认值,可以在创建线程时更改。

在具有多个服务线程处理不同操作的服务器进程中,命名线程是有用的。

准备工作

这个threading模块提供了currentThread().getName()方法,它返回当前线程的名称。

以下部分展示了我们如何使用此功能来确定哪个线程正在运行。

如何做…

让我们看看以下步骤:

  1. 要确定哪个线程正在运行,我们创建了三个target函数,并导入time模块以引入两秒的暂停执行:
import threading
import time

def function_A():
    print (threading.currentThread().getName()+str('-->\
        starting \n'))
    time.sleep(2)
    print (threading.currentThread().getName()+str( '-->\
        exiting \n'))

def function_B():
    print (threading.currentThread().getName()+str('-->\
        starting \n'))
    time.sleep(2)
    print (threading.currentThread().getName()+str( '-->\
        exiting \n'))

def function_C():
    print (threading.currentThread().getName()+str('-->\
        starting \n'))
    time.sleep(2)
    print (threading.currentThread().getName()+str( '-->\
        exiting \n'))

  1. 使用target函数创建了三个线程。然后,我们传递要打印的名称,如果没有定义,则使用默认名称。然后,对每个线程调用start()join()方法:
if __name__ == "__main__":

    t1 = threading.Thread(name='function_A', target=function_A)
    t2 = threading.Thread(name='function_B', target=function_B)
    t3 = threading.Thread(name='function_C',target=function_C) 

    t1.start()
    t2.start()
    t3.start()

    t1.join()
    t2.join()
    t3.join()

它是如何工作的…

我们将设置三个线程,每个线程都分配了一个target函数。当target函数执行并终止时,将适当地打印出函数名称。

对于这个例子,输出应该看起来像这样(即使显示的顺序可能不同):

function_A--> starting 
function_B--> starting 
function_C--> starting 

function_A--> exiting 
function_B--> exiting 
function_C--> exiting

定义线程子类

创建一个线程可能需要定义一个子类,该子类从Thread类继承。正如在定义线程部分中解释的那样,后者包含在threading模块中,然后必须导入该模块。

准备工作

我们将在下一节定义的类,它代表我们的线程,遵循一个精确的结构:我们首先必须定义**__init__**方法,但最重要的是,我们必须重写run方法。

如何做…

涉及的步骤如下:

  1. 我们定义了MyThreadClass类,我们可以使用它来创建我们想要的任何线程。这种类型的每个线程都将由run方法中定义的操作所特征化,在这个简单的例子中,它限制自己在执行开始和结束时打印一个字符串:
import time
import os
from random import randint
from threading import Thread

class MyThreadClass (Thread):
  1. 此外,在__init__方法中,我们指定了两个初始化参数,分别是nameduration,它们将在run方法中使用:
def __init__(self, name, duration):
      Thread.__init__(self)
      self.name = name
      self.duration = duration 

   def run(self):
      print ("---> " + self.name +\
             " running, belonging to process ID "\
             + str(os.getpid()) + "\n")
      time.sleep(self.duration)
      print ("---> " + self.name + " over\n")
  1. 这些参数将在创建线程时设置。特别是,duration参数使用randint函数计算,该函数输出一个介于110之间的随机整数。从MyThreadClass的定义开始,让我们看看如何实例化更多线程,如下所示:
def main():

    start_time = time.time()

    # Thread Creation
    thread1 = MyThreadClass("Thread#1 ", randint(1,10))
    thread2 = MyThreadClass("Thread#2 ", randint(1,10))
    thread3 = MyThreadClass("Thread#3 ", randint(1,10))
    thread4 = MyThreadClass("Thread#4 ", randint(1,10))
    thread5 = MyThreadClass("Thread#5 ", randint(1,10))
    thread6 = MyThreadClass("Thread#6 ", randint(1,10))
    thread7 = MyThreadClass("Thread#7 ", randint(1,10))
    thread8 = MyThreadClass("Thread#8 ", randint(1,10)) 
    thread9 = MyThreadClass("Thread#9 ", randint(1,10))

    # Thread Running
    thread1.start()
    thread2.start()
    thread3.start()
    thread4.start()
    thread5.start()
    thread6.start()
    thread7.start()
    thread8.start()
    thread9.start()

    # Thread joining
    thread1.join()
    thread2.join()
    thread3.join()
    thread4.join()
    thread5.join()
    thread6.join()
    thread7.join()
    thread8.join()
    thread9.join()

    # End 
    print("End")

    #Execution Time
    print("--- %s seconds ---" % (time.time() - start_time))

if __name__ == "__main__":
    main()

它是如何工作的…

在这个例子中,我们创建了九个线程,每个线程都有自己的nameduration属性,这些属性根据__init__方法的定义。

我们然后使用start方法运行它们,该方法仅限于执行先前定义的run方法的内容。请注意,每个线程的进程 ID 是相同的,这意味着我们处于一个多线程进程。

此外,请注意,start方法不是阻塞的:当它执行时,控制权立即转到下一行,而线程则在后台启动。实际上,正如您所看到的,线程的创建并不按照代码指定的顺序进行。同样,线程的终止受duration参数值的约束,该参数使用randint函数评估,并通过每个线程创建实例的参数传递。要等待线程完成,必须执行join操作。

输出看起来像这样:

---> Thread#1 running, belonging to process ID 13084
---> Thread#5 running, belonging to process ID 13084
---> Thread#2 running, belonging to process ID 13084
---> Thread#6 running, belonging to process ID 13084
---> Thread#7 running, belonging to process ID 13084
---> Thread#3 running, belonging to process ID 13084
---> Thread#4 running, belonging to process ID 13084
---> Thread#8 running, belonging to process ID 13084
---> Thread#9 running, belonging to process ID 13084

---> Thread#6 over
---> Thread#9 over
---> Thread#5 over
---> Thread#2 over
---> Thread#7 over
---> Thread#4 over
---> Thread#3 over
---> Thread#8 over
---> Thread#1 over

End

--- 9.117518663406372 seconds ---

还有更多…

与面向对象编程(OOP)最常相关联的特性是继承,这是将一个新类定义为已存在类的修改版本的能力。继承的主要优势是您可以在不修改原始定义的情况下向类中添加新方法。

原始类通常被称为父类,而派生类被称为子类。继承是一个强大的特性,某些程序可以编写得更加容易和简洁,提供了在不修改原始类的情况下自定义类行为的可能性。实际上,继承结构能够反映问题的结构,在某些情况下,可以使程序更容易理解。

然而(为了提醒用户注意!),继承可能会使程序更难阅读。这是因为,在调用方法时,并不总是清楚这个方法是在代码的哪个地方定义的,而这个代码需要在多个模块中追踪,而不是在一个定义良好的地方。

通常,可以使用继承完成的事情,即使没有它也可以优雅地管理,因此只有在问题的结构需要时才应该使用继承。如果使用不当,那么继承可能造成的危害可能会超过使用它的好处。

使用锁进行线程同步

threading模块还包括一个简单的锁机制,这允许我们在线程之间实现同步。

准备工作

不过是一个通常可以被多个线程访问的对象,线程在进入程序受保护部分的执行之前必须拥有它。这些锁是通过执行Lock()方法创建的,该方法定义在threading模块中。

一旦创建了锁,我们可以使用两种方法来同步两个(或更多)线程的执行:acquire()方法用于获取锁控制,release()方法用于释放它。

acquire()方法接受一个可选参数,如果未指定或设置为True,将强制线程暂停其执行,直到锁被释放并可以获取。另一方面,如果acquire()方法以等于False的参数执行,则它立即返回一个布尔结果,如果锁已被获取,则为True,否则为False

在以下示例中,我们通过修改上一节中引入的代码,定义线程子类,来展示锁机制。

如何做到…

涉及的步骤如下:

  1. 如以下代码块所示,MyThreadClass类已被修改,在**run**方法中引入了acquire()release()方法,而Lock()的定义则位于类定义之外:
import threading
import time
import os
from threading import Thread
from random import randint

# Lock Definition
threadLock = threading.Lock()

class MyThreadClass (Thread):
   def __init__(self, name, duration):
      Thread.__init__(self)
      self.name = name
      self.duration = duration
   def run(self):
      #Acquire the Lock
      threadLock.acquire() 
      print ("---> " + self.name + \
             " running, belonging to process ID "\
             + str(os.getpid()) + "\n")
      time.sleep(self.duration)
      print ("---> " + self.name + " over\n")
      #Release the Lock
      threadLock.release()
  1. 与之前的代码示例相比,main()函数没有发生变化:
def main():
    start_time = time.time()
    # Thread Creation
    thread1 = MyThreadClass("Thread#1 ", randint(1,10))
    thread2 = MyThreadClass("Thread#2 ", randint(1,10))
    thread3 = MyThreadClass("Thread#3 ", randint(1,10))
    thread4 = MyThreadClass("Thread#4 ", randint(1,10))
    thread5 = MyThreadClass("Thread#5 ", randint(1,10))
    thread6 = MyThreadClass("Thread#6 ", randint(1,10))
    thread7 = MyThreadClass("Thread#7 ", randint(1,10))
    thread8 = MyThreadClass("Thread#8 ", randint(1,10))
    thread9 = MyThreadClass("Thread#9 ", randint(1,10))

    # Thread Running
    thread1.start()
    thread2.start()
    thread3.start()
    thread4.start()
    thread5.start()
    thread6.start()
    thread7.start()
    thread8.start()
    thread9.start()

    # Thread joining
    thread1.join()
    thread2.join()
    thread3.join()
    thread4.join()
    thread5.join()
    thread6.join()
    thread7.join()
    thread8.join()
    thread9.join()

    # End 
    print("End")
    #Execution Time
    print("--- %s seconds ---" % (time.time() - start_time))

if __name__ == "__main__":
    main()

它是如何工作的…

我们通过使用锁修改了上一节的代码,以便线程将按顺序执行。

第一个线程获取锁并执行其任务,而其他八个线程则保持等待状态。在第一个线程执行结束后,即执行release()方法后,第二个线程将获取锁,而三到八个线程将继续等待,直到执行结束(即再次,只有在运行release()方法之后)。

锁获取锁释放的执行会重复进行,直到第九个线程,最终结果是,由于锁机制,此执行以顺序模式进行,如下面的输出所示:

---> Thread#1 running, belonging to process ID 10632
---> Thread#1 over
---> Thread#2 running, belonging to process ID 10632
---> Thread#2 over
---> Thread#3 running, belonging to process ID 10632
---> Thread#3 over
---> Thread#4 running, belonging to process ID 10632
---> Thread#4 over
---> Thread#5 running, belonging to process ID 10632
---> Thread#5 over
---> Thread#6 running, belonging to process ID 10632
---> Thread#6 over
---> Thread#7 running, belonging to process ID 10632
---> Thread#7 over
---> Thread#8 running, belonging to process ID 10632
---> Thread#8 over
---> Thread#9 running, belonging to process ID 10632
---> Thread#9 over

End

--- 47.3672661781311 seconds ---

还有更多…

acquire()release()方法的插入点决定了整个代码的执行。因此,您花时间分析您想要使用的线程以及您想要如何同步它们非常重要。

例如,我们可以将MyThreadClass类中release()方法的插入点更改如下:

import threading
import time
import os
from threading import Thread
from random import randint

# Lock Definition
threadLock = threading.Lock()

class MyThreadClass (Thread):
   def __init__(self, name, duration):
      Thread.__init__(self)
      self.name = name
      self.duration = duration
   def run(self):
      #Acquire the Lock
      threadLock.acquire() 
      print ("---> " + self.name + \
             " running, belonging to process ID "\ 
             + str(os.getpid()) + "\n")
      #Release the Lock in this new point
      threadLock.release()
      time.sleep(self.duration)
      print ("---> " + self.name + " over\n")

在这种情况下,输出发生了相当显著的变化:

---> Thread#1 running, belonging to process ID 11228
---> Thread#2 running, belonging to process ID 11228
---> Thread#3 running, belonging to process ID 11228
---> Thread#4 running, belonging to process ID 11228
---> Thread#5 running, belonging to process ID 11228
---> Thread#6 running, belonging to process ID 11228
---> Thread#7 running, belonging to process ID 11228
---> Thread#8 running, belonging to process ID 11228
---> Thread#9 running, belonging to process ID 11228

---> Thread#2 over
---> Thread#4 over
---> Thread#6 over
---> Thread#5 over
---> Thread#1 over
---> Thread#3 over
---> Thread#9 over
---> Thread#7 over
---> Thread#8 over

End
--- 6.11468243598938 seconds ---

如您所见,只有线程创建是以顺序模式发生的。一旦线程创建完成,新线程获取锁,而前一个线程则在后台继续计算。

使用 RLock 进行线程同步

可重入锁,或简称 RLock,是一种可以被同一线程多次获取的同步原语。

它使用了专有线程的概念。这意味着在锁定状态下,一些线程拥有锁,而在未锁定状态下,锁不被任何线程拥有。

下一个示例演示了如何通过RLock()机制管理线程。

准备就绪

RLock是通过threading.RLock()类实现的。它提供了与threading.Lock()类语法相同的acquire()release()方法。

RLock块可以被同一线程多次获取。其他线程将无法获取RLock块,直到拥有它的线程为每个之前的acquire()调用执行了release()调用。实际上,RLock块必须被释放,但只能由获取它的线程释放。

如何做到这一点…

涉及的步骤如下:

  1. 我们引入了Box类,它提供了add()remove()方法,这些方法通过访问execute()方法来执行添加或删除项目的操作。对execute()方法的访问由RLock()调节:
import threading
import time
import random

class Box:
    def __init__(self):
        self.lock = threading.RLock()
        self.total_items = 0

    def execute(self, value):
        with self.lock:
            self.total_items += value

    def add(self):
        with self.lock:
            self.execute(1)

    def remove(self):
        with self.lock:
            self.execute(-1)
  1. 以下函数由两个线程调用。它们有box类和要添加或删除的总items数量作为参数:
def adder(box, items):
    print("N° {} items to ADD \n".format(items))
    while items:
        box.add()
        time.sleep(1)
        items -= 1
        print("ADDED one item -->{} item to ADD \n".format(items))

def remover(box, items):
    print("N° {} items to REMOVE\n".format(items))
    while items:
        box.remove()
        time.sleep(1)
        items -= 1
        print("REMOVED one item -->{} item to REMOVE\
            \n".format(items))
  1. 在这里,设置要添加到或从盒子中移除的项目总数。如您所见,这两个数字将不同。当adderremover方法完成其任务时,执行结束:
def main():
    items = 10
    box = Box()

    t1 = threading.Thread(target=adder, \
                          args=(box, random.randint(10,20)))
    t2 = threading.Thread(target=remover, \
                          args=(box, random.randint(1,10)))

    t1.start()
    t2.start()

    t1.join()
    t2.join()

if __name__ == "__main__":
    main()

它是如何工作的…

main程序中,t1t2两个线程已经与adder()remover()函数相关联。如果项目的数量大于零,则函数是活跃的。

RLock()的调用是在**Box**类的__init__方法中进行的:

class Box:
    def __init__(self):
        self.lock = threading.RLock()
        self.total_items = 0

两个adder()remover()函数分别与Box类的项目交互,并调用Box类的add()remove()方法。

在每次方法调用中,通过在_init_方法中设置的lock参数捕获资源,然后释放资源。

这里是输出结果:

16 items to ADD 
N° 1 items to REMOVE 

ADDED one item -->15 item to ADD 
REMOVED one item -->0 item to REMOVE 

ADDED one item -->14 item to ADD 
ADDED one item -->13 item to ADD 
ADDED one item -->12 item to ADD 
ADDED one item -->11 item to ADD 
ADDED one item -->10 item to ADD 
ADDED one item -->9 item to ADD 
ADDED one item -->8 item to ADD 
ADDED one item -->7 item to ADD 
ADDED one item -->6 item to ADD 
ADDED one item -->5 item to ADD 
ADDED one item -->4 item to ADD 
ADDED one item -->3 item to ADD 
ADDED one item -->2 item to ADD 
ADDED one item -->1 item to ADD 
ADDED one item -->0 item to ADD 
>>>

还有更多…

lockRLock之间的区别如下:

  • 一个在必须释放之前只能被获取一次。然而,RLock可以从同一线程多次获取;必须以相同的方式释放相同次数,才能释放。

  • 另一个区别是,获取的锁可以被任何线程释放,而获取的RLock只能由获取它的线程释放。

使用信号量进行线程同步

信号量是一种由操作系统管理的高级数据类型,用于同步多个线程对共享资源和数据的访问。它包含一个内部变量,用于标识与它关联的资源并发访问的数量。

准备就绪

信号量的操作基于两个函数:acquire()release(),如这里所述:

  • 当一个线程想要访问与信号量关联的给定资源时,它必须调用 acquire() 操作,这将减少信号量的内部变量,如果这个变量的值看起来是非负的,则允许访问资源。如果值是负的,则线程将被挂起,并且另一个线程释放资源的操作将被暂停。

  • 使用完共享资源后,线程通过 release() 指令释放资源。这样,信号量的内部变量就会增加,为等待的线程(如果有的话)提供了访问新释放资源的机遇。

信号量是计算机科学历史上最古老的同步原语之一,由早期荷兰计算机科学家 Edsger W. Dijkstra 发明。

以下示例展示了如何通过信号量同步线程。

如何实现…

以下代码描述了一个问题,其中我们有两个线程,producer()consumer(),它们共享一个公共资源,即项目。producer() 的任务是生成项目,而 consumer() 线程的任务是使用已经被生产的项目。

如果项目尚未由 consumer() 线程生产,那么它必须等待。一旦项目被生产,producer() 线程通知消费者资源应该被使用:

  1. 通过将信号量初始化为 0,我们获得了一个所谓的信号量事件,其唯一目的是同步两个或更多线程的计算。在这里,一个线程必须同时使用数据或共享资源:
semaphore = threading.Semaphore(0)
  1. 这个操作与锁机制的描述非常相似。producer() 线程创建项目,然后通过调用 release() 方法释放资源:
semaphore.release()
  1. 同样,consumer() 线程通过 acquire() 方法获取数据。如果信号量的计数器等于 0,则它将阻塞条件的 acquire() 方法,直到它被另一个线程通知。如果信号量的计数器大于 0,则它将减少该值。当生产者创建一个项目时,它释放信号量,然后消费者获取它并消费共享资源:
semaphore.acquire()
  1. 通过信号量完成的同步过程在以下代码块中展示:
import logging
import threading
import time
import random

LOG_FORMAT = '%(asctime)s %(threadName)-17s %(levelname)-8s %\
              (message)s'
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)

semaphore = threading.Semaphore(0)
item = 0

def consumer():
    logging.info('Consumer is waiting')
    semaphore.acquire()
    logging.info('Consumer notify: item number {}'.format(item))

def producer():
    global item
    time.sleep(3)
    item = random.randint(0, 1000)
    logging.info('Producer notify: item number {}'.format(item))
    semaphore.release()

#Main program
def main():
    for i in range(10):
        t1 = threading.Thread(target=consumer)
        t2 = threading.Thread(target=producer)

        t1.start()
        t2.start()

        t1.join()
        t2.join()

if __name__ == "__main__":
    main()

它是如何工作的…

然后将在标准输出上打印获取的数据:

print ("Consumer notify : consumed item number %s " %item)

这是我们在 10 次运行后得到的结果:

2019-01-27 19:21:19,354 Thread-1 INFO Consumer is waiting
2019-01-27 19:21:22,360 Thread-2 INFO Producer notify: item number 388
2019-01-27 19:21:22,385 Thread-1 INFO Consumer notify: item number 388
2019-01-27 19:21:22,395 Thread-3 INFO Consumer is waiting
2019-01-27 19:21:25,398 Thread-4 INFO Producer notify: item number 939
2019-01-27 19:21:25,450 Thread-3 INFO Consumer notify: item number 939
2019-01-27 19:21:25,453 Thread-5 INFO Consumer is waiting
2019-01-27 19:21:28,459 Thread-6 INFO Producer notify: item number 388
2019-01-27 19:21:28,468 Thread-5 INFO Consumer notify: item number 388
2019-01-27 19:21:28,476 Thread-7 INFO Consumer is waiting
2019-01-27 19:21:31,478 Thread-8 INFO Producer notify: item number 700
2019-01-27 19:21:31,529 Thread-7 INFO Consumer notify: item number 700
2019-01-27 19:21:31,538 Thread-9 INFO Consumer is waiting
2019-01-27 19:21:34,539 Thread-10 INFO Producer notify: item number 685
2019-01-27 19:21:34,593 Thread-9 INFO Consumer notify: item number 685
2019-01-27 19:21:34,603 Thread-11 INFO Consumer is waiting
2019-01-27 19:21:37,604 Thread-12 INFO Producer notify: item number 503
2019-01-27 19:21:37,658 Thread-11 INFO Consumer notify: item number 503
2019-01-27 19:21:37,668 Thread-13 INFO Consumer is waiting
2019-01-27 19:21:40,670 Thread-14 INFO Producer notify: item number 690
2019-01-27 19:21:40,719 Thread-13 INFO Consumer notify: item number 690
2019-01-27 19:21:40,729 Thread-15 INFO Consumer is waiting
2019-01-27 19:21:43,731 Thread-16 INFO Producer notify: item number 873
2019-01-27 19:21:43,788 Thread-15 INFO Consumer notify: item number 873
2019-01-27 19:21:43,802 Thread-17 INFO Consumer is waiting
2019-01-27 19:21:46,807 Thread-18 INFO Producer notify: item number 691
2019-01-27 19:21:46,861 Thread-17 INFO Consumer notify: item number 691
2019-01-27 19:21:46,874 Thread-19 INFO Consumer is waiting
2019-01-27 19:21:49,876 Thread-20 INFO Producer notify: item number 138
2019-01-27 19:21:49,924 Thread-19 INFO Consumer notify: item number 138
>>>

更多内容…

信号量的一个特定用途是互斥锁。互斥锁不过是一个内部变量初始化为 1 的信号量,它允许实现对数据和资源的互斥访问。

信号量在多线程编程语言中仍然被广泛使用;然而,它们有两个主要问题,我们已经在以下内容中讨论过:

  • 它们并不能阻止线程在同一个信号量上执行更多等待操作的可能性。很容易忘记与执行等待操作的数量相关的所有必要的信号。

  • 你可能会遇到死锁的情况。例如,当 t1 线程在 s1 信号量上执行等待操作,而 t2 线程在 t1 线程上执行等待操作,然后等待 s2t2,最后等待 s1 时,就会创建一个死锁情况。

使用条件进行线程同步

条件 识别应用程序中的状态变化。它是一种同步机制,其中线程等待特定条件,而另一个线程通知该 条件已经发生

一旦条件成立,线程 获取 锁以获得对共享资源的 独占访问

准备工作

通过再次查看生产者/消费者问题,可以很好地说明这种机制。如果缓冲区不满,producer 类将写入缓冲区,如果缓冲区已满,consumer 类将从缓冲区中取出数据(从后者中删除)。producer 类将通知消费者缓冲区不为空,而消费者将向生产者报告缓冲区不为空。

如何做到这一点…

涉及的步骤如下:

  1. consumer 类通过 items[] 列表获取共享资源:
condition.acquire()
  1. 如果列表长度等于 0,则消费者将处于等待状态:
if len(items) == 0:
   condition.wait()
  1. 然后它从项目列表中执行一个 pop 操作:
items.pop()
  1. 因此,将消费者状态通知给生产者,并释放共享资源:
condition.notify()
  1. producer 类获取共享资源,然后它验证列表是否完全填满(在我们的例子中,我们放置了项目列表中可以包含的最大项目数,10)。如果列表已满,则生产者将处于等待状态,直到列表被消耗:
condition.acquire()
if len(items) == 10:
   condition.wait()
  1. 如果列表未满,则添加一个单独的项目。状态被通知,资源被释放:
condition.notify()
condition.release()
  1. 为了向您展示条件机制,我们将再次使用 消费者/生产者 模型:
import logging
import threading
import time

LOG_FORMAT = '%(asctime)s %(threadName)-17s %(levelname)-8s %\
             (message)s'
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)

items = []
condition = threading.Condition()

class Consumer(threading.Thread):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)

    def consume(self):

        with condition:

            if len(items) == 0:
                logging.info('no items to consume')
                condition.wait()

            items.pop()
            logging.info('consumed 1 item')

            condition.notify()

    def run(self):
        for i in range(20):
            time.sleep(2)
            self.consume()

class Producer(threading.Thread):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)

    def produce(self):

        with condition:

            if len(items) == 10:
                logging.info('items produced {}.\
                    Stopped'.format(len(items)))
                condition.wait()

            items.append(1)
            logging.info('total items {}'.format(len(items)))

            condition.notify()

    def run(self):
        for i in range(20):
            time.sleep(0.5)
            self.produce()

它是如何工作的…

producer 持续生成项目并将其存储在缓冲区中。同时,consumer 使用生产出的数据,并时不时地从缓冲区中移除它。

一旦 consumer 从缓冲区中取走一个对象,它将唤醒 producer,然后 producer 将开始再次填充缓冲区。

类似地,如果缓冲区为空,consumer 将会挂起。一旦 producer 将数据下载到缓冲区,consumer 将会唤醒。

如您所见,即使在在这种情况下,使用 condition 指令也能使线程得到适当的同步。

单次运行后得到的结果如下:

2019-08-05 14:33:44,285 Producer INFO total items 1
2019-08-05 14:33:44,786 Producer INFO total items 2
2019-08-05 14:33:45,286 Producer INFO total items 3
2019-08-05 14:33:45,786 Consumer INFO consumed 1 item
2019-08-05 14:33:45,787 Producer INFO total items 3
2019-08-05 14:33:46,287 Producer INFO total items 4
2019-08-05 14:33:46,788 Producer INFO total items 5
2019-08-05 14:33:47,289 Producer INFO total items 6
2019-08-05 14:33:47,787 Consumer INFO consumed 1 item
2019-08-05 14:33:47,790 Producer INFO total items 6
2019-08-05 14:33:48,291 Producer INFO total items 7
2019-08-05 14:33:48,792 Producer INFO total items 8
2019-08-05 14:33:49,293 Producer INFO total items 9
2019-08-05 14:33:49,788 Consumer INFO consumed 1 item
2019-08-05 14:33:49,794 Producer INFO total items 9
2019-08-05 14:33:50,294 Producer INFO total items 10
2019-08-05 14:33:50,795 Producer INFO items produced 10\. Stopped
2019-08-05 14:33:51,789 Consumer INFO consumed 1 item
2019-08-05 14:33:51,790 Producer INFO total items 10
2019-08-05 14:33:52,290 Producer INFO items produced 10\. Stopped
2019-08-05 14:33:53,790 Consumer INFO consumed 1 item
2019-08-05 14:33:53,790 Producer INFO total items 10
2019-08-05 14:33:54,291 Producer INFO items produced 10\. Stopped
2019-08-05 14:33:55,790 Consumer INFO consumed 1 item
2019-08-05 14:33:55,791 Producer INFO total items 10
2019-08-05 14:33:56,291 Producer INFO items produced 10\. Stopped
2019-08-05 14:33:57,791 Consumer INFO consumed 1 item
2019-08-05 14:33:57,791 Producer INFO total items 10
2019-08-05 14:33:58,292 Producer INFO items produced 10\. Stopped
2019-08-05 14:33:59,791 Consumer INFO consumed 1 item
2019-08-05 14:33:59,791 Producer INFO total items 10
2019-08-05 14:34:00,292 Producer INFO items produced 10\. Stopped
2019-08-05 14:34:01,791 Consumer INFO consumed 1 item
2019-08-05 14:34:01,791 Producer INFO total items 10
2019-08-05 14:34:02,291 Producer INFO items produced 10\. Stopped
2019-08-05 14:34:03,791 Consumer INFO consumed 1 item
2019-08-05 14:34:03,792 Producer INFO total items 10
2019-08-05 14:34:05,792 Consumer INFO consumed 1 item
2019-08-05 14:34:07,793 Consumer INFO consumed 1 item
2019-08-05 14:34:09,794 Consumer INFO consumed 1 item
2019-08-05 14:34:11,795 Consumer INFO consumed 1 item
2019-08-05 14:34:13,795 Consumer INFO consumed 1 item
2019-08-05 14:34:15,833 Consumer INFO consumed 1 item
2019-08-05 14:34:17,833 Consumer INFO consumed 1 item
2019-08-05 14:34:19,833 Consumer INFO consumed 1 item
2019-08-05 14:34:21,834 Consumer INFO consumed 1 item
2019-08-05 14:34:23,835 Consumer INFO consumed 1 item

还有更多…

看到 Python 内部的条件同步机制非常有趣。内部class _Condition在类构造函数未传递现有锁时创建一个RLock()对象。此外,当调用acquire()released()时,将管理锁:

class _Condition(_Verbose):
    def __init__(self, lock=None, verbose=None):
        _Verbose.__init__(self, verbose)
        if lock is None:
            lock = RLock()
        self.__lock = lock

使用事件进行线程同步

事件是一个用于线程间通信的对象。一个线程等待信号,而另一个线程输出它。基本上,一个event对象管理一个内部标志,可以通过clear()将其设置为false,通过set()将其设置为true,并通过is_set()进行测试。

一个线程可以通过wait()方法持有信号,该方法通过set()方法发送调用。

准备工作

要通过事件对象理解线程同步,让我们看看生产者/消费者问题。

如何做…

再次,为了解释如何通过事件同步线程,我们将参考生产者/消费者问题。该问题描述了两个进程,一个生产者和一个消费者,他们共享一个固定大小的公共缓冲区。生产者的任务是生成项并将它们存入连续的缓冲区。同时,消费者将使用生产的项,并时不时地从缓冲区中移除它们。

问题在于确保生产者在缓冲区满时不处理新数据,而消费者在缓冲区空时不寻找数据。

现在,让我们看看如何通过使用事件语句进行线程同步来实现消费者/生产者问题:

  1. 在这里,相关库按如下方式导入:
import logging
import threading
import time
import random
  1. 然后,我们定义日志输出格式。这有助于清楚地可视化正在发生的事情:
LOG_FORMAT = '%(asctime)s %(threadName)-17s %(levelname)-8s %\
             (message)s'
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)
  1. 设置items列表。此参数将由ConsumerProducer类使用:
items = []
  1. event参数定义如下。此参数将用于同步线程间的通信:
event = threading.Event()
  1. Consumer类使用项列表和Event()函数初始化。在run方法中,消费者等待要消费的新项。当项到达时,它从item列表中弹出:
class Consumer(threading.Thread):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)

    def run(self):
        while True:
            time.sleep(2)
            event.wait()
            item = items.pop()
            logging.info('Consumer notify: {} popped by {}'\
                        .format(item, self.name))
  1. Producer类使用项列表和Event()函数初始化。与使用condition对象的示例不同,项列表不是全局的,而是作为参数传递:
class Producer(threading.Thread):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)
  1. 在为每个创建的项的run方法中,Producer类将其追加到项列表中,然后通知事件:
    def run(self):
        for i in range(5):
            time.sleep(2)
            item = random.randint(0, 100)
            items.append(item)
            logging.info('Producer notify: item {} appended by\ 
                         {}'\.format(item, self.name))
  1. 您需要采取两个步骤来完成此操作,第一步如下:
            event.set()
            event.clear()
  1. t1线程将值追加到列表中,然后设置事件以通知消费者。消费者的wait()调用停止阻塞,并从列表中检索整数:
if __name__ == "__main__":
    t1 = Producer()
    t2 = Consumer()

    t1.start()
    t2.start()

    t1.join()
    t2.join()

它是如何工作的…

所有在ProducerConsumer类之间的操作都可以通过以下方案轻松恢复:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/915fcac0-899f-465a-9fd2-ce18a24bb0f0.png

使用事件对象进行线程同步

特别是,ProducerConsumer类具有以下行为:

  • Producer获取一个锁,向队列中添加一个项目,并通过set event通知Consumer。然后它休眠,直到收到要添加的新项目。

  • Consumer获取一个块并开始在一个连续循环中监听元素。当事件到达时,消费者放弃该块,从而允许其他生产者/消费者进入并获取该块。如果Consumer被重新激活,那么它将通过安全地处理队列中的新项目来重新获取锁:

2019-02-02 18:23:35,125 Thread-1 INFO Producer notify: item 68 appended by Thread-1
2019-02-02 18:23:35,133 Thread-2 INFO Consumer notify: 68 popped by Thread-2
2019-02-02 18:23:37,138 Thread-1 INFO Producer notify: item 45 appended by Thread-1
2019-02-02 18:23:37,143 Thread-2 INFO Consumer notify: 45 popped by Thread-2
2019-02-02 18:23:39,148 Thread-1 INFO Producer notify: item 78 appended by Thread-1
2019-02-02 18:23:39,153 Thread-2 INFO Consumer notify: 78 popped by Thread-2
2019-02-02 18:23:41,158 Thread-1 INFO Producer notify: item 22 appended by Thread-1
2019-02-02 18:23:43,173 Thread-1 INFO Producer notify: item 48 appended by Thread-1
2019-02-02 18:23:43,178 Thread-2 INFO Consumer notify: 48 popped by Thread-2

使用屏障进行线程同步

有时,一个应用程序可以被划分为阶段,规则是如果首先,该进程的所有线程都完成了它们自己的任务,则没有进程可以继续。一个屏障实现了这个概念:一个完成了其阶段的线程调用原始屏障并停止。当所有涉及的线程都完成了它们的执行阶段并也调用了原始屏障后,系统将解锁它们,允许线程移动到后续阶段。

准备就绪

Python 的 threading 模块通过**Barrier**类实现屏障。在下一节中,我们将学习如何在一个非常简单的示例中使用这种同步机制。

如何做到这一点…

在这个例子中,我们模拟了一场有三名参与者HueyDeweyLouie的跑步比赛,其中屏障被类比为终点线。

此外,当所有三名参赛者都越过终点线时,比赛可以自行结束。

屏障是通过Barrier类实现的,其中必须指定要完成的线程数作为参数,以便移动到下一阶段:

from random import randrange
from threading import Barrier, Thread
from time import ctime, sleep

num_runners = 3
finish_line = Barrier(num_runners)
runners = ['Huey', 'Dewey', 'Louie']

def runner():
    name = runners.pop()
    sleep(randrange(2, 5))
    print('%s reached the barrier at: %s \n' % (name, ctime()))
    finish_line.wait()

def main():
    threads = []
    print('START RACE!!!!')
    for i in range(num_runners):
        threads.append(Thread(target=runner))
        threads[-1].start()
    for thread in threads:
        thread.join()
    print('Race over!')

if __name__ == "__main__":
    main()

它是如何工作的…

首先,我们将跑者的数量设置为num_runners = 3,以便通过Barrier指令在下一行设置最终目标。跑者被设置在跑者列表中;每个跑者的到达时间将在runner函数中使用randrange指令确定。

当一名跑者到达终点线时,调用wait方法,这将阻塞所有已经调用该方法的跑者(线程)。输出如下:

START RACE!!!!
Dewey reached the barrier at: Sat Feb 2 21:44:48 2019 

Huey reached the barrier at: Sat Feb 2 21:44:49 2019 

Louie reached the barrier at: Sat Feb 2 21:44:50 2019 

Race over!

在这种情况下,Dewey赢得了比赛。

使用队列进行线程通信

当线程需要共享数据或资源时,多线程可能会变得复杂。幸运的是,threading 模块提供了许多同步原语,包括信号量、条件变量、事件和锁。

然而,使用queue模块被认为是一种最佳实践。实际上,队列更容易处理,并且使得线程编程更加安全,因为它有效地将所有对单个线程资源的访问引导到一个方向,并允许设计出更清晰、更易读的模式。

准备就绪

我们将简单地考虑这些队列方法:

  • put(): 将一个项目放入队列

  • get(): 从队列中移除并返回一个项目

  • task_done(): 需要在每次处理完一个项目时调用

  • join(): 阻塞直到所有项目都已被处理

如何做到这一点…

在这个例子中,我们将看到如何使用threading模块与queue模块。此外,这里有两个实体试图共享一个公共资源,一个队列。代码如下:

from threading import Thread
from queue import Queue
import time
import random

class Producer(Thread):
    def __init__(self, queue):
        Thread.__init__(self)
        self.queue = queue
    def run(self):
        for i in range(5):
            item = random.randint(0, 256)
            self.queue.put(item)
            print('Producer notify : item N°%d appended to queue by\ 
                  %s\n'\
                  % (item, self.name))
            time.sleep(1)

class Consumer(Thread):
    def __init__(self, queue):
        Thread.__init__(self)
        self.queue = queue

    def run(self):
        while True:
            item = self.queue.get()
            print('Consumer notify : %d popped from queue by %s'\
                  % (item, self.name))
            self.queue.task_done()

if __name__ == '__main__':
    queue = Queue()
    t1 = Producer(queue)
    t2 = Consumer(queue)
    t3 = Consumer(queue)
    t4 = Consumer(queue)

    t1.start()
    t2.start()
    t3.start()
    t4.start()

    t1.join()
    t2.join()
    t3.join()
    t4.join()

它是如何工作的…

首先,使用producer类,我们不需要传递整数列表,因为我们使用队列来存储生成的整数。

producer类中的线程生成整数,并在for循环中将它们放入队列中。producer类使用Queue.put(item[, block[, timeout]])在队列中插入数据。它具有在队列中插入数据之前获取锁的逻辑。

有两种可能性:

  • 如果可选参数blocktruetimeoutNone(这是我们示例中使用的默认情况),那么我们需要阻塞,直到有空闲槽位可用。如果timeout是一个正数,那么它最多阻塞timeout秒,如果在那个时间内没有空闲槽位,则引发满异常。

  • 如果blockfalse,那么如果立即有空闲槽位,则将项目放入队列中,否则引发满异常(在这种情况下忽略超时)。在这里,put检查队列是否已满,然后内部调用wait,之后生产者开始等待。

接下来是consumer类。线程从队列中获取整数,并通过使用task_done来指示它已完成对该整数的工作。consumer类使用Queue.get([block[, timeout]])并在从队列中移除数据之前获取锁。如果队列是空的,消费者将被置于等待状态。最后,在main函数中,我们创建了四个线程,一个用于producer类,三个用于consumer类。

输出应该像这样:

Producer notify : item N°186 appended to queue by Thread-1
Consumer notify : 186 popped from queue by Thread-2

Producer notify : item N°16 appended to queue by Thread-1
Consumer notify : 16 popped from queue by Thread-3

Producer notify : item N°72 appended to queue by Thread-1
Consumer notify : 72 popped from queue by Thread-4

Producer notify : item N°178 appended to queue by Thread-1
Consumer notify : 178 popped from queue by Thread-2

Producer notify : item N°214 appended to queue by Thread-1
Consumer notify : 214 popped from queue by Thread-3

更多…

producer类和consumer类之间的所有操作都可以使用以下方案轻松恢复:

https://github.com/OpenDocCN/freelearn-python-zh/raw/master/docs/py-prll-prog-cb-2e/img/cb11a94d-258a-485f-a1b4-8954a860b41a.png

使用队列模块进行线程同步

  • Producer线程获取锁,然后在队列数据结构中插入数据。

  • Consumer线程从队列中获取整数。这些线程在从队列中移除数据之前获取锁。

如果队列为空,那么consumer线程将进入等待状态。

使用这个配方,关于基于线程的并行主义的章节就此结束。

Logo

openvela 操作系统专为 AIoT 领域量身定制,以轻量化、标准兼容、安全性和高度可扩展性为核心特点。openvela 以其卓越的技术优势,已成为众多物联网设备和 AI 硬件的技术首选,涵盖了智能手表、运动手环、智能音箱、耳机、智能家居设备以及机器人等多个领域。

更多推荐